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

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)或显式上下文关联,并行执行时如果没有正确传递父观测的上下文,子任务会成为独立的链路,无法和父任务关联。


一、修正虚拟线程版本

  1. 先启动父观测,确保父上下文存在
  2. 异步任务中显式关联父观测的上下文,避免无差别捕获全部上下文
  3. 子观测创建时,通过父上下文绑定关联关系
  4. 等待所有异步任务完成,避免父观测提前结束导致链路断裂

修正后的代码:

@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版本

  1. 启动父观测并将上下文绑定到当前线程,确保并行任务能继承
  2. 子观测无需手动传递父实例,自动继承当前线程的观测上下文
  3. 限制并发数避免资源过载

修正后的代码:

@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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 17:03:22