流处理引擎 · Worker 链

这一篇解决什么
Source 进来后,指标怎么从"一条一条的原始事件"变成"存进数据库的聚合点"?答案是 StreamProcessor + Worker 链:L1 内存聚合攒批、按 hash 分片到集群节点、L2 缓存做增量合并、降采样到小时天、最后批量落盘。这篇拆解四种 StreamProcessor、单条 metric 的完整 Worker DAG、L1/L2 缓存策略、集群分片机制、PersistenceTimer 的批量写入。这是 OAP 性能设计最核心的一篇。


四种 StreamProcessor

StreamProcessor 接口极简(oap-server/server-core/.../analysis/StreamProcessor.java:21):只有 void in(STREAM stream)。它是所有流处理器的统一入口。

@Stream 注解(oap-server/server-core/.../analysis/Stream.java:41)标注在数据类上,声明 name(流名/存储实体名)、scopeId、builder(实体↔存储转换器)、processor(用哪个 StreamProcessor)、allowBootReshape(BanyanDB 专用)。

四种 StreamProcessor,各管一类数据:

图1

四种数据形态对应的处理器

  • MetricsStreamProcessor(.../analysis/worker/MetricsStreamProcessor.java:65)——最复杂最重要。处理 Metrics(可聚合统计量,如 QPS、延迟均值)。单例。in(Metrics)(:141)是热路径,按 metrics.getClass() 从 entryWorkers 找 MetricsAggregateWorker 转发。支持运行时热增删(removeMetric :378、suspendDispatch :478),用于 MAL/LAL runtime rule 热更新。
  • RecordStreamProcessor(RecordStreamProcessor.java:49)——处理 Record(原始记录,如 trace segment、log),不聚合、不降采样,直接持久化,受 TTL 控制(recordDataTTL :56,in 里做过期判断 :69-73)。
  • TopNStreamProcessor(TopNStreamProcessor.java:51)——处理 TopN(Top N 慢记录,如最慢的 N 条 SQL)。维护固定大小窗口,低频持久化(topNWorkerReportCycle :59,topSize :62)。
  • NoneStreamProcessor(NoneStreamProcessor.java:46)——处理 NoneStream(UI 操作配置类实体,同步写、按 Record 模式 TTL 删除)。

数据模型继承关系

图2

  • Metrics(oap-server/server-core/.../analysis/metrics/Metrics.java:41):extends StreamData implements StorageData, ToJson。核心抽象方法:combine(Metrics)(:68,合并同 ID 同 timeBucket)、calculate()(:73,计算最终值)、toHour()/toDay()(:80/83,降采样克隆)、id0()(:167,存储主键)。子类众多:CountMetrics、LongAvgMetrics、DoubleAvgMetrics、SumMetrics、PercentileMetrics、HistogramMetrics、ApdexMetrics、RateMetrics、CPMMetrics。
  • Record(.../analysis/record/Record.java:32):implements StorageData,只有 timeBucket,子类如 SegmentRecord、LogRecord、Event、LongText。
  • TopN(.../analysis/topn/TopN.java:32):extends Record implements ComparableStorageData,增加 latency、traceId、entityId、timestamp,按 latency 排序(:62)。

Worker 链如何编排

SkyWalking 的"流处理 DAG"不是用配置文件画的,而是 MetricsStreamProcessor.create 用代码组装出来的固定形态链,每个 metric 类一条链。

为什么是代码组装而非配置
因为 worker 链的形态(要不要降采样、要不要告警、要不要导出)取决于 metric 的 @MetricsExtension 属性,运行时按属性分支组装。这比配置文件灵活,代价是新增链形态要改代码。

AbstractWorker · 基类

AbstractWorker(oap-server/server-core/.../worker/AbstractWorker.java:30):所有 Worker 基类。持有 ModuleDefineHolder(:33,便于跨模块找 Service),核心方法 in(INPUT)(:42)是数据入口。

单条 metric 的完整 Worker DAG

MetricsStreamProcessor.create(MetricsStreamProcessor.java:202-307)组装的链:

图3

组装步骤(MetricsStreamProcessor.java:202-307):

  1. 拿 StorageBuilderFactory 生成 StorageBuilder,拿 StorageDAO 生成 IMetricsDAO(:207-220)。
  2. 拿 ModelRegistry 注册存储模型 Model,拿 DownSamplingConfigService/TTLStatusQuery 决定降采样和 TTL(:222-228)。
  3. 根据 @MetricsExtension 的 supportDownSampling/supportUpdate/timeRelativeID(:235-246)决定链形态。
  4. 若支持降采样:建 hour/day 的 MetricsPersistentWorker(:329),再用 MetricsTransWorker(:267)把 minute → hour/day 串起来。
  5. 建 minute 级 MetricsPersistentMinWorker(:309),内部挂 AlarmNotifyWorker(告警)、ExportMetricsWorker(导出)、transWorker(降采样转发)三个 next worker。
  6. 把 minute worker 注册到 IWorkerInstanceSetter(:298-301),名字为 stream.getName()+"_rec"——这是集群内 RPC 的接收端 worker 名。
  7. 建 MetricsRemoteWorker(:303,把数据按 hash 发到目标节点)。
  8. 建 MetricsAggregateWorker(:304,L1 入口),next worker = remote worker。
  9. 把 aggregateWorker 放入 entryWorkers(:306),in(Metrics) 据此路由。

L1 聚合:MetricsAggregateWorker

L1 是内存里的第一道聚合,把高频小写入攒成低频批量。

MetricsAggregateWorker(.../analysis/worker/MetricsAggregateWorker.java:57):

图4

  • 用共享的 BatchQueue(:98,名为 METRICS_L1_AGGREGATION,CPU 核数线程、自适应分区、20k buffer)。
  • in(Metrics)(:146)produce 到队列;L1Handler.consume(:207)调 onWork(:152)→ mergeDataCache.accept 合并相同 ID+timeBucket 的指标 → flush()(:178)按 l1FlushPeriod(默认 500ms,:129)周期性把合并结果喂给 next worker。
  • OAL metric 权重 1.0,MAL meter 权重 0.05(:141,因 MAL 每分钟才来一批,多个类型可共享分区)。

为什么要 L1
同一个服务一秒可能收到几百条 trace,每条都要更新 resp_time 指标。如果每条都走"读旧值→合并→写回",数据库扛不住。L1 先在内存把这一秒内的几百条合并成一个"增量",500ms 才往后传一次。后面 L2 再做"读旧值→叠加→写回"。

为什么 L1 用 BatchQueue 多线程消费,而不是单线程边收边合并?因为合并本身是 CPU 活(要算 ID 去重、调 combine),单线程会成为吞吐瓶颈。BatchQueue 把"入队"和"合并"解耦:接收线程只做极快的入队,多个消费线程并行合并,合并完再汇总。代价���合并顺序不确定,但 metric 的 combine 必须满足交换律和结合律——这也是为什么所有算子的 combine 都是累加式(计数相加、桶计数相加),不能用减法或除法这类顺序敏感的操作。


集群分片:MetricsRemoteWorker

L1 合并完的数据,要按 hash 决定去哪个 OAP 节点做持久化。

MetricsRemoteWorker(.../analysis/worker/MetricsRemoteWorker.java:33):in(Metrics)(:44)调 remoteSender.send(workerName, metrics, Selector.HashCode),把数据按 hash 路由到某个 OAP 节点。

为什么必须 hash 分片
同一个服务的指标可能被不同 OAP 节点收到。如果每个节点都存全量,指标会重复或冲突(combine 算两次)。HashCodeSelector(.../remote/selector/HashCodeSelector.java:25)用 Math.abs(streamData.remoteHashCode() % size)——保证同一 entity 的 metric 总是落到同一节点,保证聚合正确。

三种 Selector(.../remote/selector/):

Selector 行为 谁用
HashCodeSelector Math.abs(hash % size),同 hash 落同节点 Metrics(保证聚合正确)
RollingSelector 轮询 —
ForeverFirstSelector 永远选第一个 —

RemoteClientManager · 集群连接管理

RemoteClientManager(oap-server/server-core/.../remote/client/RemoteClientManager.java:59)管理本节点到集群内所有其他 OAP 节点的连接。

图5

  • implements Service, ClusterWatcher(:59)。
  • start()(:95)起定时线程,每 10 秒 refresh()。
  • refresh()(:105)→ 从 ClusterModule 拿 ClusterNodesQuery(:107)查询集群节点列表 → refresh(instanceList)(:112)。
  • refresh(List<RemoteInstance>)(:118)去重、排序、比对,若变了就 reBuildRemoteClients(:146)重建连接,并通知 ServerStatusService.rebalancedCluster(:147,触发所有 MetricsPersistentWorker 的 onClusterRebalanced)。
  • reBuildRemoteClients(:202)做 diff:不变的 Unchanged、新出现的 Create(本机建 SelfRemoteClient,远端建 GRPCRemoteClient 并 connect(),:238-246)、消失的 Close(:255-261)。最后排序保证顺序稳定(:252,为 RollingSelector 一致性)。

不要直接 find ClusterModule
RemoteClientManager.refresh 内部才 find(ClusterModule.NAME)。业务代码应通过 CoreModule 的 RemoteClientManager 间接使用。详见 07-集群协调与动态配置。

接收侧:RemoteServiceHandler

RemoteServiceHandler(oap-server/server-core/.../remote/RemoteServiceHandler.java:52)是接收侧 gRPC 实现:

  • call(StreamObserver)(:105)返回一个 StreamObserver<RemoteMessage>。
  • onNext(RemoteMessage)(:118):取 nextWorkerName,用 IWorkerInstanceGetter.get(nextWorkerName)(:125)找到 RemoteHandleWorker,newStreamDataInstance() 反序列化(:130-135),streamData.deserialize(remoteData)(:137),最后 nextWorker.in(streamData)(:143)——即把远端发来的数据喂给本节点对应的 MetricsPersistentMinWorker(名字就是 streamName+"_rec",见上面组装第 6 步)。

IWorkerInstanceGetter/Setter(IWorkerInstanceGetter.java:26 / IWorkerInstanceSetter.java:29):worker 名 → 实例的注册与查找表。put 在建链时调用(MetricsStreamProcessor.java:298),get 在收 RPC 时调用(RemoteServiceHandler.java:125)。

这就完成了集群内分发:发送端按 hash 选节点,接收端按 worker 名找到本地 worker 接续处理。


L2 持久化:MetricsPersistentWorker

L2 是真正落盘前的最后一道缓存,做"读旧值→合并→写回"。

PersistenceWorker(.../analysis/worker/PersistenceWorker.java:38)抽象基类:持 ReadWriteSafeCache(:40),onWork(:50)写缓存,buildBatchRequests(:64)和 endOfRound(:58)由 PersistenceTimer 驱动。

MetricsPersistentWorker(.../analysis/worker/MetricsPersistentWorker.java:53)核心逻辑:

图6

  • in(Metrics)(:172):先查 TTL 过期(:173),再 super.onWork 入缓存。
  • buildBatchRequests(:188):受 persistentMod 控制(minute 级每轮都执行,降采样级每 4 轮执行一次,见 :165)。
  • buildBatchRequestsUnconditionally(:201):读缓存 → 对每批 loadFromStorage(:249,从 DB 读已有值到 sessionCache)→ combine 合并 → calculate → 生成 prepareBatchUpdate/Insert(:271-285)→ 喂给 alarm/export/trans worker(:303)。
  • sessionCache(MetricsSessionCache :115)缓存最近写过的 metric,避免每次都读 DB;endOfRound(:357)清过期。

启动优化:ServerStatusWatcher
MetricsPersistentWorker 实现 ServerStatusWatcher(:53):onServerBooted(:428)记录启动时间,onClusterRebalanced(:434)记录 rebalance 时间,requireInitialization(:385)用这个时间避免启动后无谓的 DB 回读——这是性能优化关键。OAP 刚启动时,缓存里没数据,如果不优化,每个 metric 都要回读 DB 拿旧值合并,启动会很慢。


降采样:分钟 → 小时 → 天

指标不只存一个粒度,分钟级再聚合成小时、天。

图7

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

降采样枚举 DownSampling:Minute/Hour/Day/Second/None。Metrics 子类实现 toHour()/toDay() 做降采样克隆。


PersistenceTimer · 批量写入调度

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

PersistenceTimer(oap-server/server-core/.../storage/PersistenceTimer.java:45,单例枚举):

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

图8

start(:58)用 scheduleWithFixedDelay 起定时器(初始延迟 5s,周期 = moduleConfig.getPersistentPeriod(),:103),每周期跑 extractDataAndSave(:110):

  1. 收集所有持久化 worker:TopNStreamProcessor + MetricsStreamProcessor 的 getPersistentWorkers(:119-120)。
  2. 用 prepareExecutorService(线程数 = getPrepareThreads(),:94-96)并发跑每个 worker 的 buildBatchRequests(prepare 阶段,产出 List<PrepareRequest>,:135)——这里 worker 内部调 IMetricsDAO.prepareBatchInsert/Update 把对象转成 InsertRequest/UpdateRequest。
  3. 把所有 PrepareRequest 交给 IBatchDAO.flush(prepareRequests)(:146,execute 阶段)真正批量写入。
  4. endOfFlush(:152)给后端清理机会。

IBatchDAO.flush 是异步返回 CompletableFuture<Void>(IBatchDAO.java:47),所以 prepare 与 execute 可以重叠。


告警与导出:在 Worker 链里的位置

MetricsStreamProcessor.minutePersistentWorker(:317)创建 AlarmNotifyWorker,作为 MetricsPersistentMinWorker 的下游(:320-323)。即分钟级指标持久化后,会调用 AlarmNotifyWorker.in(metrics)。

图9

  • AlarmNotifyWorker(.../analysis/worker/AlarmNotifyWorker.java:30):in(Metrics) 若 metrics 实现 WithMetadata,调 entrance.forward(metrics)。
  • AlarmEntrance(.../core/alarm/AlarmEntrance.java:24):forward 先检查 has(AlarmModule.NAME)(没装告警模块就跳过),再懒加载 MetricsNotify 服务,调 metricsNotify.notify(metrics)。

AlarmNotifyWorker 只接在分钟级
AlarmNotifyWorker 只接在分钟级 worker 后,不接小时/天。这与 RunningRule 注释一致:“only minute dimensionality metrics are expected to process”。告警基于分钟级指标,不会用降采样后的数据。详见 06-查询层与告警。

ExportMetricsWorker 把指标外发给外部系统(OTel collector/Kafka),订阅驱动。详见 07-集群协调与动态配置。


完整数据流总图

把前面所有要素汇总:

图10


关键文件速查

子系统 关键类 路径
StreamProcessor StreamProcessor 接口 oap-server/server-core/.../analysis/StreamProcessor.java
Stream 注解 同目录 analysis/Stream.java
四种 Processor Metrics/Record/TopN/None oap-server/server-core/.../analysis/worker/{Metrics,Record,TopN,None}StreamProcessor.java
数据模型 Metrics/Record/TopN oap-server/server-core/.../analysis/metrics/Metrics.java、record/Record.java、topn/TopN.java
Worker 基类 AbstractWorker/PersistenceWorker oap-server/server-core/.../worker/AbstractWorker.java、analysis/worker/PersistenceWorker.java
L1 MetricsAggregateWorker oap-server/server-core/.../analysis/worker/MetricsAggregateWorker.java
Remote MetricsRemoteWorker oap-server/server-core/.../analysis/worker/MetricsRemoteWorker.java
L2 MetricsPersistentWorker oap-server/server-core/.../analysis/worker/MetricsPersistentWorker.java
集群 RemoteClientManager oap-server/server-core/.../remote/client/RemoteClientManager.java
RemoteSenderService oap-server/server-core/.../remote/RemoteSenderService.java
RemoteServiceHandler oap-server/server-core/.../remote/RemoteServiceHandler.java
HashCodeSelector oap-server/server-core/.../remote/selector/HashCodeSelector.java
批量写入 PersistenceTimer oap-server/server-core/.../storage/PersistenceTimer.java
告警入口 AlarmNotifyWorker / AlarmEntrance oap-server/server-core/.../analysis/worker/AlarmNotifyWorker.java、core/alarm/AlarmEntrance.java
Worker 注册 IWorkerInstanceGetter/Setter oap-server/server-core/.../remote/IWorkerInstanceGetter.java、IWorkerInstanceSetter.java

接下来

数据算好准备落盘了,具体怎么存?BanyanDB/ES/JDBC 三套实现差异在哪?详见 05-存储层。

心智模型
把流处理链想成一条"加工流水线 + 仓库调度":L1(AggregateWorker)是车间里的暂存货架,攒一会儿再往后送;RemoteWorker 是调度员,按零件编号(hash)决定送到哪个分厂;L2(PersistentWorker)是分厂的成品库,先核对旧账(multiGet)再合并入账;PersistenceTimer 是统一发货的铃铛,铃一响所有分厂一起发货。降采样是把成品再打包成"日报/周报"存到更长期的货架。