Spring 5 Reactor应用中MDC跨线程传递用户ID问题求助
解决Reactor Subscriber Context与Logback MDC跨线程同步问题
这个问题我太熟了——Reactor的Subscriber Context和Logback的MDC天生就不是一个路子:Context跟着订阅链走,MDC绑定线程局部变量,线程一切换,MDC里的用户ID直接就丢了。而且要覆盖到自己代码和外部库的所有日志,确实得花点心思,我给你一套彻底解决的方案:
核心思路
我们需要在Reactor流的每一次线程切换时,自动把Subscriber Context里的用户ID同步到当前线程的MDC中,同时在流结束后清理MDC(避免线程池复用导致的脏数据污染)。通过Reactor的全局Hook,可以让所有Mono/Flux(包括外部库的)都自动应用这个同步逻辑。
具体实现步骤
1. 自定义MDC同步Operator
先写一个工具类,实现Reactor的Operator拦截逻辑,负责Context和MDC的同步:
import org.slf4j.MDC; import reactor.core.publisher.Mono; import reactor.core.publisher.Operators; import reactor.util.context.Context; public class MdcContextSyncUtils { // 替换成你存在Subscriber Context里的用户ID键名 private static final String USER_ID_CONTEXT_KEY = "userId"; // 给单个Mono添加MDC同步逻辑 public static <T> Mono<T> syncMdcWithContext(Mono<T> mono) { return mono.transform(Operators.lift((scannable, subscriber) -> new MdcSyncSubscriber<>(subscriber))); } // 内部订阅者,负责同步Context到MDC private static class MdcSyncSubscriber<T> implements reactor.core.CoreSubscriber<T> { private final reactor.core.CoreSubscriber<T> delegate; private Context currentContext; public MdcSyncSubscriber(reactor.core.CoreSubscriber<T> delegate) { this.delegate = delegate; } @Override public void onSubscribe(reactor.core.Disposable s) { this.currentContext = delegate.currentContext(); // 初始化时把Context里的用户ID放入MDC currentContext.getOrEmpty(USER_ID_CONTEXT_KEY) .ifPresent(userId -> MDC.put(USER_ID_CONTEXT_KEY, userId.toString())); delegate.onSubscribe(s); } @Override public void onNext(T t) { // 每次处理信号前重新同步(防止线程切换后MDC丢失) currentContext.getOrEmpty(USER_ID_CONTEXT_KEY) .ifPresent(userId -> MDC.put(USER_ID_CONTEXT_KEY, userId.toString())); delegate.onNext(t); } @Override public void onError(Throwable t) { try { delegate.onError(t); } finally { // 流异常结束后清理MDC,避免线程池复用污染 MDC.remove(USER_ID_CONTEXT_KEY); } } @Override public void onComplete() { try { delegate.onComplete(); } finally { // 流正常结束后清理MDC MDC.remove(USER_ID_CONTEXT_KEY); } } @Override public Context currentContext() { return delegate.currentContext(); } } }
2. 全局配置,让所有Reactor流自动同步
通过Spring配置类,注册Reactor的全局Hook,这样所有Mono/Flux都会自动应用MDC同步逻辑:
import org.springframework.context.annotation.Configuration; import reactor.core.publisher.Hooks; import javax.annotation.PostConstruct; @Configuration public class ReactorMdcGlobalConfig { @PostConstruct public void initMdcSync() { // 全局拦截所有Reactor操作符,自动同步Context到MDC Hooks.onEachOperator(MdcContextSyncUtils::syncMdcWithContext); } }
3. 配置Logback输出MDC字段
在你的logback-spring.xml(或logback.xml)里,修改日志格式,添加MDC中的用户ID字段:
<appender name="CONSOLE" class="ch.qos.logback.core.ConsoleAppender"> <encoder> <!-- %X{userId} 就是读取MDC里的userId值 --> <pattern>%d{yyyy-MM-dd HH:mm:ss.SSS} [%thread] %-5level %logger{36} - %msg - 用户ID:%X{userId}%n</pattern> </encoder> </appender>
关键注意事项
- 线程池复用问题:一定要在
onError和onComplete中清理MDC!否则线程池里的线程下次被复用的时候,会带着之前的用户ID,导致日志串号。 - 外部库的日志适配:只要外部库的日志是在Reactor流的处理线程中打印的,这个方案就能自动生效。如果外部库用了自己的独立线程池(比如某些异步HTTP客户端),你需要手动把Context里的用户ID传递到外部线程,然后在外部线程中手动调用
MDC.put()设置。 - Context传递正确性:确保你在业务代码中正确传递了Subscriber Context,比如用
contextWrite设置用户ID:// 示例:在流中设置用户ID到Context Mono.just("业务数据") .contextWrite(Context.of("userId", currentUserId)) .flatMap(data -> businessService.process(data));
简化方案:使用Reactor-Tools库
如果你不想自定义Operator,可以直接用Reactor官方提供的reactor-tools库,里面有现成的MdcOperator:
- 添加依赖:
<dependency> <groupId>io.projectreactor.tools</groupId> <artifactId>reactor-tools</artifactId> <version>3.4.22</version> <!-- 版本要和你的Reactor版本匹配 --> </dependency>
- 修改全局配置:
@PostConstruct public void initMdcSync() { Hooks.onEachOperator(MdcOperator.lift()); }
这个默认会把Context里的所有键值对同步到MDC中,如果你只需要特定字段,还是自定义Operator更灵活。
内容的提问来源于stack exchange,提问作者Luca
相关产品推荐
相关产品推荐

