Avro DatumWriter write方法报错缺失encoder参数,求正确语法
Avro DatumWriter.write() 正确用法及Kafka Producer代码修复
错误原因
io.DatumWriter.write()方法需要两个必填参数:
- 要序列化的Avro数据
- 一个
Encoder实例(如BinaryEncoder),负责将数据编码为二进制格式
你之前的写法缺少了encoder参数,因此抛出TypeError。
正确用法
要正确序列化Avro数据,需要结合BytesIO字节流和BinaryEncoder:
- 创建
BytesIO对象,用于存储序列化后的二进制数据 - 初始化
BinaryEncoder,关联上述字节流 - 调用
DatumWriter.write(),传入数据和encoder - 从字节流中提取最终的二进制字节
修改后的核心代码
将原代码中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
相关产品推荐
相关产品推荐

