流处理引擎 · 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(流名/存储实体名)、scopeIdbuilder(实体↔存储转换器)、processor(用哪个 StreamProcessor)、allowBootReshape(BanyanDB 专用)。

四种 StreamProcessor,各管一类数据:

图1

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

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

数据模型继承关系

图2

  • Metricsoap-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,存储主键)。子类众多:CountMetricsLongAvgMetricsDoubleAvgMetricsSumMetricsPercentileMetricsHistogramMetricsApdexMetricsRateMetricsCPMMetrics
  • Record.../analysis/record/Record.java:32):implements StorageData,只有 timeBucket,子类如 SegmentRecordLogRecordEventLongText
  • TopN.../analysis/topn/TopN.java:32):extends Record implements ComparableStorageData,增加 latencytraceIdentityIdtimestamp,按 latency 排序(:62)。

Worker 链如何编排

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

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

AbstractWorker · 基类

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

单条 metric 的完整 Worker DAG

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

图3

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

  1. StorageBuilderFactory 生成 StorageBuilder,拿 StorageDAO 生成 IMetricsDAO:207-220)。
  2. ModelRegistry 注册存储模型 Model,拿 DownSamplingConfigService/TTLStatusQuery 决定降采样和 TTL(:222-228)。
  3. 根据 @MetricsExtensionsupportDownSampling/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 · 集群连接管理

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

图5

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

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

接收侧:RemoteServiceHandler

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

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

IWorkerInstanceGetter/SetterIWorkerInstanceGetter.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)。
  • sessionCacheMetricsSessionCache :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 天),天级最长(半年)。查近期看细粒度,查历史看粗粒度,省存储。

降采样枚举 DownSamplingMinute/Hour/Day/Second/NoneMetrics 子类实现 toHour()/toDay() 做降采样克隆。


PersistenceTimer · 批量写入调度

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

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

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

图8

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

  1. 收集所有持久化 worker:TopNStreamProcessor + MetricsStreamProcessorgetPersistentWorkers:119-120)。
  2. prepareExecutorService(线程数 = getPrepareThreads():94-96)并发跑每个 worker 的 buildBatchRequestsprepare 阶段,产出 List<PrepareRequest>:135)——这里 worker 内部调 IMetricsDAO.prepareBatchInsert/Update 把对象转成 InsertRequest/UpdateRequest
  3. 把所有 PrepareRequest 交给 IBatchDAO.flush(prepareRequests):146execute 阶段)真正批量写入。
  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.javarecord/Record.javatopn/TopN.java
Worker 基类 AbstractWorker/PersistenceWorker oap-server/server-core/.../worker/AbstractWorker.javaanalysis/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.javacore/alarm/AlarmEntrance.java
Worker 注册 IWorkerInstanceGetter/Setter oap-server/server-core/.../remote/IWorkerInstanceGetter.javaIWorkerInstanceSetter.java

接下来

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

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