整体架构与端到端数据流

这一篇解决什么
在拆细节之前,先用一条 trace 从产生到被看到的完整旅程,把 OAP 内部所有子系统串一遍。读完应该能回答:数据从哪来、中间被谁加工、最后落到哪、怎么被查出来。后面每一篇都是这条链上某一环的放大。


分层视图

OAP 后端是一条分层管道,每层只做一件事,层间用 Service 接口解耦。

图1

六层的分工是理解后续所有内容的前提:

  • 接收层把字节流解成 proto 对象,不碰业务语义。
  • 解析层把 proto 对象切成 Source——语义化的数据事件,比如"这是一次服务调用"。
  • 规则层把 OAL/MAL/LAL 脚本编译成 Dispatcher,决定对每类 Source 算哪些指标。
  • 聚合层用 Worker 链做内存聚合、降采样、落盘。
  • 存储层把对象翻译成具体数据库的写入请求。
  • 查询层按用户请求读存储,必要时用 MQE 做二次计算。

为什么要分这么细?因为每层解决的问题性质不同:接收要扛并发,解析要懂业务语义,聚合要算数学,存储要懂引擎。混在一起做,任何一个改动都会牵连其他层。分层后,换存储引擎不动解析逻辑,加新指标不动接收代码。


端到端时序:一条 trace 的旅程

跟一条 HTTP trace 走完整条链。假设用户调了 order-service,它又调了 inventory-service

图2

时序图里每一步背后都有具体机制,下面拆开。


数据模型:Source → Metrics → StorageData

整条链流转三种数据形态,搞清边界就理解了为什么这么分层。

图3

三种核心数据形态

  • Source:语义化的事件,带 scope()(作用域 ID)。由 Analyzer 构造,是"发生了什么"的抽象。比如 ServiceServiceRelationEndpointLogServiceInstanceJVMMemory。定义在 server-core/.../source/
  • Metrics:可聚合的统计量。combine() 合并同 ID 同时间桶、calculate() 算最终值、toHour()/toDay() 降采样克隆。由 Dispatcher 从 Source 转换得到,是"算出了什么"。比如 ServiceAvgRespTimePercentileMetrics。定义在 server-core/.../analysis/metrics/
  • Record:不可变的原始记录,只追加不聚合,受 TTL 控制。比如 SegmentRecord(trace)、LogRecordEvent。定义在 server-core/.../analysis/record/

这三种形态对应三种写入语义:Source 是中转不落盘,Metrics 要"同时间桶反复合并更新",Record 是"来一条 append 一条"。这种区分直接决定了存储层要有三套不同的 DAO(IMetricsDAO 有 update,IRecordDAO 没有),而不是一套通吃。


Scope 体系:数据的"归属维度"

每个 Source 都有一个 scope() 整数 ID,这是 SkyWalking 里很关键的概念。

什么是 Scope
Scope 回答"这条数据描述的是哪个层级的实体"。Service=1ServiceInstance=2Endpoint=3ServiceRelation=4。OAL 脚本里 from(Service.latency).longAvg()Service 就是 scope。Scope 让指标规则和数据来源解耦——同一个 Source 可以触发多条 OAL 规则,产出多个指标。

Scope ID 集中注册在 oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/source/DefaultScopeDefine.java:31,靠 @ScopeDeclaration 注解扫描进来。常量从 SERVICE=1GEN_AI_*,区间 [0,10000) 留给 SkyWalking 官方,扩展从 10000 起。这个分区设计是为了避免官方和第三方扩展的 scope ID 撞车——第三方加新 source 从 10000 往上排,不会和官方未来新增的冲突。

Scope 还按 catalog 归类(service/serviceInstance/endpoint/relation),告警和查询用它分组。见 DefaultScopeDefine.java:166-182


集群分片:为什么需要 RemoteClientManager

OAP 可以多节点部署。问题在于:同一个服务的指标可能被不同节点收到,如果每个节点都存全量,指标会重复或冲突。

SkyWalking 用按 hash 分片解决:同一实体的指标总是落到同一个 OAP 节点做持久化。

图4

  • 节点 A/B 收到数据后,在 L1 内存聚合完,由 MetricsRemoteWorkerHashCodeSelectorMath.abs(remoteHashCode % size))算出该数据应该去哪个节点。
  • 通过 gRPC 把数据推给目标节点 C。
  • 节点 C 的 RemoteServiceHandler 收到后,按 worker 名找到本地的 MetricsPersistentMinWorker 继续处理落盘。

为什么必须分片而不是各存各的?因为 metric 是"同 (entity, timeBucket) 要合并"的——同一服务同一分钟的指标分散到多个节点各存一份,就没法合并,查出来的就是错的。分片保证同一 entity 的同一时间桶只在一个节点聚合,结果才正确。

核心类在 oap-server/server-core/.../remote/

  • RemoteClientManager.java:59 每 10 秒刷新集群节点列表,维护到各节点的 gRPC 连接。
  • RemoteSenderService.java:39 发送侧,按 selector 选 client。
  • RemoteServiceHandler.java:52 接收侧,按 worker 名路由到本地 worker。

不要直接 find ClusterModule
业务代码不该直接 moduleManager.find(ClusterModule.NAME) 拿集群服务。应该走 CoreModuleRemoteClientManager——它是对外门面,内部再去找 ClusterModule。这样换集群实现(zookeeper/etcd/k8s)对业务透明。详见 07-集群协调与动态配置。


降采样:分钟、小时、天

指标不只存一个粒度。SkyWalking 把分钟级指标再聚合成小时、天级,这就是降采样(DownSampling)。

图5

  • 分钟级保留时间短(比如 7 天),小时级长一些(30 天),天级最长(半年)。查近期看细粒度,查历史看粗粒度,省存储。
  • MetricsStreamProcessor.create 在建链时按 DownSamplingConfigService 决定要不要建 hour/day worker(MetricsStreamProcessor.java:248,257)。
  • MetricsTransWorker 把上一级聚合结果喂给下一级(MetricsStreamProcessor.java:267)。

降采样不是简单地"把分钟数据相加除以 60"。不同算子降采样方式不同:Counter 类(count/sum)直接相加,Avg 类要重新用总 summation 除以总 count,Percentile 类要把各节点的桶计数合并后重算分位数。所以每个 Metrics 子类都要实现自己的 toHour()/toDay(),把"如何降采样"这个语义钉死在算子上。降采样枚举在 DownSamplingMinute/Hour/Day/Second/None


批量写入:PersistenceTimer

写入不是来一条写一条,而是攒着批量冲。PersistenceTimer 是写入侧的总调度。

马桶式写入
每个 Worker 在内存攒增量,PersistenceTimerpersistentPeriod(默认几秒)触发一次:多线程把所有 Worker 的内存对象"翻译"成数据库写入请求(prepare 阶段),再一把 flush 进去(execute 阶段)。高频小写入变成低频批量写,降低数据库压力。

  • PersistenceTimer.java:45 单例,start 起定时器。
  • extractDataAndSave:110)每周期收集所有持久化 Worker,并发 buildBatchRequests,再 IBatchDAO.flush
  • IBatchDAO.flush 返回 CompletableFuture,prepare 和 execute 可以重叠。

为什么要分 prepare 和 execute 两阶段?因为 prepare 是 CPU 密集(要把 metrics 对象序列化成存储格式、算 docId、合并旧值),execute 是 IO 密集(真正写库)。两阶段分开后,可以用多线程并发 prepare(吃满 CPU),再把 prepare 出的请求批量 execute(吃满 IO 带宽),两种资源都不浪费。


查询侧:读 + 二次计算

查询不是简单把存的数据原样取出来。用户可能要算 service_resp_time / service_cpm,或在 Grafana 里写 PromQL。所以查询侧也要一套表达式引擎,就是 MQE。

图6

  • MQE 用 ANTLR 解析表达式,用 Visitor 模式递归求值,不做代码生成(区别于 OAL/MAL/LAL 的编译方式),因为查询高频短命,生成 class 的开销不划算。
  • 叶子节点 visitMetric 才真正读存储,通过 MetricsQueryService 调底层 IMetricsQueryDAO
  • 告警规则也复用 MQE 语法,但数据来源是内存滑动窗口而非存储。详见 06-查询层与告警。

一张图收束所有概念

把前面所有要素汇总成一张关系图,作为本篇收尾,也是后续各篇的目录索引。

图7


接下来怎么读

全局观建立完毕。后面每一篇都是这张图里某一块的放大:

  • 想懂骨架怎么拼:01-模块系统与启动流程
  • 想懂数据怎么进来:02-数据采集层-Receiver
  • 想懂规则怎么编译:03-分析引擎与四大DSL
  • 想懂聚合怎么算:04-流处理引擎-Worker链
  • 想懂怎么落盘:05-存储层
  • 想懂怎么查、怎么告警:06-查询层与告警
  • 想懂集群、配置、自监控:07-集群协调与动态配置