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

Kafka Streams表-表内连接产生重复消息问题咨询

Kafka Streams KTable内连接重复输出问题

背景与问题

我需要用Kafka Streams实现两个Kafka主题的连接:

  • 最初尝试用KStream(左表)和KTable(右表)做leftJoin,但频繁出现右值为空的异常情况;窗口连接也不适用我的业务场景。
  • 改为将两个主题都建模为KTable做内连接后,发现当向两个输入主题各发送1条同键消息时,结果流会出现2条相同的连接结果,触发两次连接逻辑。

测试与运行差异

  • 用TestTopologyDriver编写单元测试时结果正常,每对同键消息对应1条输出;
  • 但用TestCompanion搭配测试代理,或者直接运行应用时,就会出现重复输出。

环境与配置

  • 运行环境:Quarkus 2.16.x,Kafka Streams 3.3.2
  • 核心配置:
kafka-streams.cache.max.bytes.buffering=10240
kafka-streams.commit.interval.ms=500
kafka-streams.metadata.max.age.ms=500
kafka-streams.auto.offset.reset=earliest
kafka-streams.consumer.heartbeat.interval.ms=200

拓扑结构

子拓扑1用于测试时将单条输入消息转发到两个输入主题;子拓扑2实现两个KTable的内连接(外连接过滤空值后也存在同样重复问题):

Sub-topology: 1
    Source: KSTREAM-SOURCE-0000000041 (topics: [trigger])
      --> KSTREAM-FILTER-0000000044, KSTREAM-FILTER-0000000042
    Processor: KSTREAM-FILTER-0000000044 (stores: [])
      --> KSTREAM-MAPVALUES-0000000045
      <-- KSTREAM-SOURCE-0000000041
    Processor: KSTREAM-FILTER-0000000042 (stores: [])
      --> KSTREAM-SINK-0000000043
      <-- KSTREAM-SOURCE-0000000041
    Processor: KSTREAM-MAPVALUES-0000000045 (stores: [])
      --> KSTREAM-SINK-0000000046
      <-- KSTREAM-FILTER-0000000044
    Sink: KSTREAM-SINK-0000000043 (topic: value)
      <-- KSTREAM-FILTER-0000000042
    Sink: KSTREAM-SINK-0000000046 (topic: state)
      <-- KSTREAM-MAPVALUES-0000000045

  Sub-topology: 2
    Source: KSTREAM-SOURCE-0000000047 (topics: [value])
      --> KSTREAM-TOTABLE-0000000048
    Source: KSTREAM-SOURCE-0000000051 (topics: [state])
      --> KTABLE-SOURCE-0000000052
    Processor: KSTREAM-TOTABLE-0000000048 (stores: [KSTREAM-TOTABLE-STATE-STORE-0000000049])
      --> KTABLE-JOINOTHER-0000000055
      <-- KSTREAM-SOURCE-0000000047
    Processor: KTABLE-SOURCE-0000000052 (stores: [status-STATE-STORE-0000000050])
      --> KTABLE-JOINTHIS-0000000054
      <-- KSTREAM-SOURCE-0000000051
    Processor: KTABLE-JOINOTHER-0000000055 (stores: [status-STATE-STORE-0000000050])
      --> KTABLE-MERGE-0000000053
      <-- KSTREAM-TOTABLE-0000000048
    Processor: KTABLE-JOINTHIS-0000000054 (stores: [KSTREAM-TOTABLE-STATE-STORE-0000000049])
      --> KTABLE-MERGE-0000000053
      <-- KTABLE-SOURCE-0000000052
    Processor: KTABLE-MERGE-0000000053 (stores: [])
      --> KTABLE-TOSTREAM-0000000056
      <-- KTABLE-JOINTHIS-0000000054, KTABLE-JOINOTHER-0000000055
    Processor: KTABLE-TOSTREAM-0000000056 (stores: [])
      --> KSTREAM-MAPVALUES-0000000057
      <-- KTABLE-MERGE-0000000053
    Processor: KSTREAM-MAPVALUES-0000000057 (stores: [])
      --> KSTREAM-SINK-0000000058
      <-- KTABLE-TOSTREAM-0000000056
    Sink: KSTREAM-SINK-0000000058 (topic: joined)
      <-- KSTREAM-MAPVALUES-0000000057

补充说明

  • 两个输入主题的同键消息到达时间间隔极短(仅数毫秒);
  • 看到有建议在结果流上加去重处理器,但需要维护第三个状态存储,感觉不合理;
  • 未在官方文档中找到该重复触发问题的原因。

问题分析与解决方案

原因定位

这种重复输出是KTable连接的正常行为特性:

  1. KTable依赖状态存储工作,当第一个KTable的同键消息到达时,若第二个消息还未写入对应状态存储,不会输出连接结果;
  2. 当第二个KTable的同键消息到达时,会检测到第一个KTable状态存储中的匹配记录,触发连接并输出一条结果;
  3. 由于两条消息到达间隔极短,结合当前缓存与提交配置,第一个KTable的状态更新可能在第二个消息处理完成后,再次触发连接处理器重新计算,从而输出第二条重复结果。

解决方案

  1. 调整缓存与提交配置:

    • 调大kafka-streams.cache.max.bytes.buffering(比如设为10485760即10MB),让更多状态更新在内存缓存中合并,减少下游处理器的触发次数;
    • 适当增大kafka-streams.commit.interval.ms(比如设为1000),让状态存储的提交更平缓,避免短时间内多次触发连接计算。
  2. 使用suppress()操作去重:

    • 在KTable转流之后,用suppress()操作对同键结果进行窗口去重,无需额外维护状态存储(底层复用现有状态):
      joinedKTable.toStream()
                   .suppress(Suppressed.untilTimeLimit(Duration.ofMillis(500), 
                                                      Suppressed.BufferConfig.unbounded().withKeySerde(keySerde)))
                   .mapValues(...)
                   .to("joined");
      
    • 该操作会在指定时间窗口内只保留同键的最后一条结果,过滤重复输出,时间窗口可根据消息到达间隔调整。
  3. 优化消息发送顺序:

    • 如果业务允许,确保两个主题的同键消息按顺序发送,或让其中一个主题的消息延迟发送,避免短时间内两条同键消息同时触发状态更新。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 13:50:04