如何通过Kafka Console Consumer与Producer使用JSON Schema收发消息及可行性确认
使用kafka-console-producer/consumer收发带JSON Schema的Kafka消息
当然没问题!用 Kafka 自带的命令行工具收发带有 JSON Schema 的消息完全可行,不过前提是你已经部署好了 Confluent Schema Registry(JSON Schema 的存储、验证和获取都依赖它)。下面我给你一步步拆解具体操作:
一、用kafka-console-producer发送带JSON Schema的消息
你需要给生产者指定 Schema Registry 的地址,同时定义要使用的 JSON Schema(或者引用已注册的 Schema ID),确保发送的消息符合 Schema 规范。
方法1:直接指定JSON Schema
如果是第一次发送该 Schema 的消息,可以直接在命令中定义 Schema:
kafka-console-producer --broker-list <你的Kafka Broker地址,比如localhost:9092> --topic <目标Topic名称> \ --property schema.registry.url=<你的Schema Registry地址,比如localhost:8081> \ --property value.schema.type=json \ --property value.schema='{"type": "object", "properties": {"name": {"type": "string"}, "age": {"type": "integer"}, "email": {"type": "string"}}}'
运行命令后,在控制台输入符合 Schema 的 JSON 消息即可发送,比如:
{"name": "张三", "age": 28, "email": "zhangsan@example.com"}
参数说明:
--broker-list:指定 Kafka Broker 的地址schema.registry.url:指定 Schema Registry 的访问地址(默认端口8081)value.schema.type=json:告诉生产者消息体使用 JSON Schema 格式value.schema:定义具体的 JSON Schema 规则
方法2:引用已注册的Schema ID
如果该 Schema 已经在 Registry 中注册过,可以直接用 Schema ID 来指定,避免重复写 Schema:
kafka-console-producer --broker-list localhost:9092 --topic user-topic \ --property schema.registry.url=localhost:8081 \ --property value.schema.type=json \ --property value.schema.id=<已注册的Schema ID,比如1001>
二、用kafka-console-consumer接收带JSON Schema的消息
消费者需要配置 Schema Registry 地址,同时指定使用 JSON Schema 的反序列化器,才能正确解析消息。
运行以下命令:
kafka-console-consumer --bootstrap-server <你的Kafka Broker地址> --topic <目标Topic名称> \ --from-beginning \ --property schema.registry.url=<你的Schema Registry地址> \ --property value.deserializer=io.confluent.kafka.serializers.json.KafkaJsonSchemaDeserializer \ --property specific.avro.reader=false
参数说明:
--from-beginning:从 Topic 的起始位置开始消费(可选,根据需求决定)value.deserializer:指定使用 JSON Schema 专用的反序列化器specific.avro.reader=false:表示使用通用的 JSON 解析模式,无需预先生成对应的 Java 实体类
运行命令后,你就能看到解析后的原始 JSON 消息内容了。
注意事项
- 确保 Kafka Broker 和 Schema Registry 网络连通,地址配置正确
- 发送的消息必须严格符合指定的 JSON Schema,否则生产者会抛出验证错误,消息无法发送
- 不同版本的 Confluent 工具可能参数略有差异,如果遇到报错,可以检查工具版本对应的官方文档调整参数
内容的提问来源于stack exchange,提问作者Tokyo
相关产品推荐
相关产品推荐

