Project Reactor subscribeOn()未按文档指定线程执行onSubscribe()问题
根据Project Reactor的JavaDoc,subscribeOn(Scheduler)应该将subscribe、onSubscribe和request操作切换到指定Scheduler的线程,但你的测试结果里,log输出的onSubscribe和request却在Test worker线程,这是因为混淆了上游订阅回调和下游订阅回调的执行线程。
核心原因:订阅流程的方向与线程切换时机
Reactor的订阅是从下游到上游触发的:
- 当StepVerifier在Test worker线程调用subscribe时,首先触发下游
log()操作的onSubscribe回调——这一步是同步执行的,所以线程是Test worker。 - 随后订阅请求传递到
subscribeOn操作,subscribeOn会把订阅上游Flux.just的逻辑提交到指定的parallel线程池。 - 在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

