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

无Key场景下基于Kafka Streams实现滑动窗口聚合的技术咨询

问题描述
  • 现有Kafka Topic仅存储JSON格式的Value,无Key,Value示例:
{
"member_no" : "123",
"item_no": "item_123",
"category_no": "Category_123",
"order_no": "Order_123",
"datetime": "2022-1-11 09:11:00"
}
  • 需求:用Kafka Streams实现5分钟滑动窗口内,按item_no分组统计member_no的去重数量(对应SQL逻辑:select item_no, count(distinct member_no) from table group by item_no window 5min sliding)
  • 开发环境:Java + Spring(STS)
  • 疑问:现有代码的groupBy部分存疑,且知晓groupBy会触发重分区影响性能,想咨询是否需要预生成Key存入新Topic后再处理?
解决方案建议

1. 直接使用groupBy的可行性

如果数据量处于常规水平(比如单分区TPS在千级以内),直接用groupBy完全可行:

  • 核心步骤:先将JSON反序列化为POJO,再通过groupBy((key, value) -> value.getItemNo())按item_no分组,最后绑定滑动窗口并执行聚合。
  • 关于性能:groupBy确实会触发重分区——因为原Topic无有效Key,Kafka Streams必须把相同item_no的数据路由到同一分区才能完成正确聚合。但只要集群资源充足,这个性能损耗在绝大多数业务场景下是可接受的。

2. 预生成Key存入新Topic的优化方案

如果数据量极大,或对延迟要求极高,预生成Key是更优的性能优化方案:

  • 实现思路:先编写一个轻量的Kafka Streams(或Consumer+Producer)任务,消费原Topic消息后,将item_no作为新Key,原Value保持不变写入新Topic。
  • 核心优势:后续聚合任务可直接调用groupByKey(),彻底避免groupBy带来的额外重分区开销,性能提升显著。
  • 注意事项:预生成Key的任务要保证Exactly-Once语义,可通过Kafka事务机制避免消息重复或丢失。

代码示例(Spring Kafka Streams)

直接使用groupBy的实现

@Bean
public KStream<String, OrderEvent> kStream(StreamsBuilder streamsBuilder) {
    // 配置JSON与POJO的序列化/反序列化器
    Serde<OrderEvent> orderEventSerde = Serdes.serdeFrom(
        new JsonSerializer<>(),
        new JsonDeserializer<>(OrderEvent.class)
    );

    KStream<String, OrderEvent> stream = streamsBuilder.stream("original-topic", Consumed.with(Serdes.String(), orderEventSerde));

    stream.groupBy((key, value) -> value.getItemNo(), Grouped.with(Serdes.String(), orderEventSerde))
          .windowedBy(SlidingWindows.of(Duration.ofMinutes(5)))
          .aggregate(
              () -> new HashSet<>(), // 用HashSet实现member_no去重
              (itemNo, event, memberSet) -> {
                  memberSet.add(event.getMemberNo());
                  return memberSet;
              },
              Materialized.<String, Set<String>, WindowStore<Bytes, byte[]>>as("member-count-store")
                          .withKeySerde(Serdes.String())
                          .withValueSerde(Serdes.serdeFrom(new JsonSerializer<>(), new JsonDeserializer<>(Set.class)))
          )
          .mapValues(memberSet -> memberSet.size()) // 转换为去重后的数量
          .toStream()
          .to("aggregated-result-topic", Produced.with(WindowedSerdes.timeWindowedSerdeFrom(String.class, 5*60*1000), Serdes.Integer()));

    return stream;
}

// 对应的POJO实体类
public class OrderEvent {
    private String memberNo;
    private String itemNo;
    private String categoryNo;
    private String orderNo;
    private LocalDateTime datetime;

    // 省略getter、setter、构造方法
}

预生成Key的前置任务+聚合任务

// 前置任务:生成以item_no为Key的新Topic
@Bean
public KStream<String, OrderEvent> keyPreprocessingStream(StreamsBuilder streamsBuilder) {
    Serde<OrderEvent> orderEventSerde = Serdes.serdeFrom(new JsonSerializer<>(), new JsonDeserializer<>(OrderEvent.class));

    KStream<String, OrderEvent> stream = streamsBuilder.stream("original-topic", Consumed.with(Serdes.String(), orderEventSerde));

    stream.selectKey((key, value) -> value.getItemNo())
          .to("topic-with-item-key", Produced.with(Serdes.String(), orderEventSerde));

    return stream;
}

// 优化后的聚合任务
@Bean
public KStream<String, OrderEvent> aggregationStream(StreamsBuilder streamsBuilder) {
    Serde<OrderEvent> orderEventSerde = Serdes.serdeFrom(new JsonSerializer<>(), new JsonDeserializer<>(OrderEvent.class));

    KStream<String, OrderEvent> stream = streamsBuilder.stream("topic-with-item-key", Consumed.with(Serdes.String(), orderEventSerde));

    stream.groupByKey(Grouped.with(Serdes.String(), orderEventSerde))
          .windowedBy(SlidingWindows.of(Duration.ofMinutes(5)))
          .aggregate(
              () -> new HashSet<>(),
              (itemNo, event, memberSet) -> {
                  memberSet.add(event.getMemberNo());
                  return memberSet;
              },
              Materialized.as("member-count-store-optimized")
                          .withKeySerde(Serdes.String())
                          .withValueSerde(Serdes.serdeFrom(new JsonSerializer<>(), new JsonDeserializer<>(Set.class)))
          )
          .mapValues(memberSet -> memberSet.size())
          .toStream()
          .to("aggregated-result-topic", Produced.with(WindowedSerdes.timeWindowedSerdeFrom(String.class, 5*60*1000), Serdes.Integer()));

    return stream;
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 00:45:40