
本文介绍一种基于 kafka streams 的混合处理方案:通过 merge 操作合并多路流,再结合自定义 processor 实现“跨流去重但保留单流内重复”的复杂业务逻辑,解决 outerjoin + aggregate 无法满足的多副本优先级保留需求。
本文介绍一种基于 kafka streams 的混合处理方案:通过 merge 操作合并多路流,再结合自定义 processor 实现“跨流去重但保留单流内重复”的复杂业务逻辑,解决 outerjoin + aggregate 无法满足的多副本优先级保留需求。
在 Kafka Streams 中,标准的 outerJoin 和 aggregate 操作适用于一对一或一对多的关联与聚合场景,但当业务要求区分来源流、保留第二流中所有同值不同键的记录(即“同一原始值在 aug stream 中多次出现,需全部保留”),同时丢弃第一流中对应的所有副本时,传统窗口连接+集合去重的方式会失效——因为 aggregate(LinkedHashSet::new, ...) 会无差别地将所有记录视为等价元素,导致 a:1aug1 和 b:1aug2 被当作重复而仅保留其一。
正确的解法是放弃纯声明式流操作,转而采用 merge() + 自定义 Processor 的组合模式,以获得对每条记录来源、内容、键和顺序的完全控制权。
✅ 核心思路:Merge 后状态化逐条决策
- 统一键映射后 merge:将两路流(raw 和 augmented)各自映射为 (commonKey, CustomRecord) 形式,并调用 mergedStream = rawMapped.merge(augMapped);
-
接入自定义 Processor:使用 process() 添加一个有状态的 Processor,内部维护两个关键结构:
- Map<String, List<CustomRecord>> pendingByCommonKey:暂存按 commonKey 分组的记录(含来源标记);
- Set<String> seenInAugStream:记录所有已在 augmented stream 中出现过的原始值(如 "1"),用于后续过滤 raw stream 的冲突项;
-
按 commonKey 批量提交:当某 commonKey 的所有记录(来自任意流)都到达(可通过时间戳或水印触发 flush),执行业务规则:
- 若该 key 仅在 raw stream 出现 → 输出 raw 记录;
- 若该 key 仅在 aug stream 出现 → 输出全部 aug 记录;
- 若该 key 在两流均出现 → 仅输出 aug stream 的全部记录,彻底忽略 raw stream 对应记录(满足条件 4);
- 保序关键:由于原始 raw stream 的顺序需作为最终输出顺序依据,可在 CustomRecord 中额外携带原始 offset 或序列号,并在 Processor 内部对输出结果按此字段排序(或借助 KTable + suppress() 实现事件时间有序输出)。
? 示例 Processor 片段(简化版)
public class DeduplicateAcrossStreamsProcessor
implements Processor<String, CustomRecord, String, CustomRecord> {
private ProcessorContext<String, CustomRecord> context;
private KeyValueStore<String, List<CustomRecord>> store;
@Override
public void init(ProcessorContext<String, CustomRecord> context) {
this.context = context;
this.store = context.getStateStore("dedup-store");
}
@Override
public void process(String commonKey, CustomRecord record) {
// 按 commonKey 聚合记录
List<CustomRecord> list = store.get(commonKey);
if (list == null) list = new ArrayList<>();
list.add(record);
store.put(commonKey, list);
// 可选:定时/基于 watermark 触发 flush(此处省略)
// 或使用 punctuate() 周期性检查并输出
}
@Override
public void punctuate(long timestamp, Punctuator punctuator) {
// 遍历 store,对每个 commonKey 应用业务规则
try (KeyValueIterator<String, List<CustomRecord>> iter = store.all()) {
while (iter.hasNext()) {
KeyValue<String, List<CustomRecord>> entry = iter.next();
List<CustomRecord> records = entry.value;
List<CustomRecord> augOnly = records.stream()
.filter(r -> r.origin == OriginStream.AUGMENTED)
.collect(Collectors.toList());
if (!augOnly.isEmpty()) {
// 条件 2 & 4:只要存在 aug 记录,就全部输出,且不输出任何 raw
augOnly.forEach(r -> context.forward(r.key(), r));
} else {
// 条件 1 & 3:仅 raw 或仅 aug(但此处 augOnly 为空,故只剩 raw)
records.stream()
.filter(r -> r.origin == OriginStream.RAW)
.forEach(r -> context.forward(r.key(), r));
}
store.delete(entry.key); // 清理已处理 key
}
}
}
}⚠️ 注意事项与最佳实践
- 状态存储必须启用:store 需在 Topology 中显式定义为 Stores.persistentKeyValueStore(...),并绑定到 Processor;
- 容错与恢复:确保 CustomRecord 可序列化,且 Processor 的 punctuate() 逻辑幂等(如使用 store.delete() 配合 context.commit());
- 顺序保证:若 raw stream 的物理顺序至关重要,建议在 CustomRecord 中嵌入 rawOffset 字段,并在 punctuate() 输出前按该字段排序;
- 性能考量:避免在 process() 中做耗时操作;高频 key 可考虑分片或 TTL 策略防止状态无限增长;
- 替代方案提示:对于低吞吐、强一致性场景,也可先将两流分别写入 KTable,再用 transform() 查询 aug 表是否存在对应 key,但该方式延迟更高且无法天然保留 aug 流内多键副本。
综上,当 Kafka Streams 的 DSL 层无法表达“按来源流差异化去重”这类复杂语义时,转向 Processor API 并辅以状态存储,是最直接、可控且可验证的工程解法。



















