Kafka Stream消费失败:IIDR流入的记录无法进入流
Kafka Streams无法读取IIDR写入的AVRO消息?这里是问题根源和修复方案
首先,你的问题很典型:IIDR往test_iidr Topic发的AVRO消息,Kafka Streams死活读不到,但手动插入的同结构数据却正常。结合你提到的Kafka Connect用了JsonConverter,这大概率就是核心问题所在,下面一步步说清楚:
为什么会出现这个问题?
你要明确两个关键点:
- IIDR作为大型机的数据源,它是直接生成带Schema Registry标识的标准AVRO二进制消息写入Kafka的
- 但你的Kafka Connect配置了
org.apache.kafka.connect.json.JsonConverter——这个转换器会把接收到的消息当成JSON来处理,要么解析失败乱码,要么直接把AVRO二进制数据转成了JSON字符串写入Topic。最终导致Topic里的消息根本不是Kafka Streams期望的AVRO格式,自然读不到。
先验证消息格式是否异常
你可以用Confluent自带的工具直接检查test_iidr Topic的消息格式,执行这条命令:
kafka-avro-console-consumer --bootstrap-server <你的Kafka地址> --topic test_iidr --from-beginning
如果执行时出现解析错误,那实锤消息被Connect的JsonConverter篡改了;如果能正常输出AVRO结构的数据,那再去排查Kafka Streams的Schema配置。
修复步骤
1. 把Kafka Connect的转换器改成AVRO专用的
修改Connect的配置文件(或者Docker环境的环境变量),把key.converter和value.converter换成AVRO转换器,同时指定Schema Registry的地址:
key.converter=io.confluent.connect.avro.AvroConverter key.converter.schema.registry.url=http://<你的Schema Registry地址>:8081 value.converter=io.confluent.connect.avro.AvroConverter value.converter.schema.registry.url=http://<你的Schema Registry地址>:8081
这样Connect会正确识别IIDR发来的AVRO消息,不会修改原始格式。
2. 确认Kafka Streams/ksqlDB的Schema Registry配置
如果你用的是ksqlDB(从你写的CREATE STREAM语句来看应该是),要确保ksqlDB服务器的ksql.schema.registry.url配置正确指向你的Schema Registry,这样它才能获取对应的AVRO Schema来解析消息。
3. 重新生成测试数据
修改完Connect配置后,让IIDR重新发一批数据到test_iidr Topic,再用Kafka Streams尝试消费,应该就能正常接收了。
额外要检查的细节
- 核对Schema Registry里
test_iidr-value和test_iidr-key的Schema,和你CREATE STREAM定义的字段类型完全一致(比如DECIMAL(2,0)和DECIMAL(1,0)的精度刻度不能错) - 确认IIDR的配置里确实是用AVRO格式发送消息,并且正确配置了Schema Registry地址
内容的提问来源于stack exchange,提问作者Sanjay Nayak
相关产品推荐
相关产品推荐

