You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.05 01:30:01