
本文介绍一种基于 kafka streams 的混合处理方案:通过 merge 操作合并多路数据流,再结合自定义 processor 实现“跨流去重但保留单流内重复”的业务逻辑,精准满足复杂消息路由与去重需求。
本文介绍一种基于 kafka streams 的混合处理方案:通过 merge 操作合并多路数据流,再结合自定义 processor 实现“跨流去重但保留单流内重复”的业务逻辑,精准满足复杂消息路由与去重需求。
在 Kafka Streams 应用中,常见的 reduce 或 outerJoin 无法直接支持“同一键下保留第二流所有副本、同时丢弃第一流对应记录”这一非幂等性去重逻辑——因为 reduce 是两两归约,join 是成对匹配,二者均无法感知某键在第二流中出现的次数及全部实例。
正确解法是放弃纯 DSL(Domain Specific Language)方式,转而采用 KStream#merge() + 自定义 Processor 的混合编程模型:
- 先合并流:将原始流(无键、带时间戳)与增强流(有键、值变形)统一映射为相同键结构后合并;
-
再状态化处理:在自定义 Processor 中维护两个状态:
- Set<String> 记录所有在第二流(augmented)中出现过的去重键(用于识别“该键是否来自第二流”);
- List<CustomMessageDetailsWithKeyAndOrigin> 缓存当前键的所有第二流消息(支持保留多份);
-
按需输出:遍历合并后的每条记录,依据其来源(RAW / AUGMENTED)和全局键状态,执行差异化路由:
- 若为 AUGMENTED:加入缓存列表,并标记该键“已见于第二流”;
- 若为 RAW:仅当该键未在第二流中出现过时才转发;
- 最终排序保障:因原始流顺序需保留,建议在 merge 前为原始流添加单调递增序列号(如 mapValues((v, idx) -> new EnrichedValue(v, idx))),并在 Processor 中按此序号做缓冲排序(或依赖 Kafka 分区局部有序 + 后续 Flink/Spark 二次排序)。
以下是关键代码片段示例:
// 步骤1:统一键映射并合并
KStream<String, CustomMsg> mappedRaw = rawInputStream
.map((k, v) -> KeyValue.pair(getCommonKey(v),
new CustomMsg(v, "", OriginStream.RAW)));
KStream<String, CustomMsg> mappedAug = augmentedInputStream
.map((k, v) -> KeyValue.pair(getCommonKey(v),
new CustomMsg(v, k, OriginStream.AUGMENTED)));
KStream<String, CustomMsg> merged = mappedRaw.merge(mappedAug);
// 步骤2:接入自定义 Processor(需注册 StateStore)
merged.process(() -> new DedupProcessor(), "dedup-store");自定义 DedupProcessor 核心逻辑(简化版):
public class DedupProcessor implements Processor<String, CustomMsg, String, CustomMsg> {
private ProcessorContext<String, CustomMsg> context;
private KeyValueStore<String, List<CustomMsg>> augStore; // 存储第二流全量消息
private KeyValueStore<String, Boolean> seenInAugStore; // 标记键是否出现在第二流
@Override
public void init(ProcessorContext<String, CustomMsg> context) {
this.context = context;
this.augStore = context.getStateStore("aug-store");
this.seenInAugStore = context.getStateStore("seen-flag-store");
}
@Override
public void process(String key, CustomMsg value) {
if (value.origin == OriginStream.AUGMENTED) {
// 保存第二流消息(允许重复 key)
List<CustomMsg> list = augStore.get(key);
if (list == null) list = new ArrayList<>();
list.add(value);
augStore.put(key, list);
seenInAugStore.put(key, true); // 标记该键存在第二流
} else { // RAW 流
if (!Boolean.TRUE.equals(seenInAugStore.get(key))) {
// 仅当该 key 从未出现在第二流时才转发原始消息
context.forward(key, value, To.all());
}
}
}
@Override
public void punctuate(long timestamp) {
// 定期遍历 augStore,将每个 key 对应的全部第二流消息逐条 forward
try (KeyValueIterator<String, List<CustomMsg>> iter = augStore.all()) {
while (iter.hasNext()) {
KeyValue<String, List<CustomMsg>> entry = iter.next();
for (CustomMsg msg : entry.value) {
context.forward(entry.key, msg, To.all());
}
}
augStore.flush(); // 清空已处理批次(按需设计清理策略)
}
}
}⚠️ 注意事项:
- 状态存储需配置 TTL(如 TimeWindowedStore)防止无限增长;
- punctuate() 触发时机影响实时性,建议结合 WallclockTimeExtractor 或事件时间窗口控制;
- 若要求严格保序,应在 CustomMsg 中嵌入原始流的序列号或时间戳,并在 punctuate 阶段排序后输出;
- 生产环境务必启用 enable.auto.commit 和 processing.guarantee=exactly_once_v2 保障端到端精确一次语义。
该方案突破了 Kafka Streams DSL 的表达边界,以可控的状态管理换取业务逻辑的完全自主权,是处理“条件性多副本保留”类场景的稳健实践路径。



















