请求提供向Kafka发送Binary与JSON编码Avro消息的测试示例
向Kafka发送Binary/JSON编码Avro消息示例
Avro支持Binary和JSON两种编码格式,以下是分别发送两种编码消息的Python示例:
1. 发送Binary编码Avro消息
confluent-kafka的AvroProducer默认使用Binary编码,结合Schema Registry使用:
from confluent_kafka import avro from confluent_kafka.avro import AvroProducer # 定义Avro Schema value_schema = avro.loads(""" { "type": "record", "name": "User", "fields": [ {"name": "name", "type": "string"}, {"name": "age", "type": ["int", "null"]} ] } """) # 生产者配置 producer_conf = { "bootstrap.servers": "localhost:9092", "schema.registry.url": "http://localhost:8081" } # 初始化AvroProducer(默认Binary编码) producer = AvroProducer(producer_conf, default_value_schema=value_schema) # 构造并发送消息 msg_data = {"name": "Alice", "age": 30} producer.produce(topic="avro-binary-topic", value=msg_data) producer.flush()
2. 发送JSON编码Avro消息
使用avro库手动将数据序列化为Avro JSON格式,通过普通KafkaProducer发送:
from confluent_kafka import Producer import avro.schema from avro.io import DatumWriter, JsonEncoder import io # 加载Avro Schema schema = avro.schema.parse(""" { "type": "record", "name": "User", "fields": [ {"name": "name", "type": "string"}, {"name": "age", "type": ["int", "null"]} ] } """) # 普通Kafka生产者配置 producer_conf = {"bootstrap.servers": "localhost:9092"} producer = Producer(producer_conf) # 构造消息数据 msg_data = {"name": "Bob", "age": 25} # 序列化为Avro JSON格式 writer = DatumWriter(schema) byte_io = io.BytesIO() json_encoder = JsonEncoder(schema, byte_io) writer.write(msg_data, json_encoder) json_encoder.flush() avro_json_str = byte_io.getvalue().decode("utf-8") # 发送JSON编码消息 producer.produce(topic="avro-json-topic", value=avro_json_str.encode("utf-8")) producer.flush()
注意事项
- Binary编码是Avro的紧凑二进制格式,性能高,是Kafka生态中Avro消息的主流选择。
- JSON编码可读性强,但体积更大,适合调试场景。
- 依赖安装:执行
pip install confluent-kafka avro安装所需库。
内容的提问来源于stack exchange,提问作者journeyman
相关产品推荐
相关产品推荐

