如何用Project Reactor计算两个有序冷Flux的集合差并处理背压?
实现支持背压的有序Flux集合差计算
要实现两个冷发布者有序Flux的集合差(B - A,即B中存在但A中不存在的元素),同时避免内存耗尽,核心是利用Reactor的背压机制控制B元素的暂存缓冲区大小,当缓冲区满时暂停B的生产。以下是具体实现方案:
核心思路
- 标记元素来源:给A、B Flux的元素打上来源标记,方便后续区分处理。
- 有界暂存缓冲区:用有界阻塞队列暂存B中未在A中出现过的元素,队列满时阻塞B的生产,触发背压。
- 流式跟踪A元素:用线程安全的集合记录A中已处理的元素,实时过滤B的元素。
- 协调流的生命周期:当A流结束时,输出缓冲区中剩余的B元素(即集合差结果),同时处理取消、错误等边界情况。
代码实现
1. 元素标记类
首先定义一个简单的标记类,区分元素来自A还是B:
import reactor.core.publisher.Flux; import reactor.core.publisher.BaseSubscriber; import java.util.concurrent.BlockingQueue; import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.ConcurrentHashMap; import java.util.Set; static class TaggedElement<T> { enum Source { A, B } final T value; final Source source; TaggedElement(T value, Source source) { this.value = value; this.source = source; } }
2. 集合差计算方法
实现核心逻辑,通过Flux.push手动控制流的生产,结合BaseSubscriber处理B流的背压:
public static <T> Flux<T> orderedSetDifference(Flux<T> fluxA, Flux<T> fluxB, int bufferCapacity) { return Flux.push(sink -> { // 有界队列暂存B元素,满时触发背压 BlockingQueue<T> bBuffer = new LinkedBlockingQueue<>(bufferCapacity); // 线程安全集合记录A中已出现的元素 Set<T> seenInA = ConcurrentHashMap.newKeySet(); // 订阅A流:记录元素并清理缓冲区中匹配的B元素 Disposable aDisposable = fluxA.subscribe( a -> { seenInA.add(a); // 移除缓冲区中所有与当前A元素匹配的B元素 bBuffer.remove(a); }, sink::error, () -> { // A流结束,输出缓冲区剩余的B元素 while (!bBuffer.isEmpty()) { sink.next(bBuffer.poll()); } sink.complete(); } ); // 订阅B流:手动处理背压,仅当缓冲区有空间时请求下一个元素 BaseSubscriber<T> bSubscriber = new BaseSubscriber<T>() { @Override protected void hookOnSubscribe(Subscription subscription) { // 初始请求1个元素 request(1); } @Override protected void hookOnNext(T b) { try { if (seenInA.contains(b)) { // 元素已在A中出现,直接丢弃,继续请求下一个 request(1); return; } // 尝试将元素加入缓冲区,队列满时阻塞,触发背压 bBuffer.put(b); // 加入成功,请求下一个元素 request(1); } catch (InterruptedException e) { Thread.currentThread().interrupt(); cancel(); sink.error(e); } } @Override protected void hookOnError(Throwable throwable) { sink.error(throwable); } @Override protected void hookOnComplete() { // B流结束,等待A流处理完剩余逻辑 } }; Disposable bDisposable = fluxB.subscribe(bSubscriber); // 当结果流被取消时,及时取消A、B流的订阅,避免资源泄漏 sink.onCancel(() -> { aDisposable.dispose(); bDisposable.dispose(); }); }); }
关键细节说明
- 背压控制:通过
LinkedBlockingQueue的put方法实现阻塞,当缓冲区满时,B流的BaseSubscriber会暂停请求下一个元素,B的生产者会因背压机制停止发送,避免内存耗尽。 - 冷发布者适配:每次订阅结果流时,A、B流都会重新订阅,符合冷发布者的特性。
- 线程安全:使用
ConcurrentHashMap.newKeySet()和LinkedBlockingQueue保证多线程环境下的安全操作。 - 结果顺序:输出的差集元素保持B流原有的顺序,仅移除在A中出现过的元素。
使用示例
public static void main(String[] args) { Flux<Integer> fluxA = Flux.just(2, 4, 6); Flux<Integer> fluxB = Flux.just(1, 2, 3, 4, 5, 6, 7); // 计算B - A,缓冲区容量设为3 orderedSetDifference(fluxA, fluxB, 3) .subscribe(System.out::println); // 输出结果:1, 3, 5, 7 }
内容的提问来源于stack exchange,提问作者Theozaurus
相关产品推荐
相关产品推荐

