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

如何实现跨Kafka Stream去重并保留同流重复记录?

Kafka Streams 双流合并问题及解决方案

需求说明

我有两个Kafka Stream,合并要求如下:

  • 流1:无Key,时间戳非唯一,示例数据:1,3,5,7,9
  • 流2:为流1消息分配Key并做可逆值修改,时间戳与流1不匹配,可能存在对应原消息的多Key重复项(如a:1aug1,b:1aug2),也可能缺失原流消息、出现乱序
  • 合并规则:
    1. 仅存在于流1的消息需保留;
    2. 同时存在于两个流的消息,保留流2的Key:Value;
    3. 仅存在于流2的消息需保留;
    4. 同一消息在流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()));

关键修改点说明

  1. 替换Outer Join为Concat:直接合并两个流的所有消息,避免join逻辑导致流2的多副本消息被遗漏;
  2. 聚合阶段的过滤逻辑:针对每个共同key的消息集合,优先保留所有流2消息,仅当无流2消息时才保留流1消息,完全满足规则4;
  3. 顺序保证:依赖Kafka分区内的消息天然有序性,流1的消息在分区内按原始顺序处理,最终输出的流1消息会保持原顺序;流2的消息会跟随对应流1消息的位置输出;
  4. 自定义Serde:聚合时使用List存储消息,需要为List<CustomMessageDetailsWithKeyAndOrigin>实现自定义Serde,或使用Kafka官方提供的Serdes.ListSerde(需确保CustomMessageDetailsWithKeyAndOrigin实现序列化接口)。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 20:14:55