如何将订阅者Context传递给doOnNext中调用的异步fireAndForget方法
问题原因
Reactor的Context是沿订阅链从下往上传播的,仅在同一条订阅链内生效。你在doOnNext中单独调用subscribe()执行的fireAndForget内部Mono是独立的订阅链,和外层主订阅链完全隔离,默认不会继承外层写入的Context,所以会抛出Context为空的异常。
两个doOnNext分别的错误原因:
- 第一个
doOnNext调用的无Context入参的fireAndForget,内部Mono完全没有被写入外层Context,肯定无法取到key对应的值。 - 第二个
doOnNext虽然写了Mono.deferContextual,但你对这个Mono.deferContextual本身单独调用了subscribe(),它作为独立订阅链没有关联外层Context,所以deferContextual本身就取不到Context,直接抛出异常。
修复方案
如果要保留fireAndForget的异步非阻塞特性,同时正确传递外层Context,修改逻辑如下:
- 先从外层主订阅链拿到Context
- 把拿到的Context手动写入内部独立的fireAndForget Mono
- 独立subscribe时添加错误回调,避免未处理异常打印报错日志
完整修复后的代码
@Test void shouldPassContextToFireAndForget() { final Mono<String> helloWorldMono = Mono.just("hello") // 修复第一个doOnNext .doOnNext(name -> Mono.deferContextual(ctx -> { // 拿到外层Context,写入内部Mono再subscribe return Mono.just(name) .flatMap(value -> Mono.deferContextual(contextView -> Mono.just(value + contextView.get("key")))) .contextWrite(ctx); }).subscribe(null, err -> { // 自行处理异步异常,避免打印ErrorCallbackNotImplemented System.out.println("异步操作异常:" + err.getMessage()); })) // 修复第二个doOnNext .doOnNext(name -> Mono.deferContextual(contextView -> { // 直接用拿到的context调用方法,再subscribe返回的Mono即可 fireAndForget(contextView, name).subscribe(null, err -> { System.out.println("异步操作异常:" + err.getMessage()); }); return Mono.empty(); }).block()) // block是为了让deferContextual在主链上下文执行,拿到外层Context .flatMap(name -> Mono.deferContextual(contextView -> Mono.just(name + " " + contextView.get("key")))) .contextWrite(Context.of("key", "world")); StepVerifier.create(helloWorldMono) .expectNext("hello world") .verifyComplete(); } private Mono<String> fireAndForget(ContextView context, String name) { return Mono.just(name) .flatMap(value -> Mono.deferContextual(contextView -> Mono.just(value + contextView.get("key")))) .contextWrite(context); }
更规范的优化建议
如果不要求严格的fireAndForget(即可以接受主链等待异步操作触发后再继续),更推荐用flatMap替换doOnNext+手动subscribe的写法,自动继承主链Context,不需要手动传递:
.flatMap(name -> Mono.deferContextual(contextView -> fireAndForget(contextView, name) .subscribeOn(Schedulers.parallel()) // 异步调度 .then(Mono.just(name)) // 不影响主链输出 ))
内容的提问来源于stack exchange,提问作者Adam
相关产品推荐
相关产品推荐

