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

线程切换后Project Reactor上下文丢失问题求助

问题分析与解决方案

这不是Reactor Hooks.enableAutomaticContextPropagation()的Bug,而是你的实现没有适配非请求驱动的无限流场景——Spring Cloud GCP Trace默认是为HTTP请求/响应链路设计的,依赖请求上下文自动初始化Trace传播,但在无请求绑定的消息流中,需要手动适配Reactor的上下文传播机制。


核心问题原因

  1. 缺少Trace上下文的传播载体:Reactor的自动上下文传播需要明确注册ThreadLocalAccessor,用来识别并跨线程传递Trace Span这类自定义上下文数据,你当前没有注册对应Accessor,导致线程切换后上下文丢失。
  2. Span上下文绑定方式错误:仅用contextWrite写入Span到Reactor上下文,但未同步绑定到当前线程的ThreadLocal,导致后续算子切换线程时无法正确恢复Span上下文。
  3. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 17:45:57