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

Resilience4j重试时丢失MDC上下文与correlationId HTTP头问题咨询

问题分析与解决方案

首先明确:这是预期行为。SLF4J的MDC基于ThreadLocal实现,线程切换后上下文自然丢失;Resilience4j默认不会自动传播MDC这类线程绑定的上下文,尤其是异步重试/线程池执行场景下,必须手动配置上下文传播逻辑。

针对你的技术栈(Resilience4j 2.0.2 + Spring Boot 3.0.2 + Java 17),以下是可行的解决方案:

1. 自定义TaskDecorator传播MDC到Resilience4j线程池

如果你的重试是通过Resilience4j的线程池(比如结合ThreadPoolBulkhead,或异步重试配置)执行的,可以通过自定义TaskDecorator包装线程任务,将原线程的MDC上下文复制到重试线程中:

步骤1:实现MDC任务装饰器

import org.slf4j.MDC;
import org.springframework.core.task.TaskDecorator;
import java.util.Map;

public class MdcTaskDecorator implements TaskDecorator {
    @Override
    public Runnable decorate(Runnable runnable) {
        // 捕获当前线程的MDC上下文
        Map<String, String> mdcContext = MDC.getCopyOfContextMap();
        return () -> {
            try {
                // 将MDC上下文设置到新线程中
                if (mdcContext != null) {
                    MDC.setContextMap(mdcContext);
                }
                runnable.run();
            } finally {
                // 清理当前线程的MDC,避免污染后续任务
                MDC.clear();
            }
        };
    }
}

步骤2:配置Resilience4j线程池使用该装饰器

通过Spring的Customizer定制Resilience4j的线程池执行器:

import io.github.resilience4j.bulkhead.ThreadPoolBulkheadConfig;
import io.github.resilience4j.bulkhead.ThreadPoolBulkheadRegistry;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;

@Configuration
public class Resilience4jConfig {

    @Bean
    public Customizer<ThreadPoolBulkheadRegistry> threadPoolBulkheadCustomizer() {
        return registry -> registry.configureDefault(id ->
                ThreadPoolBulkheadConfig.custom()
                        .executorService(threadPoolTaskExecutor())
                        .build());
    }

    private ThreadPoolTaskExecutor threadPoolTaskExecutor() {
        ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
        executor.setCorePoolSize(10);
        executor.setMaxPoolSize(20);
        executor.setQueueCapacity(100);
        executor.setThreadNamePrefix("resilience4j-pool-");
        // 绑定MDC装饰器
        executor.setTaskDecorator(new MdcTaskDecorator());
        executor.initialize();
        return executor;
    }
}

如果是直接配置异步重试的线程池,逻辑类似,只需将配置目标切换为Retry的异步执行器即可。

2. Feign客户端的MDC上下文传播补充

除了线程池装饰器,还可以在Feign拦截器中确保重试时能重新获取correlationId:

  • 若correlationId来自请求头,可以在拦截器中直接从请求头提取,而非仅依赖MDC;
  • 或者通过Resilience4j的RetryRegistry自定义重试监听器,在重试前将原请求的correlationId重新放入MDC:
import io.github.resilience4j.retry.Retry;
import io.github.resilience4j.retry.RetryRegistry;
import org.slf4j.MDC;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;

@Configuration
public class FeignRetryConfig {

    @Bean
    public RetryRegistry retryRegistry() {
        RetryRegistry registry = RetryRegistry.ofDefaults();
        registry.retry("feign-retry").getEventPublisher()
                .onRetry(event -> {
                    // 假设correlationId存储在请求属性中,根据实际场景调整获取方式
                    String correlationId = event.getRetry().getMetadata().get("correlationId").toString();
                    MDC.put("correlationId", correlationId);
                });
        return registry;
    }
}

3. Spring WebClient(Reactive场景)的MDC传播

Reactive场景下不能依赖ThreadLocal,需改用Reactor的Context传递上下文:

步骤1:在WebClient过滤器中注入correlationId到Reactor Context

import org.springframework.web.reactive.function.client.ExchangeFilterFunction;
import reactor.core.publisher.Mono;

public class CorrelationIdFilter implements ExchangeFilterFunction {
    @Override
    public Mono<ClientResponse> filter(ClientRequest request, ExchangeFunction next) {
        String correlationId = request.headers().firstHeader("X-Correlation-ID");
        return next.exchange(request)
                .contextWrite(context -> context.put("correlationId", correlationId));
    }
}

步骤2:订阅时从Context中提取并设置MDC

在调用WebClient的地方,通过doOnEach操作符将Context中的correlationId放入MDC:

webClient.get()
        .uri("/api/resource")
        .header("X-Correlation-ID", MDC.get("correlationId"))
        .retrieve()
        .bodyToMono(String.class)
        .doOnEach(signal -> {
            signal.getContextView().getOrEmpty("correlationId")
                    .ifPresent(id -> MDC.put("correlationId", id));
        })
        .doFinally(signal -> MDC.clear());

也可以封装成全局的Reactor上下文处理器,避免重复代码。


内容的提问来源于stack exchange,提问作者Nacho Martín Moreno

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 19:47:46