如何使用Project Reactor实现共享关键资源的反应式访问编排?
如何以反应式方式编排对单一关键资源的串行访问?
假设我们有一个每次只能执行一项操作的虚拟关键组件(比如文件存储、高成本远程服务等),当存在多个访问入口时,需要以反应式方式编排访问:资源空闲时立即执行操作,若已有操作正在执行,则将当前操作加入队列,待前序操作完成后完成对应的Mono,且全程无阻塞。
我最初的实现思路是将任务加入Flux队列逐个执行,并返回一个会在任务完成时结束的Mono,初步代码如下:
class CriticalResource { private final Sinks.Many<Mono<?>> taskExecutor = Sinks.many() .unicast() .onBackpressureBuffer(); private final Disposable taskExecutorDisposable = taskExecutor.asFlux() .concatMap(Function.identity()) // 按顺序执行操作 .subscribe(); public Mono<Void> resourceOperation1() { return doSomething() .as(this::sequential); } public Mono<Void> resourceOperation2() { return doSomethingElse() .as(this::sequential); } public Mono<Void> resourceOperation3() { return doSomething() .then(somethingElse()) .as(this::sequential); } private <T> Mono<T> sequential(Mono<T> action) { return Mono.defer(() -> { Sinks.One<T> actionResult = Sinks.one(); // 创建一个在action完成时结束的Mono,通过taskExecutor订阅action并传递结果 while (taskExecutor.tryEmitNext(action.doOnError(t -> actionResult.emitError(t, Sinks.EmitFailureHandler.FAIL_FAST)) .doOnSuccess(next -> actionResult.emitValue(next, Sinks.EmitFailureHandler.FAIL_FAST))) != Sinks.EmitResult.OK) { } return actionResult.asMono(); }); } }
不过这只是一个初步实现,还需要正确处理背压传递、上下文转移等问题。想请教一下,有没有基于Project Reactor官方支持的更优实现方式?
更优的官方推荐实现方案
你的核心思路(串行化任务队列)是完全正确的,但手动使用Sinks来实现确实容易在背压、上下文传递和错误处理上出问题。Project Reactor官方有两种更简洁且健壮的方案:
方案1:基于AtomicReference的轻量级任务链
这种方式通过原子引用维护当前正在执行的任务链,每次新任务都追加到链的末尾,保证串行执行,同时自动处理上下文和背压:
class CriticalResource { // 用原子引用维护当前的任务链,初始为空任务 private final AtomicReference<Mono<Void>> currentTaskChain = new AtomicReference<>(Mono.empty()); public Mono<Void> resourceOperation1() { return submit(Mono.fromRunnable(this::doSomething)); } public Mono<Void> resourceOperation2() { return submit(Mono.fromRunnable(this::doSomethingElse)); } public Mono<Void> resourceOperation3() { return submit(Mono.fromRunnable(this::doSomething).then(Mono.fromRunnable(this::somethingElse))); } private <T> Mono<T> submit(Mono<T> task) { return Mono.defer(() -> { // 缓存任务结果,避免重复执行 Mono<T> cachedTask = task.cache(); // 循环CAS更新任务链,直到成功 for (; ; ) { Mono<Void> currentChain = currentTaskChain.get(); // 将新任务追加到当前链的末尾 Mono<Void> newChain = currentChain.then(cachedTask.then()); if (currentTaskChain.compareAndSet(currentChain, newChain)) { // 返回缓存的任务,确保调用者拿到的是对应任务的结果 return cachedTask; } } }); } // 示例业务方法 private void doSomething() { /* ... */ } private void doSomethingElse() { /* ... */ } private void somethingElse() { /* ... */ } }
优势:
- 轻量级,基于CAS操作,高并发下性能优异
- 自动传递上下文:
Mono.defer会捕获当前订阅上下文,then操作会保留上下文传递 - 天然处理背压:任务链会自动排队,不会丢失任务
- 错误处理更简洁:任务的错误会自动传递到对应的
Mono,无需手动处理Sinks的错误发射
方案2:基于Sinks + concatMap的队列式实现
如果更倾向于显式的队列模型,可以简化你最初的实现,避免手动自旋发射Sinks:
class CriticalResource { private final Sinks.Many<Mono<?>> taskQueue = Sinks.many().unicast().onBackpressureBuffer(); private final Disposable queueDisposable; public CriticalResource() { // 启动队列的消费逻辑,按顺序执行任务 this.queueDisposable = taskQueue.asFlux() .concatMap(Function.identity()) .subscribe( unused -> {}, error -> { /* 处理全局队列错误(可选) */ } ); } public Mono<Void> resourceOperation1() { return enqueue(Mono.fromRunnable(this::doSomething)); } public Mono<Void> resourceOperation2() { return enqueue(Mono.fromRunnable(this::doSomethingElse)); } public Mono<Void> resourceOperation3() { return enqueue(Mono.fromRunnable(this::doSomething).then(Mono.fromRunnable(this::somethingElse))); } private <T> Mono<T> enqueue(Mono<T> task) { return Mono.create(sink -> { // 使用emitNext的异步方式,避免自旋 taskQueue.emitNext( task.doOnSuccess(sink::success) .doOnError(sink::error) .doOnTerminate(sink::onComplete), (signalType, emitResult) -> { // 处理发射失败的情况,比如队列已满(这里用重试策略) return true; } ); }); } // 示例业务方法 private void doSomething() { /* ... */ } private void doSomethingElse() { /* ... */ } private void somethingElse() { /* ... */ } }
优势:
- 队列模型更直观,适合需要监控任务队列状态的场景
concatMap是Reactor官方提供的串行化操作符,背压和错误处理都经过充分验证- 使用
emitNext的失败回调可以更优雅地处理发射失败的情况,避免自旋
关键注意事项
- 上下文传递:两种方案都自动处理了Reactor的上下文传递,无需手动干预,确保任务执行时的上下文和订阅时一致
- 背压处理:
concatMap和任务链的then操作都会自动处理背压,不会因为任务过多导致内存溢出 - 资源清理:如果
CriticalResource是可销毁的,记得在销毁时调用Disposable.dispose()来停止队列消费,避免内存泄漏
内容的提问来源于stack exchange,提问作者schananas
相关产品推荐
相关产品推荐

