为何Kafka Topic收到消息但ksqlDB Stream未同步?
问题排查与解决方案
1. 字段大小写不匹配(最可能原因)
你的流定义中,VALUE字段是NUM INTEGER(大写),但插入的JSON数据里字段是num: 2(小写)。ksqlDB默认对JSON字段名大小写敏感,会导致无法解析出NUM字段,进而消息在流中无法正常显示。
解决方法:
- 要么修改流定义,将字段改为小写:
CREATE STREAM PROVA (ID STRING KEY, num INTEGER) WITH (KAFKA_TOPIC='topic_prova', KEY_FORMAT='KAFKA', PARTITIONS=1, VALUE_FORMAT='JSON'); - 要么修改插入代码,将字段名改为大写:
const row = { id: "value", NUM: 2 }; - 或者开启ksqlDB的JSON大小写不敏感解析:在ksqlDB配置中添加
ksql.schema.registry.parse.json.case.insensitive=true,然后重启ksqlDB服务。
2. 流的KEY定义与实际消息键不匹配
流定义中ID STRING KEY表示ID对应Kafka消息的键(Key),而非Value中的字段。但你的插入代码中,id: "value"是作为Value的一部分发送的,Kafka消息的实际Key可能为空。这会导致流中的ID字段为NULL,如果你的查询过滤了非空ID,就看不到数据。
解决方法:
- 如果
ID应该是Value中的字段,修改流定义,去掉KEY修饰:CREATE STREAM PROVA (ID STRING, NUM INTEGER) WITH (KAFKA_TOPIC='topic_prova', KEY_FORMAT='KAFKA', PARTITIONS=1, VALUE_FORMAT='JSON'); - 如果
ID确实要作为Kafka消息的键,修改插入代码,指定消息键。根据ksqldb-client的用法,需要在插入时显式设置键:// 参考客户端API调整,显式指定消息键 const {data, status, error } = await client.insertInto("PROVA", row, { key: row.id });
3. 验证消息格式与偏移量策略
- 偏移量检查:执行push查询等待新消息,确认是否能获取数据:
SELECT * FROM PROVA EMIT CHANGES LIMIT 5; - 消息格式验证:直接查看Kafka Topic中的消息,确认JSON格式与流定义匹配:
# 使用kafkacat工具查看消息 kafkacat -b <kafka-broker-address> -t topic_prova -C -q -J | head -5
4. 检查ksqlDB解析日志
查看ksqlDB Pod的日志,排查是否有消息解析失败的报错:
kubectl logs <ksqldb-pod-name> -n <your-namespace> | grep -i error
如果出现Failed to deserialize value类的错误,说明数据格式与流定义不匹配,回到前面的字段大小写或类型问题排查。
内容的提问来源于stack exchange,提问作者Luana Mantovan
相关产品推荐
相关产品推荐

