无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
相关产品推荐
相关产品推荐

