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

如何在带聚合与suppress的Kafka Session时间窗口中获取最后一条窗口消息?

Kafka Streams SessionWindow结合Suppress无输出问题排查

问题场景

我需要实现一个功能:生产者向Kafka Streams发送整数列表,在指定时间窗口内对这些整数求和。使用SessionTimeWindow持续聚合直到最后一条消息到来,窗口关闭后通过Suppress获取最终求和值,但无论本地运行拓扑还是用TopologyTestDriver测试,都无法在输出主题中得到任何消息。

拓扑结构代码

StreamsBuilder builder = new StreamsBuilder();

builder.stream("input_topic", Consumed.with(Serdes.String(), new Serdes.ListSerde<Integer>()))
  .flatMapValues((value) -> value)
  .groupByKey(Grouped.with(Serdes.String(), Serdes.Integer()))
  .windowedBy(SessionWindows.ofInactivityGapWithNoGrace(Duration.ofSeconds(10)))
  .aggregate(
    ()-> 0,
    ((key, value, aggregate) -> aggregate + value),
    ((aggKey, aggOne, aggTwo) -> aggOne + aggTwo),
    Materialized.with(Serdes.String(), Serdes.Integer()))
  .suppress(Suppressed.untilWindowCloses(Suppressed.BufferConfig.unbounded()))
  .toStream()
  .map((key, value) -> new KeyValue<String, Integer>(key.key(), value))
  .to("output_topic", Produced.with(Serdes.String(), Serdes.Integer()));

Topology topology = builder.build();

测试配置与代码

TestTopologyDriver配置

Properties properties = new Properties();
properties.put(StreamsConfig.APPLICATION_ID_CONFIG, "app");        
properties.put(StreamsConfig.TOPOLOGY_OPTIMIZATION_CONFIG, StreamsConfig.OPTIMIZE);       
properties.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "test");       
properties.put(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 0);

测试代码

TopologyTestDriver topologyTestDriver = new TopologyTestDriver(topology, properties, Instant.now());

TestInputTopic<String, List<Integer>> inputTopic = topologyTestDriver.createInputTopic("input_topic", new StringSerializer(), new Serdes.ListSerde<Integer>().serializer());
TestOutputTopic<String, Integer> outputTopic = topologyTestDriver.createOutputTopic("output_topic", new StringDeserializer(), new IntegerDeserializer());

inputTopic.pipeInput("1", List.of(1,2,3));
inputTopic.pipeInput("1", List.of(1));
List<TestRecord<String, Integer>> outputList = outputTopic.readRecordsToList();

assertEquals(1, outputList.size());
assertEquals(7, outputList.get(0).value());

已尝试的操作

  • 在另一个线程运行TopologyTestDriver并让主线程休眠,无结果
  • 发送带时间戳的记录
  • 使用时间戳提取器
  • 调用TopologyTestDriver.advanceWallClockTime();

问题原因与解决方案

核心原因

SessionWindow的关闭条件是超过设置的无活动间隔(10秒)没有新消息,而Suppress.untilWindowCloses只会在窗口完全关闭后才输出聚合结果。你的测试代码没有让时间推进到窗口关闭的时间点,因此Suppress一直缓存结果,不会输出。

修正方案

  1. 测试环境下:发送完所有测试消息后,必须明确推进时钟时间超过无活动间隔,触发窗口关闭:
inputTopic.pipeInput("1", List.of(1,2,3));
inputTopic.pipeInput("1", List.of(1));
// 推进时间超过session的10秒无活动间隔,确保窗口关闭
topologyTestDriver.advanceWallClockTime(Duration.ofSeconds(11));
List<TestRecord<String, Integer>> outputList = outputTopic.readRecordsToList();

assertEquals(1, outputList.size());
assertEquals(7, outputList.get(0).value());
  1. 本地运行环境下:

    • 确保最后一条消息发送后,等待至少10秒,让窗口触发关闭逻辑
    • 如果生产者持续发送消息,窗口会保持活跃状态,不会关闭,也就不会有输出,这符合SessionWindow的设计逻辑
  2. 额外注意点:

    • SessionWindows.ofInactivityGapWithNoGrace的NoGrace参数意味着窗口关闭后不会处理延迟消息,确保你的消息时间戳符合预期
    • 确认COMMIT_INTERVAL_MS_CONFIG设置为0不会影响窗口关闭逻辑,这里的配置是正常的

内容的提问来源于stack exchange,提问作者Adrián Polo Alcaide

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 02:40:26