本文假设你已经完成了 Elasticsearch 和 Kibana 的部署与基础配置。如果你还没有搭建 ES/Kibana,请先参考官方文档或其他教程完成基础设施准备。本文聚焦于在已有 ES/Kibana 的基础上,如何通过 Elastic APM Java Agent 以最低成本接入链路追踪与日志采集。

最近在项目可观测性的重构中,我尝试了 Elastic APM Java Agent 方案,发现它完美解决了下述问题:

  • 一个 Agent 搞定自动注入链路追踪 + 直接发送应用日志到 ES
  • 零代码侵入,不需要改任何业务代码
  • 日志自动携带 trace.idspan.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;
    }
}

效果如下:
在这里插入图片描述

在这里插入图片描述

Logo

有“AI”的1024 = 2048,欢迎大家加入2048 AI社区

更多推荐