Reactor多源入队集中处理并等待完成的无阻塞优化方案咨询
Reactor 多数据源中心化处理队列实现咨询
核心需求
源端将条目加入队列,待条目处理完成后接收回调通知
请问在Reactor中是否可以实现如下逻辑:从一个流水线将条目加入队列,由另一个线程/Flux从队列拉取条目处理,且原流水线需等待该条目被另一线程处理完成后再继续执行?
我需要类似的实现,因为业务中存在多个不同的数据源,我需要一套中心化的流水线统一处理背压、重试等通用逻辑。
目前我有一套可运行的实现,但需要占用大量线程且用到了Semaphores,且不同流水线之间会丢失上下文,后续无法适配响应式事务这类特性。
当前简化版实现代码
public class Source { private final Queue queue = new Queue(); public Flux<String> produce() { return Flux.range(0, 10) .doOnNext(v -> System.out.println("Before queue processed: " + v)) .flatMap(this.queue::enqueue) .doOnNext(v -> System.out.println("After queue processed: " + v)); } } public class Queue { private class WorkEntry { int number; String word; Throwable exception; final Semaphore semaphore = new Semaphore(0); } private final LinkedBlockingQueue<WorkEntry> blockingQueue = new LinkedBlockingQueue<>(); public Mono<String> enqueue(int number) { return Mono.just(number) .flatMap(n -> { final var entry = new WorkEntry(); entry.number = number; this.blockingQueue.add(entry); return Mono.just(entry); }) .subscribeOn(Schedulers.boundedElastic()) .flatMap(entry -> { try { // 阻塞直到drain()处理完当前条目 entry.semaphore.acquire(); } catch (InterruptedException e) { e.printStackTrace(); } // 'word'字段已在drain()处理过程中赋值 return Mono.just(entry.word); }); } public void drain() { Flux .<WorkEntry>generate(sink -> { final var entry = blockingQueue.poll(); if (entry == null) { sink.complete(); } else { sink.next(entry); } }) .flatMap(entry -> Mono.just(entry) .flatMap(e -> Mono.just(e.number + "!")) .onErrorResume(ex -> { entry.exception = ex; return Mono.empty(); }) .doOnSuccess(word -> { entry.word = word; entry.semaphore.release(1); })) .subscribeOn(Schedulers.boundedElastic()) .subscribe(); } } public class Runner { public void run() { final var source = new Source(); source.produce() .subscribe(); source.queue.drain(); } }
运行输出结果
Before queue processed: 0 ... Before queue processed: 9 After queue processed: 0! ... After queue processed: 9!
待解决的优化问题
- 如何优化这套实现?
- 如何避免在
.flatMap中使用阻塞调用? - 如何避免启动过多的elastic线程?
- 如何保证全链路上下文一致?
我曾考虑过重写为单个Flux从多个数据源拉取数据的模式,但这种方案的问题是会提升测试复杂度,很难针对单个数据源推送条目后验证该数据源的处理完成状态。
更新优化方案
我已经找到了一种无需Semaphore锁的更优实现,虽然不算完美但已有很大提升,核心是通过delayUntil和Sinks实现数据源与队列之间的通信,实现代码如下:
... public class Queue { private class WorkEntry { int number; final Sink.One<String> sink; } private final LinkedBlockingQueue<WorkEntry> blockingQueue = new LinkedBlockingQueue<>(); public Mono<String> enqueue(int number) { return Mono.just(number) .flatMap(n -> { final var entry = new WorkEntry(); entry.number = number; entry.sink = Sink.one(); this.blockingQueue.add(entry); return Mono.just(entry); }) .subscribeOn(Schedulers.boundedElastic()) .delayUntil(e -> e.sink.asMono()) .flatMap(e -> e.sink.asMono()); } public void drain() { Flux .<WorkEntry>generate(sink -> { final var entry = blockingQueue.poll(); if (entry == null) { sink.complete(); } else { sink.next(entry); } }) .flatMap(entry -> Mono.just(entry) .flatMap(e -> Mono.just(e.number + "!")) .doOnError(t -> entry.sink.tryEmitError(t)) .doOnSuccess(word -> entry.sink.tryEmitValue(word)) ) .subscribeOn(Schedulers.boundedElastic()) .subscribe(); } } ...
该方案可行的原因是Sink.one()的.asMono()方法每次调用都会返回同一个实例,因此我们既可以用它做延迟触发,也可以在.flatMap()中直接返回处理结果。
内容的提问来源于stack exchange,提问作者Stmated
相关产品推荐
相关产品推荐

