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

Project Reactor subscribeOn()未按文档指定线程执行onSubscribe()问题

为什么subscribeOn的onSubscribe和request运行在调用线程而非指定Scheduler线程?

根据Project Reactor的JavaDoc,subscribeOn(Scheduler)应该将subscribe、onSubscribe和request操作切换到指定Scheduler的线程,但你的测试结果里,log输出的onSubscribe和request却在Test worker线程,这是因为混淆了上游订阅回调和下游订阅回调的执行线程。

核心原因:订阅流程的方向与线程切换时机

Reactor的订阅是从下游到上游触发的:

  1. 当StepVerifier在Test worker线程调用subscribe时,首先触发下游log()操作的onSubscribe回调——这一步是同步执行的,所以线程是Test worker。
  2. 随后订阅请求传递到subscribeOn操作,subscribeOn会把订阅上游Flux.just的逻辑提交到指定的parallel线程池。
  3. 在parallel线程中,完成上游Flux.just的订阅,触发上游的onSubscribe和request,然后Flux.just在parallel线程中发送onNext/onComplete信号,这些信号再传递回下游的log(),所以你看到onNext/onComplete在parallel线程。

JavaDoc中描述的"在指定Scheduler线程执行subscribe、onSubscribe和request",指的是上游的订阅相关操作,而非下游的回调。你当前的log()放在subscribeOn之后,捕捉的是下游的回调,自然会显示调用线程。

验证上游订阅线程的正确方式

如果你想验证subscribeOn确实切换了上游订阅的线程,可以调整log()的位置到subscribeOn之前,或者使用doOnSubscribe捕捉上游订阅事件:

方式1:将log放在subscribeOn之前

@Test
void testSubscribeOn() {
    var flux = Flux.just("alex", "adam", "andrew")
            .log("upstream") // 捕捉上游的订阅事件
            .subscribeOn(Schedulers.parallel())
            .log("downstream"); // 捕捉下游的订阅事件

    StepVerifier.create(flux)
            .expectNextCount(3)
            .verifyComplete();
}

此时upstream的onSubscribe和request会输出在parallel线程,而downstream的对应事件仍在Test worker线程。

方式2:使用doOnSubscribe观察上游订阅

@Test
void testSubscribeOn() {
    var flux = Flux.just("alex", "adam", "andrew")
            .doOnSubscribe(sub -> System.out.printf("上游订阅线程: %s%n", Thread.currentThread().getName()))
            .subscribeOn(Schedulers.parallel())
            .log();

    StepVerifier.create(flux)
            .expectNextCount(3)
            .verifyComplete();
}

运行后会看到上游订阅线程: parallel-1的输出,证明上游订阅逻辑确实在指定Scheduler线程执行。

总结

subscribeOn的作用是将上游的订阅启动逻辑切换到指定线程,下游的初始订阅回调(onSubscribe、request)仍由发起订阅的线程执行。只有上游产生的事件(onNext/onError/onComplete)以及上游的订阅相关操作,才会运行在指定Scheduler线程中。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 08:47:39