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

如何将订阅者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,修改逻辑如下:

  1. 先从外层主订阅链拿到Context
  2. 把拿到的Context手动写入内部独立的fireAndForget Mono
  3. 独立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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 04:15:04