Kafka队列中Avro Schema序列化责任方及使用相关技术咨询
Kafka Avro Schema 序列化相关问题解答
一、谁负责Avro Schema序列化?
Kafka 本身不做Avro序列化/反序列化工作,这个流程完全由生产者、消费者客户端处理:
- 生产者端:把JSON等结构化数据,按指定Avro Schema转成二进制后再发送到Kafka
- 消费者端:从Kafka拿到二进制消息,用对应的Avro Schema转成可读的结构化数据
Kafka服务器只负责存和转发二进制消息,不会碰消息内容的序列化逻辑。
二、Avro Schema的使用时机
创建Topic时不需要指定Avro Schema,Schema是在客户端发消息、消费消息时指定的。实际生产中一般配合Schema Registry统一管理:
- 生产者发消息时,要么把Schema注册到Registry(拿到Schema ID),要么直接引用已注册的ID,然后把ID和二进制消息一起发去Kafka
- 消费者消费时,通过消息里的Schema ID从Registry拉取对应Schema,再做反序列化
如果不用Registry,也可以在客户端本地加载Schema文件处理,但这种方式没法统一管理Schema版本,容易出问题。
三、命令行工具&REST Proxy示例
假设你已经有mykeyschema.avsc和myvalueschema.avsc两个Schema文件,以下是实操示例:
1. Kafka命令行工具(配合Schema Registry)
需确保你有kafka-avro-console-producer.sh和kafka-avro-console-consumer.sh(一般随Confluent Platform自带)
生产者发消息
# 加载本地键、值Schema,发送JSON格式消息 ./kafka-avro-console-producer.sh \ --broker-list localhost:9092 \ --topic test-avro-topic \ --property schema.registry.url=http://localhost:8081 \ --property key.schema.file=./mykeyschema.avsc \ --property value.schema.file=./myvalueschema.avsc
运行后直接输入JSON消息(比如{"keyField": 1, "valueField": "test"})即可发送。
消费者消费消息
./kafka-avro-console-consumer.sh \ --bootstrap-server localhost:9092 \ --topic test-avro-topic \ --property schema.registry.url=http://localhost:8081 \ --from-beginning
执行后会输出反序列化后的JSON格式消息。
2. Kafka REST Proxy的curl示例
第一步:注册Schema到Schema Registry
注册键Schema:
curl -X POST -H "Content-Type: application/vnd.schemaregistry.v1+json" \ --data '{"schema": "'"$(sed 's/"/\\"/g' mykeyschema.avsc)"'"}' \ http://localhost:8081/subjects/test-avro-topic-key/versions
注册值Schema:
curl -X POST -H "Content-Type: application/vnd.schemaregistry.v1+json" \ --data '{"schema": "'"$(sed 's/"/\\"/g' myvalueschema.avsc)"'"}' \ http://localhost:8081/subjects/test-avro-topic-value/versions
注:
sed 's/"/\\"/g'用来转义Schema文件里的双引号,避免JSON格式报错。
第二步:发送Avro格式消息
curl -X POST -H "Content-Type: application/vnd.kafka.avro.v2+json" \ --data '{ "records": [ { "key": {"keyField": 1}, "value": {"valueField": "test message"} } ] }' \ http://localhost:8082/topics/test-avro-topic
第三步:消费Avro格式消息
先创建消费者实例:
curl -X POST -H "Content-Type: application/vnd.kafka.v2+json" \ --data '{"name": "avro-consumer", "format": "avro", "auto.offset.reset": "earliest"}' \ http://localhost:8082/consumers/avro-group
订阅Topic并拉取消息:
curl -X GET -H "Accept: application/vnd.kafka.avro.v2+json" \ http://localhost:8082/consumers/avro-group/instances/avro-consumer/records?timeout=1000
用完关闭消费者:
curl -X DELETE \ http://localhost:8082/consumers/avro-group/instances/avro-consumer
内容的提问来源于stack exchange,提问作者Peter
相关产品推荐
相关产品推荐

