线程切换后Project Reactor上下文丢失问题求助
问题分析与解决方案
这不是Reactor Hooks.enableAutomaticContextPropagation()的Bug,而是你的实现没有适配非请求驱动的无限流场景——Spring Cloud GCP Trace默认是为HTTP请求/响应链路设计的,依赖请求上下文自动初始化Trace传播,但在无请求绑定的消息流中,需要手动适配Reactor的上下文传播机制。
核心问题原因
- 缺少Trace上下文的传播载体:Reactor的自动上下文传播需要明确注册
ThreadLocalAccessor,用来识别并跨线程传递Trace Span这类自定义上下文数据,你当前没有注册对应Accessor,导致线程切换后上下文丢失。 - Span上下文绑定方式错误:仅用
contextWrite写入Span到Reactor上下文,但未同步绑定到当前线程的ThreadLocal,导致后续算子切换线程时无法正确恢复Span上下文。 - Span生命周期未闭环:没有在消息处理的成功/错误分支统一结束Span,且业务代码直接侵入Trace逻辑,不符合你“后台自动处理”的需求。
具体解决方案
1. 注册Trace Span的ThreadLocalAccessor
创建自定义ThreadLocalAccessor,让Reactor能识别并跨线程传递Span上下文:
public class TraceSpanThreadLocalAccessor implements ThreadLocalAccessor { public static final String KEY = "trace-span"; @Override public Object key() { return KEY; } @Override public void restore(ContextView context) { Span span = context.getOrDefault(KEY, null); if (span != null) { Tracer.currentTracer().withSpan(span); // 同步到MDC,让日志自动携带traceId/spanId MDC.put("traceId", span.getTraceId()); MDC.put("spanId", span.getSpanId()); } } @Override public void capture(Context.Builder contextBuilder) { Span currentSpan = Tracer.currentTracer().currentSpan(); if (currentSpan != null) { contextBuilder.put(KEY, currentSpan); } } @Override public void reset() { Tracer.currentTracer().withSpan(null); MDC.remove("traceId"); MDC.remove("spanId"); } }
在应用启动时注册这个Accessor:
static { Hooks.enableAutomaticContextPropagation(); ContextPropagation.registerThreadLocalAccessor(new TraceSpanThreadLocalAccessor()); }
2. 正确初始化消息Span并绑定上下文
修改Handler.java中消息处理逻辑,确保Span先绑定到当前线程,再写入Reactor上下文,同时用doFinally保证Span生命周期闭环:
.flatMap(message -> { // 基于消息自带的traceId创建子Span TraceContext parentContext = TraceContext.newBuilder() .setTraceId(message.getTraceId()) .build(); Span processingSpan = Tracer.currentTracer() .spanBuilder("message-processing") .setParent(parentContext) .startSpan(); // 用Scope绑定Span到当前线程,同时写入Reactor上下文 try (Scope scope = Tracer.currentTracer().withSpan(processingSpan)) { return processMessage(message) .contextWrite(ctx -> ctx.put(TraceSpanThreadLocalAccessor.KEY, processingSpan)) .doFinally(signalType -> processingSpan.end()); // 无论成功/错误都结束Span } })
3. 让WebClient自动传播Trace上下文
给WebClient添加过滤器,从Reactor上下文中提取Span并注入请求头:
public ExchangeFilterFunction tracePropagationFilter() { return (request, next) -> Mono.deferContextual(ctx -> { Span span = ctx.getOrDefault(TraceSpanThreadLocalAccessor.KEY, null); if (span != null) { HttpRequest modifiedRequest = request.mutate() .header("X-Trace-Id", span.getTraceId()) .header("X-Span-Id", span.getSpanId()) .build(); return next.exchange(modifiedRequest); } return next.exchange(request); }); }
构建WebClient时添加该过滤器:
WebClient webClient = WebClient.builder() .filter(tracePropagationFilter()) .build();
4. 封装通用Trace操作符,隔离业务与Trace逻辑
为了避免业务代码侵入Trace逻辑,封装一个通用的Reactor操作符,自动处理Span的创建、传播和结束:
public <T> Function<Mono<T>, Mono<T>> traceMessageProcessing(String operationName) { return mono -> Mono.deferContextual(ctx -> { // 从上下文获取消息自带的traceId String messageTraceId = ctx.getOrDefault("message-trace-id", null); TraceContext parentContext = messageTraceId != null ? TraceContext.newBuilder().setTraceId(messageTraceId).build() : null; Span span = Tracer.currentTracer() .spanBuilder(operationName) .setParent(parentContext) .startSpan(); return mono .contextWrite(ctxInner -> ctxInner.put(TraceSpanThreadLocalAccessor.KEY, span)) .doFinally(signalType -> span.end()); }); }
业务代码只需调用该操作符即可:
// 先将消息的traceId写入上下文 processMessage(message) .contextWrite(ctx -> ctx.put("message-trace-id", message.getTraceId())) .transform(traceMessageProcessing("message-processing"))
日志适配
确保日志配置(如logback)中添加MDC占位符,自动打印traceId和spanId:
<pattern>%d{HH:mm:ss.SSS} [%thread] %-5level %logger{36} - %X{traceId:-} %X{spanId:-} - %msg%n</pattern>
内容的提问来源于stack exchange,提问作者sjanssen
相关产品推荐
相关产品推荐

