用 Elastic APM 实现分布式日志与链路追踪
本文假设你已经完成了 Elasticsearch 和 Kibana 的部署与基础配置。如果你还没有搭建 ES/Kibana,请先参考官方文档或其他教程完成基础设施准备。本文聚焦于在已有 ES/Kibana 的基础上,如何通过 Elastic APM Java Agent 以最低成本接入链路追踪与日志采集。
最近在项目可观测性的重构中,我尝试了 Elastic APM Java Agent 方案,发现它完美解决了下述问题:
- 一个 Agent 搞定自动注入链路追踪 + 直接发送应用日志到 ES
- 零代码侵入,不需要改任何业务代码
- 日志自动携带
trace.id、span.id,在 Kibana 中一键从 Trace 跳转到日志 - 链路很简单,应用 → APM Server → Elasticsearch
Elastic APM Java Agent 基于 ByteBuddy 字节码增强技术,在类加载时自动拦截 Spring MVC/WebFlux、JDBC、HTTP Client 等调用生成 Span。更重要的是,它内置了 Logging Appender Bridge,会自动桥接 SLF4J/Logback/Log4j2,将日志事件序列化为 ECS(Elastic Common Schema)格式,通过 APM 协议批量发送到 APM Server。这意味着你的日志不再是纯文本,而是结构化文档。
1. 安装与配置 APM Server
APM Server 是一个轻量的 Go 程序,负责接收 Agent 数据、校验、转换并写入 ES。
curl -L -O https://artifacts.elastic.co/downloads/apm-server/apm-server-9.4.1-x86_64.rpm
sudo rpm -vi apm-server-9.4.1-x86_64.rpm
编辑 /etc/apm-server/apm-server.yml:
apm-server:
host: "0.0.0.0:8201"
auth:
secret_token: "YOUR_SECRET_TOKEN"
output.elasticsearch:
hosts: ["https://127.0.0.1:9200"]
protocol: "https"
username: "elastic"
password: "${ES_PASSWORD}"
ssl.verification_mode: none
bulk_max_size: 2048 # 单次批量写入最大文档数
worker: 3 # 并发写入 worker 数
compression_level: 3 # 启用压缩减少网络开销
启动服务:
sudo systemctl enable apm-server
sudo systemctl start apm-server
2. 理解数据存储:Data Stream 与 ILM
Elastic APM 默认使用 Data Stream(数据流),它是专为时序数据设计的抽象层,底层由多个隐藏的 backing index 组成,写入时自动路由到最新的 backing index,查询时自动聚合所有 backing index,天然支持滚动更新,配合 ILM 实现自动 Hot→Warm→Cold→Delete 生命周期。
APM 预置的数据流:
| Data Stream | 用途 |
|---|---|
traces-apm.app-* |
应用链路追踪(Span/Transaction) |
logs-apm.app-* |
Agent 桥接的应用日志 |
metrics-apm.app.*-* |
应用指标(JVM、自定义指标) |
logs-apm.error-* |
错误事件 |
使用GET _data_stream/logs-apm.*,traces-apm.*,metrics-apm.*查询下APM相关的数据流。
配置 ILM 策略
APM Server 安装时会自动创建默认的 ILM 策略,但建议根据实际量级调整:
PUT _ilm/policy/logs-apm.app-custom
{
"policy": {
"phases": {
"hot": {
"min_age": "0ms",
"actions": {
"rollover": {
"max_primary_shard_size": "50gb", // 按大小滚动,而非仅按时间
"max_age": "7d"
},
"set_priority": { "priority": 100 }
}
},
"warm": {
"min_age": "7d",
"actions": {
"shrink": { "number_of_shards": 1 }, // 合并分片节省资源
"forcemerge": { "max_num_segments": 1 },
"set_priority": { "priority": 50 }
}
},
"cold": {
"min_age": "30d",
"actions": {
"searchable_snapshot": { // 可选:冷数据用快照存储
"snapshot_repository": "my_s3_repo"
}
}
},
"delete": {
"min_age": "90d",
"actions": { "delete": {} }
}
}
}
}
更新原来的索引模板以应用自定义生命周期函数。
3. 接入 Java Agent
从Maven仓库下载最新 Agent JAR。
-javaagent:/opt/agent/elastic-apm-agent-1.56.0.jar
-Delastic.apm.service_name=cdc-consumer-service
-Delastic.apm.secret_token=${APM_SECRET_TOKEN}
-Delastic.apm.server_url=http://192.168.0.201:8201
-Delastic.apm.environment=production
-Delastic.apm.application_packages=com.aj.cdc
-Delastic.apm.log_sending=true
-Delastic.apm.log_level=INFO # Agent 自身日志级别,排查问题时改为 DEBUG
-Delastic.apm.transaction_sample_rate=0.1 # 生产环境务必采样!全量采集会拖垮性能
-Delastic.apm.span_compression_enabled=true # 压缩重复 span,减少数据量
-Delastic.apm.capture_body=off # 避免采集请求体中的敏感信息
-Delastic.apm.sanitize_field_names=password,token,authorization # 字段脱敏
当 log_sending=true 时,Agent 会在 MDC(Mapped Diagnostic Context)中自动注入
trace.id 和 span.id。即使你仍然使用 Filebeat 采集文件日志(而非 Agent 直发),只要日志格式中包含这两个字段,Kibana 也能自动建立关联。
4. Kibana 观测效果
启动应用后,进入 Kibana → Observability → APM
整体看板:
服务拓扑与事务详情:

日志视图
解决 APIGateway 在 ElasticAPM 上显示 GET unknown route 的问题
<dependency>
<groupId>co.elastic.apm</groupId>
<artifactId>apm-agent-api</artifactId>
<version>1.56.0</version>
</dependency>
import co.elastic.apm.api.ElasticApm;
import co.elastic.apm.api.Transaction;
import org.springframework.cloud.gateway.filter.GatewayFilterChain;
import org.springframework.cloud.gateway.filter.GlobalFilter;
import org.springframework.cloud.gateway.route.Route;
import org.springframework.cloud.gateway.support.ServerWebExchangeUtils;
import org.springframework.core.Ordered;
import org.springframework.http.HttpHeaders;
import org.springframework.http.server.reactive.ServerHttpRequest;
import org.springframework.stereotype.Component;
import org.springframework.web.server.ServerWebExchange;
import reactor.core.publisher.Mono;
@Component
public class TracePropagationFilter implements GlobalFilter, Ordered {
@Override
public Mono<Void> filter(ServerWebExchange exchange, GatewayFilterChain chain) {
// 1. 获取当前 Elastic APM 事务
Transaction transaction = ElasticApm.currentTransaction();
// 2. 获取当前 traceparent 字符串(格式:00-<traceid>-<spanid>-01)
String traceparent = transaction.getTraceId();
String spanId = transaction.getId();
String traceparentHeader = "00-" + traceparent + "-" + spanId + "-01";
// 解决 APIGateway 在 ElasticAPM 上显示 GET unknown route 的问题
Route route = exchange.getAttribute(ServerWebExchangeUtils.GATEWAY_ROUTE_ATTR);
if (route != null) {
String method = exchange.getRequest().getMethod().name();
transaction.setName(method + " " + route.getId());
}
// 3. 将 traceparent 添加到转发的请求头中
ServerHttpRequest mutatedRequest = exchange.getRequest().mutate()
.header("traceparent", traceparentHeader)
.header("elastic-apm-traceparent", traceparentHeader)
.build();
ServerWebExchange mutatedExchange = exchange.mutate()
.request(mutatedRequest)
.build();
return chain.filter(mutatedExchange);
}
@Override
public int getOrder() {
return Ordered.HIGHEST_PRECEDENCE;
}
}
效果如下:

更多推荐

所有评论(0)