首次启动Kafka Streams应用时状态存储丢失消息问题咨询
Kafka Streams CDC应用首次启动丢数据问题咨询
问题背景
我们基于Kafka Streams实现了CDC应用,相关子拓扑结构对应指定示意图。table2主题由连接SQL数据库的Debezium创建,包含26000条数据。我们在应用中将该主题的消息key从string类型转为int类型,理论上table2的消息量、repartition内部主题的消息量、状态存储(state-store)的记录数应完全一致,但实际repartition主题的消息量大于状态存储的记录数,状态存储丢消息导致应用状态异常(已暂停Debezium连接器)。
触发条件
该问题仅发生在**首次启动应用(内部主题未创建)**时,重启已部署的应用无此问题。
环境配置
- Kafka Broker版本:3.2
- Kafka Streams客户端版本:2.8 / 3.2
- 仅配置两个参数:
CACHE_MAX_BYTES_BUFFERING_CONFIG设为0,NUM_STREAM_THREADS_CONFIG>1
已验证有效解决方案
- 首次启动使用单线程(
NUM_STREAM_THREADS_CONFIG=1),三者数量完全一致 - 预创建所有内部主题,避免启动时的重平衡操作,解决丢数据问题
问题流程推测
结合日志分析,多线程模式下,消费者协调器选定的Leader线程所分配的分区出现丢数据,推测流程如下:
- 多个消费者线程启动并通知协调器
- 协调器完成主题分区分配并选定Leader线程
- 应用自动创建所需内部主题
- Leader线程消费repartition主题消息,处理后未将变更刷入changelog主题就删除了repartition主题的消费位点
- Leader线程收到新的分区分配通知,触发重平衡
- Leader线程暂停当前负责的分区
- 重平衡完成后恢复分区消费
- Leader线程从中断位点拉取数据,导致早期未刷入changelog的消息丢失
单线程模式下无此重平衡触发流程,因此未出现问题。
咨询问题
- 上述流程推测是否正确?
- 该问题的根源是什么?
内容的提问来源于stack exchange,提问作者Youcef SEBIAT
相关产品推荐
相关产品推荐

