
本文介绍如何在 Reactor 中对无限热流(如 Kafka 消息流)进行低内存开销、高吞吐的重复值识别与转发——不依赖 cache() 或 groupBy,而是通过双 Set 状态机实现 O(1) 判重与增量输出。
本文介绍如何在 reactor 中对无限热流(如 kafka 消息流)进行低内存开销、高吞吐的重复值识别与转发——不依赖 `cache()` 或 `groupby`,而是通过双 set 状态机实现 o(1) 判重与增量输出。
在响应式编程中,对无限热流(如 Kafka 主题)做“仅保留重复项(含首次重复及后续所有出现)”的操作,是一个典型但易踩坑的场景。原始方案使用 groupBy(identity).flatMap(...cache().buffer(2)...) 虽逻辑可行,却存在三大硬伤:
- ✅ 内存不可控:cache() 会为每个分组缓存全部历史数据,分组数越多、流越长,OOM 风险越高;
- ❌ 热源失效:groupBy 内部依赖冷流语义,对 publish() 后的热流可能丢失事件或触发竞态;
- ⚠️ 延迟不可预测:buffer(2) + skip(1) 引入隐式缓冲与序列依赖,违背背压友好原则。
更优解是将去重逻辑下沉至订阅层,用轻量状态机替代流式分组:
核心思想:双状态集(Unique + Duplicate)
我们维护两个线程安全的 ConcurrentHashMap(或 Collections.synchronizedSet(new HashSet<>())):
- uniques:记录首次出现且尚未被判定为重复的值;
- dupes:记录已确认为重复(至少出现两次) 的值。
每收到一个新元素 n,按以下状态迁移处理:
- 若 n ∈ uniques → 第二次出现 → 移出 uniques,加入 dupes,立即向下游发射 n;
- 若 n ∈ dupes → 后续重复 → 直接发射 n;
- 若 n ∉ uniques ∪ dupes → 首次出现 → 加入 uniques,不发射。
该逻辑时间复杂度 O(1),空间复杂度 O(U),其中 U 是当前活跃唯一值数量(远小于总历史量),天然适配无限流。
实现示例(生产就绪版)
import reactor.core.publisher.Flux;
import java.util.Collections;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
public class DedupeStream {
// 使用 ConcurrentHashMap 保证线程安全(Kafka Flux 多线程推送常见)
private static final Set<Integer> uniques = Collections.newSetFromMap(new ConcurrentHashMap<>());
private static final Set<Integer> dupes = Collections.newSetFromMap(new ConcurrentHashMap<>());
public static Flux<Integer> filterToDuplicates(Flux<Integer> source) {
return Flux.from(source)
.doOnNext(n -> {
if (uniques.contains(n)) {
uniques.remove(n);
dupes.add(n);
// 立即转发首次重复项
downstream.accept(n);
} else if (dupes.contains(n)) {
// 转发后续重复项
downstream.accept(n);
} else {
// 首次出现,暂存
uniques.add(n);
}
})
.thenMany(Flux.empty()); // 忽略原始流数据,只走副作用路径
}
// 示例下游消费者(可替换为 KafkaProducer、DB Sink 等)
private static final java.util.function.Consumer<Integer> downstream = n -> {
System.out.println("→ Duplicate emitted: " + n);
};
// 测试入口
public static void main(String[] args) {
Flux<Integer> test = Flux.just(1, 2, 3, 4, 1, 1, 1, 2, 4);
filterToDuplicates(test).blockLast(); // 输出: 1,1,1,2,4
}
}关键注意事项
- 线程安全必须显式保障:Flux 订阅者可能被多个线程调用(尤其 Kafka binder 默认启用多线程拉取),务必使用 ConcurrentHashMap 或同步包装器;
- 避免阻塞下游影响上游:若 downstream 是慢操作(如 HTTP 调用),应使用 publishOn(scheduler) 解耦,防止背压中断;
- 状态清理策略:若业务允许“窗口去重”,可定期清空 uniques/dupes(如基于时间/计数);长期运行需监控 Set 大小,防内存泄漏;
- 扩展性提示:单机 Set 不适用于分布式场景,此时应替换为 Redis Set(SISMEMBER/SADD)或布隆过滤器 + 外部存储。
此方案剥离了流式算子的隐式开销,以明确的状态管理换取极致性能与可控性,是处理实时重复检测类需求的推荐实践。

















