如何实现跨Kafka Stream去重并保留同流重复记录?
Kafka Streams 双流合并问题及解决方案
需求说明
我有两个Kafka Stream,合并要求如下:
- 流1:无Key,时间戳非唯一,示例数据:
1,3,5,7,9 - 流2:为流1消息分配Key并做可逆值修改,时间戳与流1不匹配,可能存在对应原消息的多Key重复项(如
a:1aug1,b:1aug2),也可能缺失原流消息、出现乱序 - 合并规则:
- 仅存在于流1的消息需保留;
- 同时存在于两个流的消息,保留流2的Key:Value;
- 仅存在于流2的消息需保留;
- 同一消息在流2有多个副本时,保留所有流2副本并丢弃流1的对应消息;
- 额外要求:合并结果的顺序需与流1一致
现有代码问题
当前使用.reduce的代码可满足前三个规则,但无法实现第四个规则(无法保留流2的所有重复副本并丢弃对应流1消息),现有Java代码如下:
// pull in the two input streams (raw and augmented) KStream<String, String> rawInputStream = builder.stream(rawTopic, Consumed.with(Serdes.String(), Serdes.String())); KStream<String, String> augmentedInputStream = builder.stream(augTopic, Consumed.with(Serdes.String(), Serdes.String())); // map to a common key, so we can easily compare the messages. Store the original keys in the value also, so we can reuse them later. // The raw input won't have any original key, so use a blank string. KStream<String, CustomMessageDetailsWithKeyAndOrigin> mappedRawInputStream = rawInputStream .map((key, value) -> KeyValue.pair(getCommonKeyFromRawInputStream(value) , new CustomMessageDetailsWithKeyAndOrigin(getValueFromRawInputStream(value),key == null ? "" : key, OriginStream.RAW))); KStream<String, CustomMessageDetailsWithKeyAndOrigin> mappedAugmentedInputStream = augmentedInputStream .map((key, value) -> KeyValue.pair(getCommonKeyFromAugmentedInputStream(value) , new CustomMessageDetailsWithKeyAndOrigin(value, key == null ? "" : key, OriginStream.AUGMENTED))); // the outer join here will do a pairwise comparison across all records with a matching key, and just keep the records from the aggregated feed unless no agg value exists. KStream<String, CustomMessageDetailsWithKeyAndOrigin> mergedStream = mappedRawInputStream.outerJoin(mappedAugmentedInputStream, (value1,value2)-> { if (value2 == null) { // no augmented message // log return value1; } else if(value1 == null) {} // no raw message - log. return value2; } // Add a time-based join window to allow for time differences and sequence issues ,JoinWindows.ofTimeDifferenceAndGrace(window, windowGrace)); // We'll potentially have duplicates now - e.g. one from each input stream, or two from one?; so group by key to bring together the records that share a key KGroupedStream<String, CustomMessageDetailsWithKeyAndOrigin> groupedStream = mergedStream.groupByKey(); // ungroup the records again, reducing to remove duplicates. KStream<String, CustomMessageDetailsWithKeyAndOrigin> reducedStream = groupedStream.aggregate(LinkedHashSet<CustomMessageDetailsWithKeyAndOrigin>::new, (key, value, aggregate) -> { if (value != null) { boolean added = aggregate.add(value); // won't add again if it's a duplicate if (!added){} // duplicate - log it. } return aggregate; }).toStream().flatMapValues(value->value); // grab the original key from the key-value pair stored in the value field to use as the final key, and grab the value from the key-value pair to use as the final value reducedStream.selectKey((key, value)->value.getOriginalKey()) .mapValues((value)->value.getRawValue()) .to(outputTopicName, Produced.with(Serdes.String(), Serdes.String()));
修改后的解决方案
核心思路:放弃outer join的两两配对逻辑,改为直接合并两个流,再按共同key分组聚合;在聚合阶段判断是否存在流2的消息——若存在则保留所有流2消息,否则保留流1消息;同时依赖Kafka分区内的天然有序性,保证结果顺序与流1一致。
修改后的代码
// 1. 读取两个输入流 KStream<String, String> rawInputStream = builder.stream(rawTopic, Consumed.with(Serdes.String(), Serdes.String())); KStream<String, String> augmentedInputStream = builder.stream(augTopic, Consumed.with(Serdes.String(), Serdes.String())); // 2. 映射两个流到统一格式,保留共同key、原始key、值、来源标识 // 流1映射:无原始key用空字符串,记录来源为RAW,同时标记流1的处理时间用于顺序参考 KStream<String, CustomMessageDetailsWithKeyAndOrigin> mappedRaw = rawInputStream .map((key, value) -> KeyValue.pair( getCommonKeyFromRawInputStream(value), new CustomMessageDetailsWithKeyAndOrigin( getValueFromRawInputStream(value), key == null ? "" : key, OriginStream.RAW, System.currentTimeMillis() ) )); // 流2映射:保留分配的Key,记录来源为AUGMENTED KStream<String, CustomMessageDetailsWithKeyAndOrigin> mappedAug = augmentedInputStream .map((key, value) -> KeyValue.pair( getCommonKeyFromAugmentedInputStream(value), new CustomMessageDetailsWithKeyAndOrigin( value, key == null ? "" : key, OriginStream.AUGMENTED, 0 ) )); // 3. 直接合并两个流,避免join导致的消息遗漏 KStream<String, CustomMessageDetailsWithKeyAndOrigin> combinedStream = mappedRaw.concat(mappedAug); // 4. 按共同key分组,聚合时收集所有消息并按规则过滤 KGroupedStream<String, CustomMessageDetailsWithKeyAndOrigin> grouped = combinedStream.groupByKey(); KStream<String, List<CustomMessageDetailsWithKeyAndOrigin>> aggregatedStream = grouped .aggregate( ArrayList::new, // 初始化空列表存储同key的所有消息 (commonKey, msg, aggregateList) -> { aggregateList.add(msg); return aggregateList; }, // 需为List<CustomMessageDetailsWithKeyAndOrigin>实现自定义Serde,或使用Kafka提供的ListSerde Materialized.with(Serdes.String(), new ListSerde<>(CustomMessageDetailsWithKeyAndOrigin.class)) ) .toStream() .mapValues(aggregateList -> { // 过滤逻辑:优先保留所有流2消息,无流2消息时保留流1消息 List<CustomMessageDetailsWithKeyAndOrigin> augMessages = aggregateList.stream() .filter(msg -> OriginStream.AUGMENTED.equals(msg.getOrigin())) .collect(Collectors.toList()); return augMessages.isEmpty() ? aggregateList.stream().filter(msg -> OriginStream.RAW.equals(msg.getOrigin())).collect(Collectors.toList()) : augMessages; }); // 5. 恢复原始Key和值,依赖分区有序性保证输出顺序与流1一致 aggregatedStream .flatMapValues(Collection::stream) .selectKey((commonKey, msg) -> msg.getOriginalKey()) .mapValues(CustomMessageDetailsWithKeyAndOrigin::getRawValue) .to(outputTopicName, Produced.with(Serdes.String(), Serdes.String()));
关键修改点说明
- 替换Outer Join为Concat:直接合并两个流的所有消息,避免join逻辑导致流2的多副本消息被遗漏;
- 聚合阶段的过滤逻辑:针对每个共同key的消息集合,优先保留所有流2消息,仅当无流2消息时才保留流1消息,完全满足规则4;
- 顺序保证:依赖Kafka分区内的消息天然有序性,流1的消息在分区内按原始顺序处理,最终输出的流1消息会保持原顺序;流2的消息会跟随对应流1消息的位置输出;
- 自定义Serde:聚合时使用List存储消息,需要为
List<CustomMessageDetailsWithKeyAndOrigin>实现自定义Serde,或使用Kafka官方提供的Serdes.ListSerde(需确保CustomMessageDetailsWithKeyAndOrigin实现序列化接口)。
内容的提问来源于stack exchange,提问作者simonalexander2005
相关产品推荐
相关产品推荐

