在Micrometer Tracing(brave桥)中跨CompletableFuture传递追踪信息
问题:Spring Boot Micrometer Tracing在CompletableFuture多线程场景下TraceID/SpanID丢失
在Spring Boot中使用Micrometer Tracing实现分布式追踪时,单线程环境下功能正常,但使用CompletableFuture这类多线程操作时,子线程日志中的traceId/spanID为空,无法传递追踪上下文。
环境依赖(Gradle)
implementation 'org.springframework.boot:spring-boot-starter-actuator' implementation 'org.springframework.boot:spring-boot-starter-web' implementation 'io.micrometer:micrometer-tracing-bridge-brave'
相关代码
MainApplication类
import lombok.extern.slf4j.Slf4j; import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.RestController; @SpringBootApplication @RestController @Slf4j public class MainApplication { @GetMapping("async") String helloAsync() { log.info("hello from rest endpoint[GET]"); CompletableFutureUtil.performAsyncOperation(); return "hi"; } public static void main(String[] args) { SpringApplication.run(MainApplication.class, args); } }
CompletableFutureUtil类
import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; import lombok.extern.slf4j.Slf4j; @Slf4j public class CompletableFutureUtil { public static void performAsyncOperation() { // 创建CompletableFuture模拟异步流程 CompletableFuture<String> future = CompletableFuture.supplyAsync( () -> { // 模拟异步处理逻辑 log.info("inside CompletableFuture.supplyAsync.."); simulateAsyncOperation(); log.info("returning from CompletableFuture.supplyAsync.."); return "Simulated result"; }); // 成功回调 future.thenAccept( result -> { logInfo("Async process completed successfully. Result: " + result); }); // 异常回调 future.exceptionally( exception -> { logError("Async process failed. Reason: " + exception.getMessage()); return null; }); // 主线程其他逻辑 logInfo("Some other work is being done..."); logInfo("Main thread ended."); } private static void simulateAsyncOperation() { logInfo("Simulating async operation..."); try { TimeUnit.SECONDS.sleep(2); // 模拟延迟 } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } private static void logInfo(String message) { log.info("[INFO] " + message); } private static void logError(String message) { log.error("[ERROR] " + message); } }
问题日志
2023-12-09T00:29:29.117+05:30 " INFO [content-publisher-service,6573679104ea6b2868558e3949953a1e,68558e3949953a1e]" 42181 --- [nio-8080-exec-5] c.c.a.c.p.service.BootsApplication : hello from rest endpoint[GET] 2023-12-09T00:29:29.121+05:30 " INFO [content-publisher-service,6573679104ea6b2868558e3949953a1e,68558e3949953a1e]" 42181 --- [nio-8080-exec-5] c.c.a.c.p.s.util.CompletableFutureUtil : [INFO] Some other work is being done... 2023-12-09T00:29:29.121+05:30 " INFO [content-publisher-service,,]" 42181 --- [onPool-worker-1] c.c.a.c.p.s.util.CompletableFutureUtil : inside CompletableFuture.supplyAsync.. 2023-12-09T00:29:29.122+05:30 " INFO [content-publisher-service,6573679104ea6b2868558e3949953a1e,68558e3949953a1e]" 42181 --- [nio-8080-exec-5] c.c.a.c.p.s.util.CompletableFutureUtil : [INFO] Main thread ended. 2023-12-09T00:29:29.122+05:30 " INFO [content-publisher-service,,]" 42181 --- [onPool-worker-1] c.c.a.c.p.s.util.CompletableFutureUtil : [INFO] Simulating async operation... 2023-12-09T00:29:31.122+05:30 " INFO [content-publisher-service,,]" 42181 --- [onPool-worker-1] c.c.a.c.p.s.util.CompletableFutureUtil : returning from CompletableFuture.supplyAsync.. 2023-12-09T00:29:31.122+05:30 " INFO [content-publisher-service,,]" 42181 --- [onPool-worker-1] c.c.a.c.p.s.util.CompletableFutureUtil : [INFO] Async process completed successfully. Result: Simulated result
解决方案
问题根源:默认的CompletableFuture.supplyAsync()使用JDK自带的ForkJoinPool,该线程池不会自动传递Micrometer Tracing的上下文(Trace/Span信息),导致子线程无法获取父线程的追踪信息。
方案1:使用ContextSnapshot包装任务,手动传递上下文
通过Micrometer提供的ContextSnapshot工具,将父线程的追踪上下文包装到异步任务中,确保子线程能继承上下文。
修改CompletableFutureUtil类中的异步任务部分:
import io.micrometer.context.ContextSnapshot; // ... 其他代码不变 CompletableFuture<String> future = CompletableFuture.supplyAsync( ContextSnapshot.wrap(() -> { log.info("inside CompletableFuture.supplyAsync.."); simulateAsyncOperation(); log.info("returning from CompletableFuture.supplyAsync.."); return "Simulated result"; }) );
同时回调部分也需要包装:
future.thenAccept(ContextSnapshot.wrap(result -> { logInfo("Async process completed successfully. Result: " + result); })); future.exceptionally(ContextSnapshot.wrap(exception -> { logError("Async process failed. Reason: " + exception.getMessage()); return null; }));
方案2:使用Spring封装的TraceableExecutorService
自定义线程池并使用TraceableExecutorService包装,让线程池自动传递追踪上下文,再将该线程池传入supplyAsync方法。
- 配置线程池Bean:
import io.micrometer.tracing.Tracer; import io.micrometer.tracing.context.TraceableExecutorService; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; @Configuration public class AsyncConfig { @Bean public ExecutorService traceableExecutorService(Tracer tracer) { ExecutorService executorService = Executors.newFixedThreadPool(5); return TraceableExecutorService.create(executorService, tracer); } }
- 修改
CompletableFutureUtil为Spring组件并注入线程池:
import java.util.concurrent.ExecutorService; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Component; @Component @Slf4j public class CompletableFutureUtil { private final ExecutorService traceableExecutorService; @Autowired public CompletableFutureUtil(ExecutorService traceableExecutorService) { this.traceableExecutorService = traceableExecutorService; } public void performAsyncOperation() { CompletableFuture<String> future = CompletableFuture.supplyAsync( () -> { log.info("inside CompletableFuture.supplyAsync.."); simulateAsyncOperation(); log.info("returning from CompletableFuture.supplyAsync.."); return "Simulated result"; }, traceableExecutorService // 使用带追踪的线程池 ); // 回调同样使用该线程池 future.thenAcceptAsync( result -> logInfo("Async process completed successfully. Result: " + result), traceableExecutorService ); future.exceptionallyAsync( exception -> { logError("Async process failed. Reason: " + exception.getMessage()); return null; }, traceableExecutorService ); // 其他逻辑不变 logInfo("Some other work is being done..."); logInfo("Main thread ended."); } // 其他方法不变 }
方案3:使用Spring的@Async注解(推荐)
如果业务场景允许,直接使用Spring提供的@Async注解,Spring会自动处理追踪上下文的传递,无需手动处理线程池或包装任务。
- 启用异步支持:
在启动类添加@EnableAsync注解:
@SpringBootApplication @RestController @Slf4j @EnableAsync public class MainApplication { // ... 代码不变 }
- 重构异步逻辑为Spring异步方法:
import org.springframework.scheduling.annotation.Async; import org.springframework.stereotype.Component; import lombok.extern.slf4j.Slf4j; import java.util.concurrent.TimeUnit; @Component @Slf4j public class AsyncService { @Async public String performAsyncOperation() { log.info("inside async method.."); simulateAsyncOperation(); log.info("returning from async method.."); return "Simulated result"; } private void simulateAsyncOperation() { log.info("[INFO] Simulating async operation..."); try { TimeUnit.SECONDS.sleep(2); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } }
- 修改MainApplication中的调用:
@Autowired private AsyncService asyncService; @GetMapping("async") String helloAsync() { log.info("hello from rest endpoint[GET]"); asyncService.performAsyncOperation(); return "hi"; }
内容的提问来源于stack exchange,提问作者Arvind Kumar
相关产品推荐
相关产品推荐

