项目概况

解决异步挑战:Reactor Context 实现响应式上下文传递

2026-04-077 min readAgent 工程化

最近在研究 AgentScope 的代码,从中发现 Reactor 编程中,链路追踪的用法,如果对这块只是看不懂,那么多半看不懂 AgentScope 的用法。这块的只是在这里补充下

在 Spring MVC 时代,ThreadLocal 是我们的“随身口袋”,随手存个 UserIdLogId 简直不要太爽。

但当你踏入 WebFlux 或 Project Reactor 的地盘,你会发现这里的线程就像渣男,说换就换。一个请求可能在 A 线程接客,B 线程查库,C 线程给响应。这时候,基于线程绑定的 ThreadLocal 瞬间就变成了“薛定谔的变量”——你永远不知道这一行代码跑完,下一个操作符里数据还在不在。

为了救火,Reactor 给我们准备了 Context


一、 核心痛点:为什么 ThreadLocal 必死?

我们先看一眼惨案现场:

java
1// 这是一个必挂的写法 2ThreadLocal<String> traceId = new ThreadLocal<>(); 3traceId.set("UUID-999"); 4 5Mono.just("Hello") 6 .subscribeOn(Schedulers.boundedElastic()) // 切换线程了 7 .map(data -> { 8 // 这里的 traceId.get() 必拿 null 9 log.info("Trace: {}", traceId.get()); 10 return data + " World"; 11 }) 12 .subscribe();

真相: 响应式编程的本质是异步流水线。操作符(Operator)不关心线程,它们只关心数据。线程只是被招来干活的“临时工”,干完这步就撤了,指望 ThreadLocal 跨步传递数据,无异于刻舟求剑。


二、 Reactor Context 的“逆流”设计

Reactor Context 不是存放在线程里的,它是存在 Subscriber(订阅者) 里的。

这里的逻辑非常硬核:

  1. 装配期: 你写 .contextWrite() 时,其实什么都没发生,只是在链条上挂了个钩子。
  2. 订阅期: 只有当你执行 .subscribe() 时,一个“订阅信号”会从下往上传递。
  3. 注入: 当信号路过 .contextWrite() 时,Context 数据就像被吸附在信号上一样,一路带到了流的最顶端。

重点: 这种“自底向上”的设计决定了,上游的操作符只能看到它下游定义的 Context。


三、 代码实力:读写姿势与实战

1. 基础读写:如何优雅地“带货”

java
1public Mono<String> processData() { 2 // 使用 deferContextual 延迟获取上下文 3 return Mono.deferContextual(ctx -> { 4 String correlationId = ctx.get("X-Correlation-ID"); 5 return Mono.just("Processed with ID: " + correlationId); 6 }); 7} 8 9@Test 10public void testContextPropagation() { 11 processData() 12 .map(String::toUpperCase) 13 // 写入 Context,注意它在下方,但能作用于上方 14 .contextWrite(Context.of("X-Correlation-ID", "REQ-12345")) 15 // ⚠️⚠️ 注意 下面是无法获取 “X-Correlation-ID” 的,因为Context,只能作用于上方。 16 .subscribe(System.out::println); 17} 18 19 // 如果想要下面也能拿到就必须这样写 20 @Test 21 public void testContext2() { 22 getSecureData() 23 .contextWrite(Context.of("user", "Admin")) // 写入 Context 24 .flatMap(data -> Mono.deferContextual(ctx -> { 25 // 在这里又可以掏口袋了,注意拿到的是下面定义的 Admin X,而不是上面的Admin 26 String user = ctx.getOrDefault("user", "unknown"); 27 // Secure data for: Admin processed by Admin X 28 return Mono.just(data + " processed by " + user); 29 })) 30 31 .contextWrite(Context.of("user", "Admin X")) 32 .subscribe(System.out::println); 33 }

2. 结合 Spring Security(最强实战)

在 WebFlux 中,获取当前用户不需要传参,直接从 Context 里掏:

java
1public Mono<Void> updateBusiness() { 2 return ReactiveSecurityContextHolder.getContext() 3 .map(SecurityContext::getAuthentication) 4 .flatMap(auth -> { 5 log.info("当前操作人: {}", auth.getName()); 6 return doBusiness(auth.getPrincipal()); 7 }); 8}

原理: Spring Security 的拦截器会在流的最下游悄悄执行 .contextWrite(),把认证信息塞进去,供你全局调用。


四、 高级进阶:如何跟 Log4j2 (MDC) 这种“老顽固”对接?

很多老旧库依然死磕 ThreadLocal(比如日志 MDC)。在 Reactor 里,我们可以利用 handle 或自定义 Hook 来做一个“时空转换”。

推荐方案:使用 Micrometer Context Propagation 库 这是目前工业级的标准做法:

java
1 2// 1. 初始化(全局执行一次) 3Hooks.enableAutomaticContextPropagation(); 4ContextRegistry.getInstance() 5 .registerThreadLocalAccessor("traceId", 6 MDC::get, MDC::put, MDC::remove); 7 8// 2. 在 Reactor 链条中开启自动同步 9Mono.just("Task") 10 .handle((val, sink) -> { 11 // 这里可以直接打日志,MDC 里的 traceId 已经被自动还原了! 12 log.info("日志里现在有 traceId 了"); 13 sink.next(val); 14 }) 15 .contextWrite(Context.of("traceId", "T-800")) 16 .contextCapture(); // 关键:捕捉当前环境到 Context

五、生产环境下的应用

5.1 注册“搬运规则” (Registry)

java
1@Configuration 2public class ContextPropagationConfig { 3 @PostConstruct 4 public void init() { 5 // 注册 MDC 的搬运规则 6 ContextRegistry.getInstance() 7 .registerThreadLocalAccessor( 8 "TRACE_ID", // Context 中的 Key 9 MDC::get, // 如何从 ThreadLocal 读 10 val -> MDC.put("traceId", val), // 如何写回 ThreadLocal 11 () -> MDC.remove("traceId") // 如何清理 12 ); 13 } 14}

5.2 开启“全自动外挂” (Global Hook)

在启动类里一行代码搞定,让 Reactor 在切换线程时自动同步。

java
1public static void main(String[] args) { 2 Hooks.enableAutomaticContextPropagation(); 3 SpringApplication.run(DemoApplication.class, args); 4}

5.3 自动搬运

java
1@Component 2@Order(Ordered.HIGHEST_PRECEDENCE) // 优先级最高,越早抓取越好 3public class ContextSnapshotFilter implements WebFilter { 4 @Override 5 public Mono<Void> filter(ServerWebExchange exchange, WebFilterChain chain) { 6 // 1. 此时是在 HTTP 线程,ThreadLocal 还在 7 // 2. 这里的 contextCapture() 会把 ThreadLocal 存入 Reactor Context 8 return chain.filter(exchange) 9 .contextCapture(); 10 } 11}

5.4 注意点

很多人觉得开了 enableAutomaticContextPropagation() 就万事大吉了,其实不然: Automatic Propagation 负责的是:数据已经在 Context 里了,如何在 flatMap 切换线程时帮你自动搬运。 contextCapture() 负责的是:数据还在 ThreadLocal 里,如何把它“第一手”抓进响应式世界。


六、 使用指南

  1. 不可变性警告: ContextImmutable 的。ctx.put(k, v) 不会改变原对象,它会返回一个全新的 Context。记得像用 String 一样用它。
  2. 位置决定生死: 记住口诀——下层写,上层读。如果你把 contextWrite 放在了链条的最顶端,那么它下面的 flatMap 永远拿不到值。
  3. 别塞大对象: Context 会随着整个订阅生命周期常驻内存。别把整个数据库查询结果存进去,那不是上下文,那是 OOM 的温床。
java
1final class Context0 implements CoreContext { 2 3 static final Context0 INSTANCE = new Context0(); 4 5 @Override 6 public Context put(Object key, Object value) { 7 Objects.requireNonNull(key, "key"); 8 Objects.requireNonNull(value, "value"); 9 return new Context1(key, value); 10 } 11}

七、简单的Demo案例

7.1 线程切换导致无法获取 ThreadLocal 中的数据

  • RequestContextWebFilter 将上下文通过contextWrite写入到 Reactor Context 中
  • ContextPropagationService中只能获取 Reactor Context 信息,无法从 ThreadLocal 中获取,因为线程变了
java
1 2@Component 3public class RequestContextWebFilter implements WebFilter { 4 5 private static final String TRACE_ID_HEADER = "X-Trace-Id"; 6 private static final String USER_ID_HEADER = "X-User-Id"; 7 8 @Override 9 public Mono<Void> filter(ServerWebExchange exchange, WebFilterChain chain) { 10 ServerHttpRequest request = exchange.getRequest(); 11 String traceId = headerOrDefault(request, TRACE_ID_HEADER, "trace-" + UUID.randomUUID()); 12 String userId = headerOrDefault(request, USER_ID_HEADER, "anonymous"); 13 14 return chain.filter(exchange) 15 .contextWrite(context -> context 16 .put(ReactorContextKeys.TRACE_ID, traceId) 17 .put(ReactorContextKeys.USER_ID, userId)); 18 } 19 20 private String headerOrDefault(ServerHttpRequest request, String headerName, String defaultValue) { 21 String value = request.getHeaders().getFirst(headerName); 22 return value == null || value.isBlank() ? defaultValue : value; 23 } 24} 25 26@Service 27public class ContextPropagationService { 28 29 private static final ThreadLocal<String> TRACE_ID_HOLDER = new ThreadLocal<>(); 30 private static final ThreadLocal<String> USER_ID_HOLDER = new ThreadLocal<>(); 31 32 public Mono<ContextDemoResponse> inspectContext() { 33 return Mono.deferContextual(contextView -> { 34 ContextHopSnapshot beforeSwitch = snapshot("before-publishOn", contextView); 35 36 return Mono.just(beforeSwitch) 37 .publishOn(Schedulers.boundedElastic()) 38 .flatMap(ignored -> readAfterThreadSwitch(beforeSwitch)); 39 }); 40 } 41 42 private Mono<ContextDemoResponse> readAfterThreadSwitch(ContextHopSnapshot beforeSwitch) { 43 return Mono.deferContextual(contextView -> Mono.just(new ContextDemoResponse( 44 "Automatic context propagation is active. Reactor Context values are restored into ThreadLocal even after thread switching.", 45 beforeSwitch, 46 snapshot("after-publishOn", contextView) 47 ))); 48 } 49 50 private ContextHopSnapshot snapshot(String stage, ContextView contextView) { 51 // 是无法拿到,TRACE_ID_HOLDER和USER_ID_HOLDER的。因为切换了线程 52 return new ContextHopSnapshot( 53 stage, 54 Thread.currentThread().getName(), 55 contextView.get(ReactorContextKeys.TRACE_ID), 56 contextView.get(ReactorContextKeys.USER_ID), 57 TRACE_ID_HOLDER.get(), 58 USER_ID_HOLDER.get() 59 ); 60 } 61} 62

7.2 添加上自动映射

  • 添加上自动转换了逻辑,自动将 Reactor Context 上下文转换到 ThreadLocal 上下文中。
xml
1 <dependency> 2 <groupId>io.micrometer</groupId> 3 <artifactId>context-propagation</artifactId> 4 <version>1.1.0</version> 5 </dependency>
java
1@Service 2public class ContextPropagationService { 3 4 private static final ThreadLocal<String> TRACE_ID_HOLDER = new ThreadLocal<>(); 5 private static final ThreadLocal<String> USER_ID_HOLDER = new ThreadLocal<>(); 6 7 @PostConstruct 8 public void init() { 9 Hooks.enableAutomaticContextPropagation(); 10 ContextRegistry contextRegistry = ContextRegistry.getInstance(); 11 contextRegistry.removeThreadLocalAccessor(ReactorContextKeys.TRACE_ID); 12 contextRegistry.removeThreadLocalAccessor(ReactorContextKeys.USER_ID); 13 contextRegistry.registerThreadLocalAccessor( 14 ReactorContextKeys.USER_ID, 15 USER_ID_HOLDER::get, 16 USER_ID_HOLDER::set, 17 USER_ID_HOLDER::remove 18 ); 19 contextRegistry.registerThreadLocalAccessor( 20 ReactorContextKeys.TRACE_ID, 21 TRACE_ID_HOLDER::get, 22 TRACE_ID_HOLDER::set, 23 TRACE_ID_HOLDER::remove 24 ); 25 26 } 27} 28

八、核心原理

其实归根到底,我们是在研究,Reactor 中异步编程 Context 上下文,如何跟 ThreadLocal 进行数据交换的逻辑。 要么是 ThreadLocal 转移到 Context , 要么是 Context 转换到 ThreadLocal。

如果要我们自己来设计如何才能真正的做到无缝迁移,其实如果我们能在Reactor 的操作符中每次进行转移,其实就能无缝实现自动数据交换了。其实 Reactor 也提供了这样的能力,下面我们就研究这种方式。

8.1 思路分析

  • web 过滤器中,将用户信息,保存到 Context 中。
  • 在每个 Reactor 每个操作符中,从 Context中恢复到 ThreadLocal 中。
  • 当每个操作符执行完成后,在从 ThreadLocal 中删除。

8.2 过滤器中维护 Reactor Context

请看代码注释,写的很清楚了。RequestContextRegistry 看不懂没关系,马上就说。

java
1@Component 2public class RequestContextWebFilter implements WebFilter { 3 4 private static final String TRACE_ID_HEADER = "X-Trace-Id"; 5 private static final String USER_ID_HEADER = "X-User-Id"; 6 7 @Override 8 public Mono<Void> filter(ServerWebExchange exchange, WebFilterChain chain) { 9 ServerHttpRequest request = exchange.getRequest(); 10 // 从请求头中拿到用户信息 11 String traceId = headerOrDefault(request, TRACE_ID_HEADER, "trace-" + UUID.randomUUID()); 12 String userId = headerOrDefault(request, USER_ID_HEADER, "anonymous"); 13 // 生成模拟的上下文需要的用户对象 14 RequestContextSnapshot requestContext = new RequestContextSnapshot(traceId, userId); 15 // 通过contextWrite,添加到 Reactor Context 中。RequestContextRegistry 封装了统一的操作。属于无状态类,所以是线程安全的。 16 return chain.filter(exchange) 17 .contextWrite(context -> RequestContextRegistry.get().storeRequestContext(context, requestContext)); 18 } 19 20 private String headerOrDefault(ServerHttpRequest request, String headerName, String defaultValue) { 21 String value = request.getHeaders().getFirst(headerName); 22 return value == null || value.isBlank() ? defaultValue : value; 23 } 24}

8.3 注册统一的Reactor操作符处理

代码比较长,主要看 RequestContextRegistry#enablePropagationHook 方法,信号处理的方法: onNext/onError/onComplete中封装转换的逻辑,其他方法都透传即可,这是 api,没必要死记硬背。会用即可。

java
1public final class RequestContextRegistry { 2 private static final String HOOK_KEY = "react-ctx-request-context"; 3 4 private static volatile boolean hookEnabled = false; 5 6 private static final RequestContextTracer TRACER = new RequestContextTracer(); 7 8/** 9 * 启用 Reactor 全局 hook。 10 * 11 * <p>实现机制: 12 * <p>对于每一个 Reactor operator,Reactor 都会把原始 subscriber 提升为我们自己的包装 subscriber。 13 * 这个包装层不会改变业务逻辑,只负责在 {@code onNext}{@code onError}{@code onComplete} 14 * 被调用之前,先把当前请求上下文恢复到 {@link ThreadLocal} 中。 15 * 16 * <p>为什么这样可以工作: 17 * <p>下游 subscriber 会一直携带 Reactor 的 {@link Context},可以通过 18 * {@link CoreSubscriber#currentContext()} 取到。这个上下文能够跨异步边界传播。 19 * 因此我们只需要在 signal 回调前,从当前 Reactor Context 中取出请求数据,恢复到 20 * {@code ThreadLocal},执行原始回调,然后再恢复旧的 ThreadLocal 状态即可。 21 */ 22 public static synchronized void enablePropagationHook() { 23 if (!hookEnabled) { 24 Hooks.onEachOperator( 25 HOOK_KEY, 26 Operators.lift( 27 (scannable, subscriber) -> 28 new CoreSubscriber<Object>() { 29 /** 30 * 订阅阶段。 31 * 32 * <p>当前 demo 在这个阶段不需要恢复上下文, 33 * 直接把订阅事件透传给原始 subscriber 即可。 34 */ 35 @Override 36 public void onSubscribe(Subscription subscription) { 37 subscriber.onSubscribe(subscription); 38 } 39 40 /** 41 * 数据信号。 42 * 43 * <p>在把数据传递给下游逻辑之前,先把 Reactor {@link Context} 44 * 里的请求上下文恢复到 {@link ThreadLocal}45 * 这是整个实现的关键步骤,它保证旧的命令式代码依然可以通过 46 * {@link RequestContextThreadLocalHolder} 读取上下文。 47 */ 48 @Override 49 public void onNext(Object value) { 50 TRACER.runWithContext( 51 subscriber.currentContext(), 52 () -> { 53 subscriber.onNext(value); 54 return null; 55 }); 56 } 57 58 /** 59 * 异常信号。 60 * 61 * <p>异常处理逻辑通常也会读取 traceId、userId 这类请求上下文, 62 * 所以这里和正常数据信号一样,也需要先恢复 ThreadLocal。 63 */ 64 @Override 65 public void onError(Throwable throwable) { 66 TRACER.runWithContext( 67 subscriber.currentContext(), 68 () -> { 69 subscriber.onError(throwable); 70 return null; 71 }); 72 } 73 74 /** 75 * 完成信号。 76 * 77 * <p>完成回调里也可能访问请求级 ThreadLocal, 78 * 因此这里同样保持一致,先恢复上下文再调用下游。 79 */ 80 @Override 81 public void onComplete() { 82 TRACER.runWithContext( 83 subscriber.currentContext(), 84 () -> { 85 subscriber.onComplete(); 86 return null; 87 }); 88 } 89 90 /** 91 * 原样暴露下游 subscriber 的 Reactor Context。 92 * 93 * <p>这一点很重要:包装 subscriber 自己不维护独立的 Context, 94 * 而是始终委托给原始 subscriber。 95 * 这样上游放进去的请求上下文才能在整条链路里持续可见。 96 */ 97 @Override 98 public Context currentContext() { 99 return subscriber.currentContext(); 100 } 101 })); 102 hookEnabled = true; 103 } 104 } 105} 106

8.4 真正的转换逻辑-RequestContextTracer

主要负责两件事:

  1. 从 Reactor {@link ContextView} 中读取当前请求上下文;
  2. 在执行指定逻辑前,把请求上下文恢复到 {@link ThreadLocal},执行完成后再恢复旧值。

这里逻辑很简单,其实就是 ThreadLocal 的多线程上下文切换的原理。

java
1public class RequestContextTracer { 2 3 private static final String CONTEXT_KEY = "request-context"; 4 5 /** 6 * 从 Reactor Context 中读取请求上下文。 7 * 8 * <p>如果当前链路里没有放入请求上下文,则返回一个空对象, 9 * 而不是返回 {@code null}。这样调用方在读取时可以少做一层判空处理。 10 */ 11 public RequestContextSnapshot getRequestContextFromContextView(ContextView reactorCtx) { 12 return reactorCtx.getOrDefault(CONTEXT_KEY, RequestContextSnapshot.empty()); 13 } 14 15 // 从Reactor Context读取上下文, 16 public <TResp> TResp runWithContext(ContextView reactorCtx, Supplier<TResp> inner) { 17 RequestContextSnapshot requestContext = getRequestContextFromContextView(reactorCtx); 18 // 先记录现在的上下文。 19 RequestContextSnapshot previous = RequestContextThreadLocalHolder.get(); 20 // 将 Reactor Context 放到当前的ThreadLocal 中 21 RequestContextThreadLocalHolder.set(requestContext); 22 try { 23 // 然后执行操作 24 return inner.get(); 25 } finally { 26 // 执行完后将原来的上下文,重新恢复进去 27 RequestContextThreadLocalHolder.restore(previous); 28 } 29 } 30} 31
java
1package cn.duoduo.reactctx.context; 2public final class RequestContextThreadLocalHolder { 3 private static final ThreadLocal<RequestContextSnapshot> REQUEST_CONTEXT_HOLDER = new ThreadLocal<>(); 4 private RequestContextThreadLocalHolder() { 5 } 6 public static RequestContextSnapshot get() { 7 return REQUEST_CONTEXT_HOLDER.get(); 8 } 9 public static void set(RequestContextSnapshot requestContext) { 10 if (requestContext == null) { 11 REQUEST_CONTEXT_HOLDER.remove(); 12 return; 13 } 14 REQUEST_CONTEXT_HOLDER.set(requestContext); 15 } 16 public static void restore(RequestContextSnapshot requestContext) { 17 set(requestContext); 18 } 19 public static String getTraceId() { 20 RequestContextSnapshot requestContext = get(); 21 return requestContext == null ? null : requestContext.traceId(); 22 } 23 public static String getUserId() { 24 RequestContextSnapshot requestContext = get(); 25 return requestContext == null ? null : requestContext.userId(); 26 } 27 public static void clearAll() { 28 REQUEST_CONTEXT_HOLDER.remove(); 29 } 30} 31

总结

Reactor Context 是响应式架构里的“隐形传送门”。它虽然看起来有点反直觉(自底向上),但却是解决异步链路数据丢失的唯一正解。

作为开发者,我们要做的就是:放下对 ThreadLocal 的执念,拥抱响应式流的生命周期。 觉得有用?点个赞,咱们评论区讨论一下你被响应式坑过的瞬间。