迁移到命名内部主题后Kafka Streams Join出现空指针异常问题
问题解答
1. 问题根本原因
你遇到的是Kafka Streams旧版本的已知缺陷:在KStream与KTable外键join的场景下,处理订阅状态的SubscriptionWrapperSerializer序列化逻辑没有做空值校验。当手动变更内部主题命名、执行应用重置后,内部订阅主题中会出现主键哈希字段为null的无效记录,代码直接读取该字段序列化时就会抛出空指针异常。
2. 升级到Kafka 2.7+是否可以解决
可以解决。对应缺陷修复确实已经正式合入Kafka 2.7及后续版本,修复逻辑就是在序列化前增加了空值判断,遇到无效空字段时会做兼容处理,不会再触发崩溃。对应你使用的Confluent Platform版本,升级到6.1及以上版本即可覆盖该修复。
3. 其他可采用的规避方案
如果暂时无法完成版本升级,可以尝试以下方案解决问题:
- 彻底清理所有手动命名的FK-JOIN相关内部主题和对应的本地状态存储,重新执行应用重置流程后再启动应用,确保没有残留的无效历史记录触发空指针
- 在KStream和KTable执行join操作的上游,提前过滤掉外键为null的记录,从源头上避免生成无效的订阅记录写入内部主题
- 临时替换kafka-streams依赖包中的
SubscriptionWrapperSerde类,自行给序列化逻辑增加空值判断的补丁,在升级前临时规避问题
内容的提问来源于stack exchange,提问作者qlangman
相关产品推荐
相关产品推荐

