Spring Boot 3.4+Micrometer并行Span链路追踪关联问题
问题描述
我有一个基于Spring Boot 3.4的Web服务,使用Micrometer进行观测。该服务的接口接收消息列表,需调用一个无法控制的外部REST API(处理耗时较长)。
请求示例Payload(实际包含大量消息):
[ "somefirst message", "some second message", "etc" ]
为处理该场景,我编写了两个版本的Spring Rest Controller:
1. 虚拟线程实现
@PostMapping("/question2") public String question2(@RequestBody List<String> messages) { Observation parent = Observation.createNotStarted("parent", registry); try (ExecutorService executor = Executors.newVirtualThreadPerTaskExecutor()) { for (String message : messages) { parent.observe(() -> { ContextExecutorService.wrap(executor, () -> ContextSnapshotFactory.builder().build().captureAll()) .submit(() -> { Observation.createNotStarted("child" + message, registry).observe(() -> { messageProducer.sendRequest("topic-loom-micrometer", message + Thread.currentThread().getName()); }); }); }); } executor.shutdown(); } return "it seems everything went fine"; }
2. Stream + Gatherers mapConcurrent实现
@PostMapping("/question3") public List<String> question3(@RequestBody List<String> messages) { Observation parent = Observation.createNotStarted("parent", registry); return messages.stream() .gather(Gatherers.mapConcurrent(messages.size() + 1, oneMessage -> sendRequestInParallel(oneMessage, parent))) .toList(); } private String sendRequestInParallel(String oneMessage, Observation parent) { return parent.observe(() -> { return Observation.createNotStarted("child" + oneMessage, registry).observe(() -> { return restClient.post() .uri("http://localhost:8081/justString?name=" + oneMessage) .retrieve() .body(String.class); }); }); }
问题:普通串行调用时链路呈现层级关联结构,但使用虚拟线程或mapConcurrent并行调用时,无法得到预期的「主任务下包含所有并行子任务」的两层链路结构,链路追踪结果不符合预期。
解决方案
核心问题分析
Micrometer Observation的上下文传递依赖线程本地(ThreadLocal)或显式上下文关联,并行执行时如果没有正确传递父观测的上下文,子任务会成为独立的链路,无法和父任务关联。
一、修正虚拟线程版本
- 先启动父观测,确保父上下文存在
- 异步任务中显式关联父观测的上下文,避免无差别捕获全部上下文
- 子观测创建时,通过父上下文绑定关联关系
- 等待所有异步任务完成,避免父观测提前结束导致链路断裂
修正后的代码:
@PostMapping("/question2") public String question2(@RequestBody List<String> messages) { // 创建并启动父观测,用try-with-resources自动管理生命周期 try (Observation parent = Observation.createNotStarted("parent", registry).start()) { try (ExecutorService executor = Executors.newVirtualThreadPerTaskExecutor()) { Observation.Context parentContext = parent.context(); for (String message : messages) { // 仅传递观测上下文给异步线程 ContextExecutorService wrappedExecutor = ContextExecutorService.wrap(executor, () -> ContextSnapshot.builder().add(ObservationThreadLocalAccessor.KEY, parentContext).build() ); wrappedExecutor.submit(() -> { // 子观测显式关联父上下文 try (Observation child = Observation.createNotStarted("child-" + message, registry) .parent(parentContext) .start()) { messageProducer.sendRequest("topic-loom-micrometer", message + Thread.currentThread().getName()); } }); } executor.shutdown(); executor.awaitTermination(10, TimeUnit.SECONDS); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new RuntimeException("Task interrupted", e); } } return "it seems everything went fine"; }
二、修正Stream + mapConcurrent版本
- 启动父观测并将上下文绑定到当前线程,确保并行任务能继承
- 子观测无需手动传递父实例,自动继承当前线程的观测上下文
- 限制并发数避免资源过载
修正后的代码:
@PostMapping("/question3") public List<String> question3(@RequestBody List<String> messages) { // 创建并启动父观测 try (Observation parent = Observation.createNotStarted("parent", registry).start()) { // 将父观测上下文绑定到当前线程,并行任务自动继承 try (Observation.Scope scope = parent.openScope()) { return messages.stream() .gather(Gatherers.mapConcurrent(Math.min(messages.size(), 10), oneMessage -> sendRequestInParallel(oneMessage) )) .toList(); } } } private String sendRequestInParallel(String oneMessage) { // 子观测自动继承父上下文 try (Observation child = Observation.createNotStarted("child-" + oneMessage, registry).start()) { return restClient.post() .uri("http://localhost:8081/justString?name=" + oneMessage) .retrieve() .body(String.class); } }
通用注意事项
- 确保Spring Boot 3.x默认的Micrometer Observation自动配置生效
- 父观测的生命周期必须覆盖所有子任务的执行时间,避免父观测提前结束导致子链路独立
- 并行场景下优先使用Micrometer提供的上下文传递机制,减少手动捕获全部上下文的开销
内容的提问来源于stack exchange,提问作者PatPanda
相关产品推荐
相关产品推荐

