SkyWalking 源码阅读
SkyWalking 的价值在于"无侵入":Java Agent 通过字节码增强,业务代码零改动就能拿到完整链路。本文拆解从 Agent 采集到 OAP 分析的完整链路。
整体架构与模块
┌──────────────────────────── 应用侧 ────────────────────────────┐
│ 业务代码(零改动) │
│ ↑ 字节码增强 │
│ Agent(skywalking-agent.jar) │
│ ├─ 插件体系(采集点) │
│ ├─ Segment 构建(链路数据) │
│ └─ gRPC 客户端(上报) │
└────────────────────────────┬───────────────────────────────────┘
│ gRPC
┌────────────────────────────▼───────────────────────────────────┐
│ OAP Server(分析引擎) │
│ ├─ Receiver(接收数据) │
│ ├─ 分析管道(聚合指标) │
│ ├─ 告警引擎 │
│ └─ 存储(ES / MySQL / H2) │
└────────────────────────────────────────────────────────────────┘
│
▼
Web UI(拓扑/追踪/指标/告警)Agent 启动:Premain
Agent 入口
java
// skywalking-agent 的入口
public class SkyWalkingAgent {
public static void premain(String agentArgs, Instrumentation instrumentation) {
// 1. 初始化配置
Config.Initialize(agentArgs);
// 2. 安装字节码增强
AgentInstaller.install(instrumentation);
}
}启动时序
JVM 启动(-javaagent)
│
├─ premain() 执行
│
├─ 1. 解析配置(agent.service_name、采样率、上报地址)
│
├─ 2. 加载插件定义(plugins/ 目录下的 jar)
│
├─ 3. 字节码增强:拦截关键类的关键方法
│
└─ 4. 启动 gRPC 上报客户端字节码增强:插件体系
增强原理
不是修改源码,而是 JVM 加载类时在内存中改写字节码:
类加载 → Agent 拦截 → 在方法前后插入采集代码 → 类被加载
例如增强 HttpClient.execute():
调用前:创建 Span、记录开始时间
调用后:记录耗时、响应状态、上报插件模型
java
// 插件定义:增强哪个类的哪个方法
public class HttpClientPlugin extends AbstractClassEnhancePluginDefine {
@Override
protected ClassMatch enhanceClass() {
// 匹配要增强的类
return NameMatch.byName("org.apache.http.client.HttpClient");
}
@Override
protected InstanceMethodsInterceptPoint[] getInstanceMethodsInterceptPoints() {
// 匹配要增强的方法 + 拦截器
return new InstanceMethodsInterceptPoint[]{
new InstanceMethodsInterceptPoint() {
@Override
public ElementMatcher<MethodDescription> getMethodsMatcher() {
return named("execute"); // 增强 execute 方法
}
@Override
public String getMethodsInterceptor() {
return "HttpClientExecuteInterceptor"; // 拦截器
}
}
};
}
}拦截器
java
// 拦截器:在方法前后插入逻辑
public class HttpClientExecuteInterceptor implements InstanceMethodsAroundInterceptor {
@Override
public void beforeMethod(EnhancedInstance objInst, Method method,
Object[] allArguments, Class<?>[] argumentsTypes,
MethodInterceptResult result) {
// 创建子 Span
AbstractSpan span = ContextManager.createExitSpan(
"/http/client/execute",
remotePeer, contextCarrier);
// 注入追踪上下文到请求头(传播)
injectContext(carrier, allArguments);
}
@Override
public Object afterMethod(EnhancedInstance objInst, Method method,
Object[] allArguments, Class<?>[] argumentsTypes,
Object ret) {
// 结束 Span(记录耗时、状态)
ContextManager.stopSpan();
return ret;
}
}常见插件
插件按技术栈分类:
Web:Tomcat、Spring MVC、Spring WebFlux
RPC:Dubbo、gRPC、Feign、HTTP Client
DB:JDBC、MyBatis、Redis、MongoDB
消息:Kafka、RabbitMQ、RocketMQ数据模型:Context / Segment / Span
Context:请求上下文
java
// 每个请求一个 Context(ThreadLocal 保存)
public class Context {
private final String traceId; // 链路 ID
private List<AbstractSpan> spanStack; // Span 栈
private Segment segment; // 当前 Segment
}Segment:一次进程内的追踪片段
java
// 服务 A 内的所有 Span 组成一个 Segment
public class Segment {
private String traceSegmentId; // Segment ID
private String serviceId; // 服务 ID
private String serviceInstanceId; // 实例 ID
private List<Span> spans; // 本进程内 Span
}Span 数据结构
java
public abstract class AbstractSpan {
private int spanId; // Span ID
private int parentSpanId; // 父 Span ID
private long startTime; // 开始时间(微秒)
private long endTime; // 结束时间
private String operationName; // 操作名
private String peer; // 对端地址
private String spanType; // Entry / Exit / Local
}Span 类型:
Entry:服务入口(接收请求)
Exit:服务出口(调用下游)
Local:本地操作(方法调用、异步任务)三种 Span
服务A(segment-A)
├─ Entry Span:GET /api/order/1001 ← 入口
├─ Exit Span: 调用 order-service ← 出口(传播上下文)
│
服务B(segment-B)
├─ Entry Span:接收请求(提取上下文)
└─ Exit Span:查询数据库跨进程传播:ContextCarrier
传播机制
java
// 调用下游前:把上下文写入请求头
ContextCarrier carrier = new ContextCarrier();
ContextManager.inject(carrier); // traceId + segmentId + spanId → 请求头
// HTTP 头:
// sw8: <traceId>-<segmentId>-<spanId>-...
// 下游收到请求:提取上下文
ContextCarrier carrier = new ContextCarrier();
ContextManager.extract(carrier); // 从请求头还原
ContextManager.continued(carrier); // 续接父 Span传播过程
服务A:创建 Exit Span → inject(carrier) → 请求头携带
│
▼ HTTP 请求
服务B:extract(carrier) → 创建 Entry Span(parentSpanId = 上游 spanId)gRPC 上报
上报协议
proto
// SkyWalking 上报协议(proto)
service TraceSegmentReportService {
// 上报 Segment(批量)
rpc collect (stream SegmentObject) returns (Commands) {}
}
message SegmentObject {
string traceSegmentId = 1;
repeated SpanObject spans = 2;
string service = 3;
string serviceInstance = 4;
}Agent 端上报流程
Segment 完成(进程内请求结束)
│
├─ 1. Segment 序列化为 protobuf
│
├─ 2. 放入队列(异步、批量)
│
├─ 3. gRPC 客户端定时批量发送(stream)
│
└─ 4. 失败重连(退避重试)采样控制
不是所有请求都上报:
采样率配置(agent.sample_n_per_3_secs)
达到阈值 → 标记 sampled=false → 丢弃 Span 数据(减少开销)OAP 分析引擎
接收链路
OAP 启动
│
├─ gRPC Server(接收 Trace 数据)
│
├─ Receiver 反序列化 → 数据进入分析管道
│
├─ 分析管道(拓扑/指标/慢请求)
│
└─ 持久化 → UI 查询分析管道
java
// OAP 中数据的处理模型(流式)
// Trace 数据 → 多个 Analysis
// ├─ TopologyAnalysis:服务调用关系 → 拓扑图
// ├─ EndpointAnalysis:端点指标(QPS、响应时间、成功率)
// ├─ ServiceAnalysis:服务指标
// └─ TraceAnalysis:慢请求/错误追踪输入:Segment(span 列表)
↓
按 service/endpoint 聚合 → 指标(直方图、计数器)
↓
按 traceId 关联 → Trace 视图
↓
拓扑分析 → 服务间调用关系
↓
写入存储(ES 索引)指标类型
指标(Metrics):
service_p99 / service_p95 / service_cpm(每分钟调用数)
endpoint_success_rate / endpoint_response_time
......
Trace:
traceId 检索、慢请求定位存储选型
| 存储 | 适用 | 特点 |
|---|---|---|
| H2 | 演示/开发 | 零依赖、内存 |
| Elasticsearch | 生产标配 | 高吞吐、支持复杂查询、可扩展 |
| MySQL | 中小规模 | 易维护、容量有限 |
| PostgreSQL | 中型 | 支持 JSON 查询 |
| BanyanDB | 大型 | SkyWalking 自研时序库 |
ES 存储结构(索引):
segment:Trace 数据
metrics:聚合指标
alarm_record:告警记录完整链路回顾
1. 请求进入服务A
├─ 字节码增强拦截 Controller 方法
├─ 创建 Entry Span(traceId 生成)
└─ Context 存入 ThreadLocal
2. 服务A 调用服务B(HTTP)
├─ 增强拦截 HTTP Client
├─ 创建 Exit Span + inject 上下文到请求头
└─ 方法返回后结束 Span
3. 服务B 接收请求
├─ 增强拦截 Controller
├─ extract 上下文 + 创建 Entry Span
└─ ...
4. Segment 生成
├─ 请求结束 → ContextManager.stop()
└─ Segment 序列化 → 队列 → gRPC
5. OAP 处理
├─ 接收 Segment
├─ 分析管道聚合指标
└─ 存储(ES)→ UI 展示常见问题
- 增强不生效? 检查插件是否支持该框架版本、类是否被 Agent 加载、采样配置。
- traceId 对不上? 跨进程传播失败(异步线程未传递 Context),需排查线程池场景。
- 上报延迟高? gRPC 批量发送,调大批量大小/上报间隔,或检查网络。
- 存储压力大? 调低采样率、配置数据保留期(TTL)、分片优化。