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

如何使用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的失败回调可以更优雅地处理发射失败的情况,避免自旋

关键注意事项

  1. 上下文传递:两种方案都自动处理了Reactor的上下文传递,无需手动干预,确保任务执行时的上下文和订阅时一致
  2. 背压处理:concatMap和任务链的then操作都会自动处理背压,不会因为任务过多导致内存溢出
  3. 资源清理:如果CriticalResource是可销毁的,记得在销毁时调用Disposable.dispose()来停止队列消费,避免内存泄漏

内容的提问来源于stack exchange,提问作者schananas

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.06 06:55:40