You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Avro DatumWriter write方法报错缺失encoder参数,求正确语法

Avro DatumWriter.write() 正确用法及Kafka Producer代码修复

错误原因

io.DatumWriter.write()方法需要两个必填参数:

  • 要序列化的Avro数据
  • 一个Encoder实例(如BinaryEncoder),负责将数据编码为二进制格式

你之前的写法缺少了encoder参数,因此抛出TypeError。

正确用法

要正确序列化Avro数据,需要结合BytesIO字节流和BinaryEncoder:

  1. 创建BytesIO对象,用于存储序列化后的二进制数据
  2. 初始化BinaryEncoder,关联上述字节流
  3. 调用DatumWriter.write(),传入数据和encoder
  4. 从字节流中提取最终的二进制字节

修改后的核心代码

将原代码中KafkaProducer的value_serializer替换为以下逻辑(建议单独定义序列化函数,比lambda更清晰):

import io as py_io  # 重命名标准库io,避免与avro.io冲突
from avro import io

def avro_serializer(data, avro_schema):
    writer = io.DatumWriter(avro_schema)
    bytes_writer = py_io.BytesIO()
    encoder = io.BinaryEncoder(bytes_writer)
    writer.write(data, encoder)
    return bytes_writer.getvalue()

# 初始化Producer时传入序列化逻辑
producer = KafkaProducer(
    # ... 其他配置保持不变 ...
    value_serializer=lambda m: avro_serializer(m, avro_schema)
)

完整修复后的代码

import os
import json
import time
import random
import io as py_io  # 重命名标准库io,避免与avro.io冲突
from kafka import KafkaProducer
from dotenv import load_dotenv
from avro import schema, io
from avro.datafile import DataFileWriter

load_dotenv()

BOOTSTRAP_SERVERS = os.getenv("KAFKA_BOOTSTRAP_SERVERS").split(",")
SSL_TRUSTSTORE = os.getenv("KAFKA_SSL_TRUSTSTORE")
TOPIC_NAME = os.getenv("KAFKA_TOPIC")
LINGER_DURATION = int(os.getenv("LINGER_DURATION"))
MESSAGE_BATCH = int(os.getenv("MESSAGE_BATCH"))


def random_with_N_digits(n):
    range_start = 10 ** (n - 1)
    range_end = (10 ** n) - 1
    return random.randint(range_start, range_end)


def produce(message):
    print(f'Sending {message}')
    producer.send(TOPIC_NAME, value=message)


def avro_serializer(data, avro_schema):
    writer = io.DatumWriter(avro_schema)
    bytes_writer = py_io.BytesIO()
    encoder = io.BinaryEncoder(bytes_writer)
    writer.write(data, encoder)
    return bytes_writer.getvalue()


# Define the Avro schema for the message
schema_str = """
{
    "namespace": "example.avro",
    "type": "record",
    "name": "Message",
    "fields": [
        {"name": "id", "type": "int"},
        {"name": "name", "type": "string"}
    ]
}
"""
avro_schema = schema.Parse(schema_str)

# 初始化Producer时传入序列化函数
producer = KafkaProducer(
    security_protocol="SSL",
    ssl_cafile=SSL_TRUSTSTORE,
    ssl_check_hostname=False,
    bootstrap_servers=BOOTSTRAP_SERVERS,
    linger_ms=LINGER_DURATION,
    value_serializer=lambda m: avro_serializer(m, avro_schema)
)


def startStream():
    msg_count = 0

    try:
        while True:
            # Create a message using the Avro schema
            message = {"id": random.randint(0, 9999999), "name": str(random_with_N_digits(12))}

            # Publish the message to a Kafka topic
            produce(message)

            msg_count += 1
            if msg_count % MESSAGE_BATCH == 0:
                time.sleep(1)

    except KeyboardInterrupt:
        print('interrupted!')
        print(f"Published {msg_count} messages")


startStream()

内容的提问来源于stack exchange,提问作者william007

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.27 19:23:12