从 SkyWalking 专用实现迁移到 OpenTelemetry:搭建可复刻的分布式 Trace、日志中心与异步访问日志系统

本文记录一次真实 Java 微服务项目的可观测性改造:使用 OpenTelemetry Java Agent 自动采集 Trace 和
SLF4J/Logback 日志,经 OpenTelemetry Collector 写入 Elasticsearch 9.x,再用 Kibana 按 Trace ID
串联一次请求涉及的 HTTP、Feign、JDBC、RocketMQ 和业务日志;同时将数据库 API 访问日志从同步
Feign 调用改为 RocketMQ 异步、幂等落库。

文中的代码、配置和验证结果来自一个 Spring Boot 4.1、JDK 25、Spring Cloud 微服务项目。版本基线
截至 2026-08-23。中间件升级很快,复制到新项目时应再次核对官方发布页。

本文对应的核心文件

如果读者正在查看配套仓库,可以用下面的索引快速定位完整实现:

内容 文件
总体部署说明 docs/observability.md
OpenTelemetry 版本管理 handao-dependencies/pom.xml
Trace 自动配置 HandaoTracerAutoConfiguration.java
Trace ID 工具 TracerUtils.java
响应 Trace Header TraceFilter.java
CORS Header 暴露 HandaoWebAutoConfiguration.java
访问日志生产者 ApiAccessLogProducer.java
访问日志消费者 ApiAccessLogMessageConsumer.java
Collector 配置 otel-collector-config.yaml
Docker Compose docker-compose.yml
RocketMQ Broker 配置 broker.conf

一、为什么要重做分布式日志方案

微服务系统出现线上问题时,我们通常需要回答下面几个问题:

  1. 用户的这次请求经过了哪些服务?
  2. 每一跳分别耗时多久?慢在网关、数据库、远程调用还是消息队列?
  3. 某一条异常日志属于哪一次请求?
  4. 请求返回给前端的 Trace ID,能否直接检索出完整调用链和相关日志?
  5. API 访问日志落库是否会增加业务接口延迟?
  6. 消息重复投递时,数据库是否会出现重复访问日志?

旧实现存在两个典型问题。

第一,业务代码直接依赖 SkyWalking Toolkit,并在网关、异常处理器等位置调用专有 API。这样做会让
业务代码与具体 APM 产品绑定,后续切换后端或采集方案时改动面很大。

第二,API 访问日志通过 Feign 同步调用 infra 服务落库。访问日志本质上属于旁路审计数据,业务线程
不应该等待日志服务和数据库;一旦 infra 服务抖动,日志链路还可能反向拖慢甚至影响业务请求。

本次改造确立了以下目标:

  • 使用厂商中立的 OpenTelemetry 标准;
  • 尽量通过 Java Agent 零侵入采集,减少手工埋点;
  • Trace 与 SLF4J 日志通过同一个 trace_id 关联;
  • 应用只连接 Collector,不直接连接 Elasticsearch;
  • API 访问日志通过 RocketMQ 异步落库;
  • 考虑 RocketMQ 至少一次投递,实现消费者幂等;
  • 保留多租户上下文;
  • 本地、IDE、Docker 和生产环境使用同一套协议和拓扑;
  • Elasticsearch 使用 9.x 原生 OTLP Endpoint,不再额外维护自定义日志索引映射。

二、先区分三类“日志”

在开始搭建之前,必须先区分下面三类数据。很多方案的问题,正是因为把它们混为一谈。

2.1 Trace 与 Span

Trace 表示一次端到端请求,Trace 内部由多个 Span 组成。例如:

Browser
  -> Gateway HTTP Server Span
  -> System HTTP Server Span
  -> Feign Client Span
  -> Infra HTTP Server Span
  -> JDBC Client Span

整个调用链共享一个 32 位十六进制 Trace ID,每个 Span 有自己的 16 位 Span ID。

2.2 应用日志

应用日志是代码通过 SLF4J 输出的内容:

log.info("[getPermissionInfo][OpenTelemetry 链路验证,userId({})]", userId);

Java Agent 会捕获 Logback 日志事件,并在事件发生时附加当前 trace_idspan_id。因此同一请求中的
业务日志可以和 Trace 关联。

2.3 API 访问日志

API 访问日志是结构化审计记录,通常包含:

  • 用户与租户;
  • 应用名;
  • URL、HTTP Method、IP、User-Agent;
  • 脱敏后的请求参数;
  • 返回码、错误信息;
  • 开始时间、结束时间、耗时;
  • Trace ID。

它需要分页查询、合规留存和业务报表,所以仍然落在关系型数据库中。但落库不应该阻塞业务线程,
因此改为 RocketMQ 异步处理。

最终形成两条互相独立、又通过 Trace ID 关联的数据链路:

flowchart LR U["浏览器 / 调用方"] --> G["Gateway"] G --> S["System / 其他业务服务"] S --> I["Infra / 下游服务"] G -. "Trace + SLF4J Log" .-> C["OpenTelemetry Collector"] S -. "Trace + SLF4J Log" .-> C I -. "Trace + SLF4J Log" .-> C C --> E["Elasticsearch 9.x /_otlp"] E --> K["Kibana"] S -- "API Access Log Event" --> R["RocketMQ"] I -- "API Access Log Event" --> R R --> AC["Infra Access Log Consumer"] AC --> DB[("MySQL infra_api_access_log")]

三、技术选型

3.1 组件与版本

本次实现采用以下版本:

组件 版本 作用
JDK 25 应用运行时
Spring Boot 4.1.x 应用框架
OpenTelemetry Java API 1.64.0 少量业务手工埋点使用稳定 API
OpenTelemetry Java Agent 2.28.1 HTTP、Feign、JDBC、RocketMQ、日志自动采集与上下文传播
OpenTelemetry Collector Contrib 0.159.0 OTLP 网关、批处理、限流、重试、队列
Elasticsearch / Kibana 9.5.2 Trace 与日志存储、查询和展示
RocketMQ Server 5.5.0 API 访问日志异步投递
RocketMQ Spring Boot Starter 2.3.6 Java 应用生产和消费消息

版本资料:

3.2 为什么选择 Java Agent,而不是在应用里初始化完整 SDK

Java Agent 的优势是:

  • 在应用 main 方法前加载;
  • 自动创建 HTTP Server/Client Span;
  • 自动处理 Feign、JDBC、RocketMQ 等常用组件;
  • 自动注入和提取 W3C traceparentbaggage
  • 自动捕获 Logback 日志;
  • 不需要每个服务重复编写 SDK、Exporter、Processor 初始化代码;
  • 采集配置可以通过环境变量统一管理。

应用代码只保留 opentelemetry-api,用于读取当前 Trace ID 或创建少量业务 Span:

<properties>
    <opentelemetry.version>1.64.0</opentelemetry.version>
</properties>

<dependencyManagement>
    <dependencies>
        <dependency>
            <groupId>io.opentelemetry</groupId>
            <artifactId>opentelemetry-api</artifactId>
            <version>${opentelemetry.version}</version>
        </dependency>
    </dependencies>
</dependencyManagement>

不要同时在应用中初始化一套 OpenTelemetry SDK,又加载 Java Agent。两套 SDK/Exporter 并存容易导致:

  • 重复 Span;
  • 重复日志;
  • GlobalOpenTelemetry 注册冲突;
  • 版本依赖复杂;
  • 资源属性和采样策略不一致。

3.3 为什么必须保留 Collector

应用理论上可以直接把 OTLP 发给 Elasticsearch,但生产环境不建议这么做。Collector 提供了应用与存储
后端之间的隔离层:

  • 批量发送,降低后端请求数;
  • 内存保护;
  • 发送队列;
  • 失败重试;
  • 统一认证和 TLS;
  • 数据过滤、脱敏和属性补充;
  • 后续替换或并行导出到其他后端时,不需要修改所有应用。

应用只知道 Collector 的 OTLP 地址,不需要知道 Elasticsearch 的地址和凭证。

3.4 为什么访问日志选 RocketMQ,而不是线程池异步

@Async 或本地线程池只能把数据库操作移出请求线程,但无法解决:

  • 应用进程崩溃导致内存任务丢失;
  • infra 服务临时不可用;
  • 消费速度低于生产速度;
  • 跨服务削峰;
  • 重试和死信管理。

RocketMQ 提供持久化、消费组、重试和死信队列,更适合审计日志事件。但 RocketMQ 是至少一次投递,
所以必须额外设计幂等。

四、总体调用过程

下面是一条用户请求从前端到 ES 和访问日志数据库的完整时序。

sequenceDiagram autonumber participant B as Browser participant G as Gateway participant S as System Service participant DB as MySQL participant MQ as RocketMQ participant I as Infra Consumer participant C as OTel Collector participant ES as Elasticsearch B->>G: HTTP Request Note over G: Agent 创建 Gateway Server Span G->>S: HTTP / Feign + traceparent Note over S: Agent 提取上下文并创建 Server Span S->>DB: JDBC Query Note over DB,S: Agent 创建 JDBC Client Span S-->>S: log.info(...) Note over S: Logback 事件附带 trace_id/span_id S-->>MQ: asyncSend API Access Log Note over MQ: 消息携带 Trace Context 与 tenant-id S-->>G: Response + trace-id G-->>B: Response + trace-id S-->>C: OTLP/HTTP Traces + Logs G-->>C: OTLP/HTTP Traces + Logs C-->>ES: /_otlp/v1/traces + /_otlp/v1/logs MQ->>I: Consume Access Log I->>DB: 幂等 INSERT infra_api_access_log

五、应用侧接入 OpenTelemetry

5.1 下载并校验 Java Agent

不要在构建时下载一个不固定版本的 latest 文件。应固定版本,并校验 SHA-256:

mkdir -p .local/otel

curl -fL \
  https://github.com/open-telemetry/opentelemetry-java-instrumentation/releases/download/v2.28.1/opentelemetry-javaagent.jar \
  -o .local/otel/opentelemetry-javaagent.jar

echo "faa89bdeebf9b1f52be4a4374689176717b02a59df2d8f8b6eb9aa39f9292589  .local/otel/opentelemetry-javaagent.jar" \
  | sha256sum -c -

校验成功应输出:

.local/otel/opentelemetry-javaagent.jar: OK

5.2 IDE 启动参数

System 服务的 VM options:

-javaagent:/absolute/path/opentelemetry-javaagent.jar
-Dotel.service.name=system-server
-Dotel.exporter.otlp.endpoint=http://127.0.0.1:4318
-Dotel.exporter.otlp.protocol=http/protobuf
-Dotel.propagators=tracecontext,baggage
-Dotel.traces.exporter=otlp
-Dotel.logs.exporter=otlp
-Dotel.metrics.exporter=none

Infra 服务只需要修改服务名:

-Dotel.service.name=infra-server

Gateway 服务使用:

-Dotel.service.name=gateway-server

几个关键配置的含义:

配置 含义
otel.service.name 服务的稳定标识,查询和聚合的重要维度
otel.exporter.otlp.endpoint Collector 地址,不是 Elasticsearch 地址
otel.exporter.otlp.protocol Agent 2.x 推荐使用 http/protobuf
otel.propagators 使用 W3C Trace Context 与 Baggage
otel.traces.exporter 开启 Trace OTLP 导出
otel.logs.exporter 开启日志 OTLP 导出
otel.metrics.exporter 本方案暂时关闭 OTel Metrics,避免与现有 Micrometer 重复

启动日志中应出现类似内容:

[otel.javaagent] opentelemetry-javaagent - version: 2.28.1

如果完全看不到 Agent 启动信息,后续所有 Trace ID 问题都应先检查 -javaagent 路径,而不是检查业务代码。

5.3 Docker 镜像内置 Agent

容器镜像使用多阶段构建,构建阶段下载并校验 Agent,运行阶段只复制最终 JAR:

FROM alpine:3.22 AS otel-agent
ARG OTEL_JAVA_AGENT_VERSION=2.28.1
ARG OTEL_JAVA_AGENT_SHA256=faa89bdeebf9b1f52be4a4374689176717b02a59df2d8f8b6eb9aa39f9292589

RUN wget -q -O /opentelemetry-javaagent.jar \
    "https://github.com/open-telemetry/opentelemetry-java-instrumentation/releases/download/v${OTEL_JAVA_AGENT_VERSION}/opentelemetry-javaagent.jar" \
    && echo "${OTEL_JAVA_AGENT_SHA256}  /opentelemetry-javaagent.jar" | sha256sum -c -

FROM eclipse-temurin:25-jre
WORKDIR /app
COPY ./target/application.jar app.jar
COPY --from=otel-agent /opentelemetry-javaagent.jar /opt/opentelemetry-javaagent.jar

ENV JAVA_TOOL_OPTIONS="-javaagent:/opt/opentelemetry-javaagent.jar" \
    OTEL_SERVICE_NAME="system-server" \
    OTEL_EXPORTER_OTLP_ENDPOINT="http://otel-collector:4318" \
    OTEL_EXPORTER_OTLP_PROTOCOL="http/protobuf" \
    OTEL_PROPAGATORS="tracecontext,baggage" \
    OTEL_TRACES_EXPORTER="otlp" \
    OTEL_LOGS_EXPORTER="otlp" \
    OTEL_METRICS_EXPORTER="none"

CMD java ${JAVA_OPTS} -jar app.jar

这里使用 JAVA_TOOL_OPTIONS,JVM 会自动读取它,因此不依赖启动脚本是否把 -javaagent 拼进命令行。

5.4 读取当前 Trace ID

业务代码不要再调用 SkyWalking Toolkit,统一通过稳定的 OpenTelemetry API 读取当前上下文:

public final class TracerUtils {

    public static final String HEADER_TRACE_ID = "trace-id";

    private TracerUtils() {
    }

    public static String getTraceId() {
        SpanContext context = Span.current().getSpanContext();
        return context.isValid() ? context.getTraceId() : "";
    }
}

必须检查 SpanContext#isValid()。如果 Agent 未加载,Span.current() 仍然可以调用,但得到的是无效上下文;
此时不能把全零 Trace ID 或其他占位值返回给前端。

5.5 把 Trace ID 返回给前端

Agent 在 Servlet Filter 之前创建 HTTP Server Span,所以 Filter 中可以读取当前 Trace ID:

public class TraceFilter extends OncePerRequestFilter {

    @Override
    protected void doFilterInternal(HttpServletRequest request,
                                    HttpServletResponse response,
                                    FilterChain chain)
            throws IOException, ServletException {
        String traceId = TracerUtils.getTraceId();
        if (!traceId.isEmpty()) {
            response.setHeader(TracerUtils.HEADER_TRACE_ID, traceId);
        }
        chain.doFilter(request, response);
    }
}

注册 Filter,并保证它早于访问日志 Filter:

@Bean
public FilterRegistrationBean<TraceFilter> traceFilter() {
    FilterRegistrationBean<TraceFilter> bean = new FilterRegistrationBean<>();
    bean.setFilter(new TraceFilter());
    bean.setOrder(WebFilterOrderEnum.TRACE_FILTER);
    return bean;
}

响应示例:

HTTP/1.1 200
trace-id: 8a7c479d422fa7bc0fe1aade0d703a24
Content-Type: application/json;charset=UTF-8

5.6 不要忘记 CORS Expose Headers

这是一个很隐蔽的问题:浏览器开发者工具可能看到 trace-id,但前端 JavaScript 读取不到,因为
trace-id 不是 CORS safelisted response header。

需要显式暴露:

CorsConfiguration config = new CorsConfiguration();
config.setAllowCredentials(true);
config.addAllowedOriginPattern("*");
config.addAllowedHeader("*");
config.addAllowedMethod("*");
config.addExposedHeader(TracerUtils.HEADER_TRACE_ID);

响应中会增加:

Access-Control-Expose-Headers: trace-id

前端读取方式:

// Axios
const traceId = response.headers['trace-id']

// Fetch
const traceId = response.headers.get('trace-id')

HTTP Header 名不区分大小写,但 Axios 通常会把 Header key 规范化成小写,因此读取时使用
response.headers['trace-id']

5.7 Logback 中打印 Trace ID 和 Span ID

Java Agent 会在 Logback 事件快照中注入 MDC 字段。日志格式直接读取这些字段:

<property name="CONSOLE_LOG_PATTERN"
          value="%d{yyyy-MM-dd HH:mm:ss.SSS} [%thread] [trace_id=%X{trace_id:-} span_id=%X{span_id:-}] %highlight(%-5level) %cyan(%logger{50}:%L) - %msg%n"/>

<property name="FILE_LOG_PATTERN"
          value="%d{yyyy-MM-dd HH:mm:ss.SSS} [%thread] [trace_id=%X{trace_id:-} span_id=%X{span_id:-}] %-5level %logger{50}:%L - %msg%n"/>

业务代码不需要手动向 MDC 写入 Trace ID:

log.info("[getPermissionInfo][OpenTelemetry 链路验证,userId({})]", userId);

输出类似:

2026-08-23 12:47:10.123 [http-nio-48081-exec-4]
[trace_id=8a7c479d422fa7bc0fe1aade0d703a24 span_id=12b91feec45048af]
INFO  c.h.h.m.s.c.a.a.AuthController:104 -
[getPermissionInfo][OpenTelemetry 链路验证,userId(1)]

Logback 仍然保留控制台和滚动文件 Appender,Java Agent 额外捕获日志事件并通过 OTLP 批量发送。
不需要再配置一个自定义 Elasticsearch Appender,也不建议让每个应用直接写 ES。

5.8 少量业务 Span

自动埋点负责基础设施 Span,业务关键动作可以使用注解和 AOP 创建业务 Span:

@Around("@annotation(trace)")
public Object around(ProceedingJoinPoint joinPoint, BizTrace trace) throws Throwable {
    Span span = tracer.spanBuilder(buildOperationName(joinPoint, trace))
            .setAttribute("component", "biz")
            .startSpan();
    try (Scope ignored = span.makeCurrent()) {
        return joinPoint.proceed();
    } catch (Throwable throwable) {
        span.recordException(throwable);
        span.setStatus(StatusCode.ERROR, throwable.getMessage());
        throw throwable;
    } finally {
        span.end();
    }
}

不要把完整异常堆栈再作为一个字符串 Attribute 写入 Span。recordException 已经使用标准事件表示异常;
额外写大字符串会显著增加 Span 体积和 ES 存储成本。

六、OpenTelemetry Collector 配置

6.1 完整的开发环境配置

extensions:
  health_check:
    endpoint: 0.0.0.0:13133

receivers:
  otlp:
    protocols:
      grpc:
        endpoint: 0.0.0.0:4317
      http:
        endpoint: 0.0.0.0:4318

processors:
  memory_limiter:
    check_interval: 1s
    limit_mib: 384
    spike_limit_mib: 64
  batch:
    timeout: 1s
    send_batch_size: 8192
    send_batch_max_size: 10000

exporters:
  otlphttp/elasticsearch:
    endpoint: http://elasticsearch:9200/_otlp
    compression: gzip
    timeout: 30s
    sending_queue:
      enabled: true
      sizer: bytes
      queue_size: 50000000
      block_on_overflow: true
      batch:
        flush_timeout: 1s
        min_size: 1000000
        max_size: 4000000
    retry_on_failure:
      enabled: true
      initial_interval: 1s
      max_interval: 30s
      max_elapsed_time: 300s

service:
  extensions: [health_check]
  pipelines:
    traces:
      receivers: [otlp]
      processors: [memory_limiter, batch]
      exporters: [otlphttp/elasticsearch]
    logs:
      receivers: [otlp]
      processors: [memory_limiter, batch]
      exporters: [otlphttp/elasticsearch]

6.2 配置解释

otlp Receiver 同时开放:

  • 4317:OTLP/gRPC;
  • 4318:OTLP/HTTP。

本项目 Java Agent 使用 http/protobuf,因此主要使用 4318。

memory_limiter 必须放在处理链前面。当 Collector 内存达到阈值时,它会施加反压,避免进程因 OOM
直接退出。

batch 将小批数据合并后再导出,降低网络请求和 ES 写入开销。

sending_queueretry_on_failure 用于处理 Elasticsearch 短暂不可用。但开发配置中的队列仍然是
内存队列,Collector 进程退出后数据会丢失;生产环境应使用持久队列,后文会说明。

Elasticsearch OTLP Endpoint 的基础地址是:

http://elasticsearch:9200/_otlp

Collector 的 otlphttp Exporter 会在其后追加标准 signal path:

/_otlp/v1/traces
/_otlp/v1/logs

6.3 用官方二进制验证 Collector 配置

不要只依靠 YAML 解析器。YAML 语法正确,不代表 Collector 组件配置正确。应使用目标版本执行:

otelcol-contrib validate --config script/docker/otel-collector-config.yaml

命令退出码为 0 且没有错误输出,才表示该版本 Collector 接受配置。

七、Elasticsearch 9.x 原生 OTLP 存储

7.1 为什么使用原生 Endpoint

Elasticsearch 9.x 可以直接接收 OTLP/HTTP。这样不需要:

  • 在 Collector 中把 OTel 字段手工转换成自定义 JSON;
  • 自己维护 Trace 和 Log index template;
  • 自己定义 trace_idspan_id 等字段类型;
  • 使用老旧或非标准的日志 Appender 直写 ES。

Elasticsearch 会自动创建 OTel Data Stream 和映射。

7.2 实际生成的数据流

本次环境实际生成了:

logs-generic.otel-default
traces-generic.otel-default

查询数据流:

curl 'http://127.0.0.1:9200/_data_stream?pretty'

7.3 字段名与常见误区

在本方案的 Elastic 原生 OTLP Data Stream 中,实际字段是:

含义 ES 字段
Trace ID trace_id
Span ID span_id
Parent Span ID parent_span_id
服务名 resource.attributes.service.name
日志正文 body.text
日志级别 severity_text
Logger 名 scope.name
Span 名 name
Span 属性 attributes.*

特别注意:这里查询的是 trace_id,不是很多 ECS/APM 教程里的 trace.id。不同接入路径使用的模板不同,
不要凭经验猜字段名,应该先查看实际 _source 和 mapping。

查看最新一条日志:

curl -H 'Content-Type: application/json' \
  'http://127.0.0.1:9200/logs-generic.otel-default/_search?pretty' \
  -d '{
    "size": 1,
    "sort": [{"@timestamp": "desc"}],
    "query": {"match_all": {}}
  }'

八、将 API 访问日志改成 RocketMQ 异步落库

8.1 改造前后的区别

改造前:

flowchart LR A["业务请求线程"] --> F["ApiAccessLogFilter"] F --> Feign["Feign 同步调用 Infra"] Feign --> DB[("访问日志表")] DB --> A

只要 Infra、网络或数据库慢,业务响应就会等待。

改造后:

flowchart LR A["业务请求线程"] --> F["ApiAccessLogFilter"] F --> P["RocketMQ asyncSend"] P --> A P -. "Broker 持久化" .-> MQ["api-access-log Topic"] MQ --> C["Infra Consumer"] C --> DB[("infra_api_access_log")]

8.2 添加 RocketMQ 依赖

Web Starter 将 RocketMQ 依赖标记为 optional,未使用 MQ 的应用不会被强制传递:

<dependency>
    <groupId>org.apache.rocketmq</groupId>
    <artifactId>rocketmq-spring-boot-starter</artifactId>
    <optional>true</optional>
</dependency>

真正发送或消费访问日志的服务显式引入该 Starter。

8.3 自动配置条件

只有同时满足以下条件才创建 Producer 和访问日志 Filter:

  • classpath 中存在 RocketMQTemplate
  • handao.access-log.enable=true
  • 配置了 rocketmq.name-server
  • 配置了 rocketmq.producer.group
@Bean
@ConditionalOnMissingBean
@ConditionalOnProperty(
        prefix = "handao.access-log",
        name = "enable",
        havingValue = "true",
        matchIfMissing = true)
@ConditionalOnProperty(
        prefix = "rocketmq",
        name = {"name-server", "producer.group"})
public ApiAccessLogProducer apiAccessLogProducer(
        RocketMQTemplate rocketMQTemplate,
        @Value("${handao.access-log.rocketmq.topic:api-access-log}") String topic) {
    return new ApiAccessLogProducer(rocketMQTemplate, topic);
}

这样可以避免未配置 RocketMQ 的轻量应用在启动时失败。

8.4 生产端异步发送

@Slf4j
public class ApiAccessLogProducer {

    private final RocketMQTemplate rocketMQTemplate;
    private final String topic;

    public void send(ApiAccessLogCreateReqDTO dto) {
        String message = null;
        try {
            message = JsonUtils.toJsonString(dto);
            String finalMessage = message;
            rocketMQTemplate.asyncSend(topic, message, new SendCallback() {
                @Override
                public void onSuccess(SendResult result) {
                    // 旁路日志成功时不再打印日志,避免日志放大
                }

                @Override
                public void onException(Throwable throwable) {
                    log.error("[send][访问日志异步发送失败,topic({}) message({})]",
                            topic, StrUtils.maxLength(finalMessage, 500), throwable);
                }
            });
        } catch (Throwable throwable) {
            log.error("[send][访问日志提交发送失败,topic({}) message({})]",
                    topic, StrUtils.maxLength(message, 500), throwable);
        }
    }
}

这里有两个设计点:

  1. 使用 asyncSend,业务线程不等待 Broker 返回;
  2. 访问日志属于旁路数据,发送失败只记录错误,不影响业务响应。

第二点是业务取舍。如果访问日志属于强合规数据、要求绝不丢失,就不能只使用 best-effort 异步发送,
应改为本地事务 Outbox、CDC 或其他可靠事件方案。

8.5 在生产端预生成日志 ID

RocketMQ 是至少一次投递。消费者处理成功但 ACK 丢失时,同一消息可能再次到达。如果数据库继续使用
自增 ID,每次消费都会生成新行,无法识别重复消息。

解决方法是在生产端生成稳定 ID,并把它放进消息:

// 生产端预生成主键,同一消息无论消费多少次,ID 都不变
accessLog.setId(IdUtil.getSnowflakeNextId());
accessLog.setTraceId(TracerUtils.getTraceId());

DTO 中把 ID 设为必填:

@NotNull(message = "日志编号不能为空")
private Long id;

8.6 消费者落库与重试

@Component
@ConditionalOnProperty(
        prefix = "handao.access-log",
        name = "enable",
        havingValue = "true",
        matchIfMissing = true)
@ConditionalOnProperty(prefix = "rocketmq", name = "name-server")
@RocketMQMessageListener(
        topic = "${handao.access-log.rocketmq.topic:api-access-log}",
        consumerGroup = "${handao.access-log.rocketmq.consumer-group:api-access-log-consumer}")
public class ApiAccessLogMessageConsumer implements RocketMQListener<String> {

    @Resource
    private ApiAccessLogService apiAccessLogService;

    @Override
    public void onMessage(String message) {
        try {
            ApiAccessLogCreateReqDTO dto =
                    JsonUtils.parseObject(message, ApiAccessLogCreateReqDTO.class);
            apiAccessLogService.createApiAccessLog(dto);
        } catch (RuntimeException exception) {
            log.error("[onMessage][API 访问日志落库失败,message({})]", message, exception);
            throw exception;
        }
    }
}

消费者异常必须继续抛出。吞掉异常会让 RocketMQ 误以为消费成功,失去自动重试和死信能力。

8.7 数据库幂等

重复消息第二次插入时会命中相同主键。确认数据库中确实已有该 ID 后,将其视为成功:

try {
    apiAccessLogMapper.insert(apiAccessLog);
} catch (DuplicateKeyException exception) {
    ApiAccessLogDO existing = TenantUtils.executeIgnore(
            () -> apiAccessLogMapper.selectById(apiAccessLog.getId()));
    if (existing == null) {
        // 可能是其他唯一约束冲突,不能错误地吞掉
        throw exception;
    }
    log.debug("[createApiAccessLog][访问日志({}) 已存在,忽略重复消息]",
            apiAccessLog.getId());
}

为什么捕获异常后还要查询一次?因为 DuplicateKeyException 不一定来自主键,也可能来自其他唯一索引。
只有相同访问日志 ID 已经存在时,才能确认这是同一消息的重复投递。

8.8 多租户上下文跨 MQ 传播

HTTP 请求中的租户 ID 通常保存在 ThreadLocal 中。消息进入 RocketMQ 后切换了线程和进程,ThreadLocal
不会自动存在。

生产 Hook 把租户 ID 写入消息属性:

public void sendMessageBefore(SendMessageContext context) {
    Long tenantId = TenantContextHolder.getTenantId();
    if (tenantId != null) {
        context.getMessage().putUserProperty("tenant-id", tenantId.toString());
    }
}

消费 Hook 在执行监听器之前恢复,之后必须清理:

public void consumeMessageBefore(ConsumeMessageContext context) {
    String tenantId = context.getMsgList().getFirst().getUserProperty("tenant-id");
    if (StrUtil.isNotEmpty(tenantId)) {
        TenantContextHolder.setTenantId(Long.parseLong(tenantId));
    }
}

public void consumeMessageAfter(ConsumeMessageContext context) {
    TenantContextHolder.clear();
}

如果不在 after 中清理,线程池复用时可能发生严重的跨租户数据污染。

8.9 配置项

rocketmq:
  name-server: 192.168.1.117:9876
  producer:
    group: ${spring.application.name}_PRODUCER

handao:
  access-log:
    enable: true
    rocketmq:
      topic: api-access-log
      consumer-group: api-access-log-consumer

Topic 和 Consumer Group 应保持稳定,不要把随机实例 ID 放入 Consumer Group,否则每个实例都会收到一份
消息,导致重复消费压力。

九、Docker Compose 部署中间件

下面是核心编排的精简版本:

services:
  elasticsearch:
    image: docker.elastic.co/elasticsearch/elasticsearch:9.5.2
    environment:
      - discovery.type=single-node
      - xpack.security.enabled=false
      - ES_JAVA_OPTS=-Xms1g -Xmx1g
    ports:
      - "9200:9200"
    volumes:
      - elasticsearch-data:/usr/share/elasticsearch/data

  kibana:
    image: docker.elastic.co/kibana/kibana:9.5.2
    environment:
      - ELASTICSEARCH_HOSTS=http://elasticsearch:9200
    ports:
      - "5601:5601"
    depends_on:
      elasticsearch:
        condition: service_healthy

  otel-collector:
    image: otel/opentelemetry-collector-contrib:0.159.0
    command: ["--config=/etc/otelcol-contrib/config.yaml"]
    volumes:
      - ./otel-collector-config.yaml:/etc/otelcol-contrib/config.yaml:ro
    ports:
      - "4317:4317"
      - "4318:4318"
      - "13133:13133"

  rocketmq-namesrv:
    image: apache/rocketmq:5.5.0
    command: sh mqnamesrv
    ports:
      - "9876:9876"

  rocketmq-broker:
    image: apache/rocketmq:5.5.0
    command: sh mqbroker -c /home/rocketmq/conf/broker.conf --enable-proxy
    environment:
      - NAMESRV_ADDR=rocketmq-namesrv:9876
    ports:
      - "8080:8080"
      - "8081:8081"
      - "10909:10909"
      - "10911:10911"
      - "10912:10912"

本地启动:

cd script/docker
docker compose up -d \
  elasticsearch kibana otel-collector rocketmq-namesrv rocketmq-broker

使用 Podman 时:

podman compose up -d \
  elasticsearch kibana otel-collector rocketmq-namesrv rocketmq-broker

9.1 RocketMQ 的 brokerIP1 陷阱

Broker 向 NameServer 注册的不只是名字,还包含客户端后续连接的 Broker 地址。

如果所有应用和 Broker 都在同一台机器,可以使用:

brokerIP1=127.0.0.1

如果应用在另一台机器,而 Broker 部署在 192.168.1.117,必须使用:

brokerIP1=192.168.1.117

否则会出现非常迷惑的现象:

  • NameServer 9876 可连接;
  • 应用能查到 Topic 路由;
  • 但路由返回 127.0.0.1:10911
  • 应用最终连接的是自己的本机 10911,发送失败。

因此远程排查 RocketMQ 时,不要只测试 9876,还要测试 10911,并确认 Broker 广播地址。

十、本地应用连接远程中间件

假设中间件统一位于:

MIDDLEWARE_HOST=192.168.1.117

推荐先检查端口:

for port in 3306 6379 8848 9876 10911 4318 13133 9200 5601; do
  nc -z -w 2 192.168.1.117 "$port" \
    && echo "OPEN $port" \
    || echo "CLOSED $port"
done

常用端口:

端口 服务
3306 MySQL
6379 Redis
8848 Nacos
9876 RocketMQ NameServer
10911 RocketMQ Broker
4317 Collector OTLP/gRPC
4318 Collector OTLP/HTTP
13133 Collector Health Check
9200 Elasticsearch
5601 Kibana

10.1 远端只有 ES,本机启动 Collector

复制一份临时配置,把 Docker DNS 名 elasticsearch 替换为远端 IP:

sed 's#http://elasticsearch:9200#http://192.168.1.117:9200#' \
  script/docker/otel-collector-config.yaml \
  > /tmp/handao-otel-collector.yaml

运行 Collector:

podman run --rm \
  --name handao-otel-collector \
  --network host \
  -v /tmp/handao-otel-collector.yaml:/etc/otelcol-contrib/config.yaml:ro,Z \
  otel/opentelemetry-collector-contrib:0.159.0 \
  --config=/etc/otelcol-contrib/config.yaml

验证:

curl http://127.0.0.1:13133

10.2 启动顺序

微服务模式推荐:

  1. Elasticsearch;
  2. OpenTelemetry Collector;
  3. RocketMQ NameServer 和 Broker;
  4. Nacos、MySQL、Redis;
  5. Infra;
  6. System;
  7. Gateway;
  8. 前端。

Infra 先启动,是因为它包含 API 访问日志消费者。System 即使先启动也不会影响业务启动,但早期产生的
访问日志需要等待消费者上线后才能落库。

十一、端到端验证

11.1 增加一条唯一的 SLF4J 验证日志

选择登录后前端稳定调用的接口:

GET /admin-api/system/auth/get-permission-info

在 Controller 中添加:

public CommonResult<AuthPermissionInfoRespVO> getPermissionInfo() {
    Long userId = getLoginUserId();
    log.info("[getPermissionInfo][OpenTelemetry 链路验证,userId({})]", userId);
    // ...
}

重启 System,登录或刷新前端页面。

11.2 获取 Trace ID

浏览器 Network 面板或 curl 查看响应头:

curl -i \
  -H 'Authorization: Bearer YOUR_TOKEN' \
  http://127.0.0.1:48081/admin-api/system/auth/get-permission-info

记录:

trace-id: 8a7c479d422fa7bc0fe1aade0d703a24

11.3 Kibana 查询

为下面两个 Data Stream 创建 Data View:

logs-generic.otel-*
traces-generic.otel-*

KQL:

trace_id : "8a7c479d422fa7bc0fe1aade0d703a24"

查验证日志:

body.text : "*OpenTelemetry 链路验证*"

11.4 直接使用 Elasticsearch API 查询

同时查询日志和 Trace:

curl -H 'Content-Type: application/json' \
  'http://192.168.1.117:9200/logs-generic.otel-default,traces-generic.otel-default/_search?pretty' \
  -d '{
    "size": 100,
    "sort": [{"@timestamp": "asc"}],
    "query": {
      "term": {
        "trace_id": "8a7c479d422fa7bc0fe1aade0d703a24"
      }
    }
  }'

先按日志内容反查 Trace ID:

curl -H 'Content-Type: application/json' \
  'http://192.168.1.117:9200/logs-generic.otel-default/_search?pretty' \
  -d '{
    "size": 20,
    "sort": [{"@timestamp": "desc"}],
    "query": {
      "match_phrase": {
        "body.text": "OpenTelemetry 链路验证"
      }
    }
  }'

一个成功的日志文档应同时包含:

{
  "trace_id": "8a7c479d422fa7bc0fe1aade0d703a24",
  "span_id": "12b91feec45048af",
  "severity_text": "INFO",
  "resource": {
    "attributes": {
      "service.name": "system-server"
    }
  },
  "scope": {
    "name": "com.example.system.AuthController"
  },
  "body": {
    "text": "[getPermissionInfo][OpenTelemetry 链路验证,userId(1)]"
  }
}

11.5 验证访问日志异步落库

接口调用后查询:

SELECT id,
       trace_id,
       application_name,
       request_method,
       request_url,
       duration,
       result_code,
       begin_time
FROM infra_api_access_log
ORDER BY id DESC
LIMIT 20;

表中的 trace_id 应与响应头和 ES 中的 trace_id 一致。

11.6 如何判断整个链路真正成功

不要只看“ES 有数据”。完整验收应该同时满足:

  • 响应头存在有效 trace-id
  • 跨域前端能读取 trace-id
  • Logback 控制台日志含相同 trace_id
  • logs-generic.otel-default 能按该 ID 查到业务日志;
  • traces-generic.otel-default 能查到 HTTP、JDBC、Feign 等 Span;
  • 多服务 Span 共享同一 Trace ID;
  • RocketMQ 消费 Span 能与生产端上下文关联;
  • infra_api_access_log.trace_id 一致;
  • 重复消费同一访问日志消息不会产生第二行数据库记录。

十二、自动化测试与配置验证

12.1 Trace ID 工具测试

@Test
void getTraceIdReturnsCurrentOpenTelemetryTraceId() {
    SpanContext context = SpanContext.create(
            "0123456789abcdef0123456789abcdef",
            "0123456789abcdef",
            TraceFlags.getSampled(),
            TraceState.getDefault());

    try (Scope ignored = Span.wrap(context).makeCurrent()) {
        assertEquals(
                "0123456789abcdef0123456789abcdef",
                TracerUtils.getTraceId());
    }
}

还应验证没有当前 Span 时返回空字符串。

12.2 CORS 测试

@Test
void corsFilterExposesTraceIdHeader() throws Exception {
    MockHttpServletRequest request =
            new MockHttpServletRequest("GET", "/admin-api/test");
    request.addHeader(HttpHeaders.ORIGIN, "http://localhost:3000");
    MockHttpServletResponse response = new MockHttpServletResponse();

    CorsFilter filter = new HandaoWebAutoConfiguration()
            .corsFilterBean().getFilter();

    filter.doFilter(request, response,
            (req, resp) -> ((HttpServletResponse) resp)
                    .setHeader("trace-id", "0123456789abcdef0123456789abcdef"));

    assertEquals("trace-id",
            response.getHeader(HttpHeaders.ACCESS_CONTROL_EXPOSE_HEADERS));
}

12.3 MQ 幂等测试

用同一个 DTO 连续调用两次:

apiAccessLogService.createApiAccessLog(createDTO);
apiAccessLogService.createApiAccessLog(createDTO);

assertEquals(1L, apiAccessLogMapper.selectCount());

12.4 构建验证

mvn -pl \
  handao-gateway,\
  handao-module-system/handao-module-system-server,\
  handao-module-infra/handao-module-infra-server,\
  handao-server \
  -am -DskipTests compile

12.5 Compose 与 Collector 验证

docker compose -f script/docker/docker-compose.yml config

otelcol-contrib validate \
  --config script/docker/otel-collector-config.yaml

两种验证解决不同问题:Compose 验证编排,Collector 验证组件配置。

十三、常见故障排查

13.1 响应没有 trace-id

检查顺序:

  1. JVM 命令行是否真的包含 -javaagent:/.../opentelemetry-javaagent.jar
  2. 启动日志是否打印 Agent 版本;
  3. Agent 文件路径是否存在;
  4. handao.tracer.enable 是否被关闭;
  5. TraceFilter 自动配置是否加载;
  6. 请求是否经过应用,而不是被前置代理直接返回。

快速检查进程:

ps -ef | grep opentelemetry-javaagent

13.2 浏览器 Network 能看到 trace-id,但 Axios 读不到

原因几乎一定是缺少:

Access-Control-Expose-Headers: trace-id

后端 CORS 配置增加 addExposedHeader("trace-id"),重启应用。

13.3 控制台有日志,但 ES 没有日志

检查:

  • otel.logs.exporter=otlp
  • Agent 是否加载;
  • Collector 4318 是否可达;
  • Collector 是否启用了 logs pipeline;
  • otlp Receiver、batch Processor、Exporter 是否都在 logs pipeline 中;
  • Collector 日志是否出现 401、403、404、429 或连接超时;
  • Elasticsearch /_otlp/v1/logs 是否可用。

13.4 Trace 有数据,但日志没有 trace_id

可能原因:

  • 日志发生在 HTTP Span 建立前,例如应用启动日志;
  • 日志发生在 Span 结束后的异步线程,且上下文未传播;
  • 手工创建 Span 后没有 makeCurrent()
  • 线程池被不兼容的包装器绕过 Agent;
  • 日志采集配置被关闭。

并不是所有日志都应该有 Trace ID。启动日志、定时任务全局日志没有当前 Trace 时,trace_id 为空是正常的。

13.5 Kibana 用 trace.id 查不到

检查实际 _source。本方案使用 Elastic 原生 OTLP Data Stream,字段是:

trace_id

不是:

trace.id

13.6 Collector 启动失败

典型原因是拿新版本文档配置旧版本 Collector,或反过来。必须使用和容器完全相同版本的二进制执行:

otelcol-contrib validate --config config.yaml

13.7 RocketMQ NameServer 能连接但消息发送失败

检查 Broker 返回的地址。远程 Broker 的 brokerIP1 不能是 127.0.0.1。同时检查防火墙是否开放:

9876
10909
10911
10912

13.8 访问日志出现重复数据

检查:

  • ID 是否在生产端生成;
  • 消息重试时是否复用了相同消息体;
  • DTO 到 DO 的转换是否保留 ID;
  • 数据库是否真的以该 ID 为主键;
  • 消费者是否错误地重新生成 ID。

13.9 访问日志不落库

检查:

  • handao.access-log.enable=true
  • rocketmq.name-serverrocketmq.producer.group 是否存在;
  • Topic 是否自动创建或已经预创建;
  • Infra Consumer Group 是否启动;
  • Broker 是否返回正确 IP;
  • 消息是否进入重试或死信队列;
  • 租户上下文是否恢复;
  • 数据库表和主键类型是否匹配 Snowflake Long。

十四、生产环境必须补齐的能力

开发环境跑通不等于可以直接上线。至少应处理以下事项。

14.1 Elasticsearch 安全

开发 Compose 中为了快速验证使用:

xpack.security.enabled=false

生产环境必须启用 TLS 和认证。Collector 可以集中持有 API Key:

exporters:
  otlphttp/elasticsearch:
    endpoint: https://elasticsearch.example.com:9200/_otlp
    headers:
      Authorization: "ApiKey ${env:ELASTICSEARCH_API_KEY}"
    tls:
      ca_file: /etc/otelcol/certs/ca.crt

API Key 不应写进 Git,也不应下发给所有业务应用。

14.2 Collector 持久队列

内存队列在 Collector 重启时会丢失。生产环境可以启用 file_storage

extensions:
  health_check:
  file_storage:
    directory: /var/lib/otelcol/storage

exporters:
  otlphttp/elasticsearch:
    sending_queue:
      enabled: true
      storage: file_storage

service:
  extensions: [health_check, file_storage]

对应目录必须挂载持久卷,并监控磁盘容量。

14.3 Collector 高可用

生产环境应运行多个 Collector Gateway 实例,通过 Service 或负载均衡器提供统一 OTLP 地址。

需要监控:

  • Receiver 接收速率;
  • Exporter 成功和失败数量;
  • Queue 容量和使用量;
  • Retry 次数;
  • Dropped spans/log records;
  • Collector 内存、CPU、GC;
  • Elasticsearch 429 和 ingest latency。

14.4 Trace 采样

开发环境通常 100% 采样:

OTEL_TRACES_SAMPLER=always_on

高流量生产环境可以使用比例采样:

OTEL_TRACES_SAMPLER=parentbased_traceidratio
OTEL_TRACES_SAMPLER_ARG=0.1

但头部采样无法提前知道请求是否会失败。若希望保留全部错误和慢请求,可在 Collector 中使用 tail
sampling。Tail sampling 必须看到一条 Trace 的所有 Span,因此通常要求相同 Trace 路由到同一 Collector,
并评估等待窗口带来的内存开销。

14.5 日志保留与 ILM

Trace、应用日志和访问日志的保留周期应分别设计:

数据 常见保留策略示例
Trace 7~30 天,按采样率和调用量调整
INFO 应用日志 15~30 天
ERROR 日志 30~90 天
API 访问审计日志 根据业务与合规要求,可能 180 天以上

不要直接修改 Elastic 内置的 OTel 模板。应通过 logs-otel@customtraces-otel@custom 等官方扩展点
定制生命周期和映射。

14.6 敏感数据与成本控制

禁止采集或输出:

  • 密码;
  • Access Token、Refresh Token;
  • Cookie、Authorization Header;
  • 身份证号、银行卡号等敏感信息;
  • 完整大请求体和大响应体;
  • 无限制 SQL 参数;
  • 完整异常对象的重复字符串副本。

API 访问日志应在进入 MQ 前脱敏和截断,否则超大消息不仅增加 Broker 压力,还可能超过 RocketMQ 消息
大小限制。现有脱敏逻辑如果解析失败后选择保留原字符串,属于 fail-open 策略;强合规场景应改成解析失败
直接丢弃请求体或仅记录摘要。

14.7 访问日志可靠性等级

本方案的语义是:

  • 请求线程不等待 Broker ACK;
  • Broker 接受消息后,依赖 RocketMQ 至少一次投递;
  • 消费端通过预生成主键幂等;
  • 发送失败只记录错误,不影响业务请求。

它适合大多数运营和排障访问日志,但不是严格“零丢失”审计。如果业务要求日志和业务事务原子一致,
推荐:

flowchart LR T["业务本地事务"] --> B[("业务表")] T --> O[("Outbox 事件表")] O --> CDC["CDC / 定时发布器"] CDC --> MQ["RocketMQ"] MQ --> A[("审计日志库")]

14.8 时间同步

跨服务 Trace 高度依赖时间。所有宿主机、容器节点应启用 NTP/chrony。时钟漂移会导致:

  • Span 时间线顺序混乱;
  • 负耗时或异常耗时;
  • 日志与 Trace 在时间窗口中看似不相关;
  • Kibana 查询遗漏。

14.9 版本升级策略

不要盲目使用浮动标签。推荐:

  1. 固定 Agent、Collector、Elasticsearch、Kibana、RocketMQ 版本;
  2. Agent JAR 固定 SHA-256;
  3. Elasticsearch 与 Kibana 保持相同版本;
  4. Collector 升级前运行 validate
  5. 阅读版本 Breaking Changes;
  6. 在测试环境回放真实 Trace/Log 流量;
  7. 检查 Data Stream、模板、ILM 和字段变化;
  8. 灰度升级 Collector,再升级 Agent;
  9. 保留快速回滚镜像。

十五、最终验收清单

上线前可以直接按下面的清单逐项确认。

应用

  • 所有服务加载相同基线的 OpenTelemetry Java Agent;
  • 每个服务有唯一且稳定的 service.name
  • 使用 tracecontext,baggage
  • Trace 和 Log 均导出到 Collector;
  • Logback Pattern 包含 trace_idspan_id
  • 响应 Header 返回有效 trace-id
  • CORS 暴露 trace-id
  • 前端能够读取并展示/上报 Trace ID。

Collector 与 Elasticsearch

  • Collector 配置通过目标版本 validate
  • 4318 和 13133 健康可达;
  • logs-generic.otel-default 有数据;
  • traces-generic.otel-default 有数据;
  • 同一请求的日志和 Span 使用相同 trace_id
  • Collector 配置队列、重试和内存限制;
  • 生产环境启用 TLS、API Key、持久队列;
  • 配置 ILM 和容量告警。

RocketMQ 访问日志

  • 生产端使用 asyncSend
  • Topic 和 Consumer Group 稳定;
  • Broker 广播地址对客户端可达;
  • 生产端生成访问日志 ID;
  • 消费端主键幂等;
  • 消费异常继续抛出;
  • 已验证重试与死信;
  • 租户 ID 能传播且消费后清理;
  • 请求参数已经脱敏和截断;
  • 明确 best-effort 与强审计可靠性边界。

十六、总结

这套方案最重要的不是“把日志写进 Elasticsearch”,而是建立清晰的职责边界:

  • OpenTelemetry Java Agent 负责自动埋点、上下文传播和 Logback 捕获;
  • 业务代码只依赖稳定的 OpenTelemetry API,避免绑定具体 APM 产品;
  • Collector 负责接收、保护、批处理、重试和统一导出;
  • Elasticsearch 9.x 原生 OTLP Endpoint 负责标准化存储;
  • Kibana 通过 trace_id 把 Trace 与 SLF4J 日志串起来;
  • RocketMQ 负责将 API 访问审计从业务请求线程解耦;
  • 生产端预生成主键、消费端检查重复主键,实现至少一次投递下的幂等;
  • Trace ID 同时出现在响应 Header、应用日志、Trace Data Stream 和访问日志表中,成为排障的统一入口。

最终,排查一次请求只需要一个动作:复制前端响应中的 trace-id,在 Kibana 中查询同值的 trace_id
从用户请求、跨服务调用、SQL、消息队列到业务日志,都可以沿着同一条链路还原。