You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

首次启动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线程所分配的分区出现丢数据,推测流程如下:

  1. 多个消费者线程启动并通知协调器
  2. 协调器完成主题分区分配并选定Leader线程
  3. 应用自动创建所需内部主题
  4. Leader线程消费repartition主题消息,处理后未将变更刷入changelog主题就删除了repartition主题的消费位点
  5. Leader线程收到新的分区分配通知,触发重平衡
  6. Leader线程暂停当前负责的分区
  7. 重平衡完成后恢复分区消费
  8. Leader线程从中断位点拉取数据,导致早期未刷入changelog的消息丢失

单线程模式下无此重平衡触发流程,因此未出现问题。

咨询问题

  1. 上述流程推测是否正确?
  2. 该问题的根源是什么?

内容的提问来源于stack exchange,提问作者Youcef SEBIAT

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.13 01:35:24