流处理引擎 · 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,各管一类数据:

四种数据形态对应的处理器
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 删除)。
数据模型继承关系

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)组装的链:

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

- 用共享的
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 节点的连接。

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)核心逻辑:

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 拿旧值合并,启动会很慢。
降采样:分钟 → 小时 → 天
指标不只存一个粒度,分钟级再聚合成小时、天。

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 阶段)。这样把高频小写入变成低频批量写,降低数据库压力。

start(:58)用 scheduleWithFixedDelay 起定时器(初始延迟 5s,周期 = moduleConfig.getPersistentPeriod(),:103),每周期跑 extractDataAndSave(:110):
- 收集所有持久化 worker:
TopNStreamProcessor+MetricsStreamProcessor的getPersistentWorkers(:119-120)。 - 用
prepareExecutorService(线程数 =getPrepareThreads(),:94-96)并发跑每个 worker 的buildBatchRequests(prepare 阶段,产出List<PrepareRequest>,:135)——这里 worker 内部调IMetricsDAO.prepareBatchInsert/Update把对象转成InsertRequest/UpdateRequest。 - 把所有
PrepareRequest交给IBatchDAO.flush(prepareRequests)(:146,execute 阶段)真正批量写入。 endOfFlush(:152)给后端清理机会。
IBatchDAO.flush 是异步返回 CompletableFuture<Void>(IBatchDAO.java:47),所以 prepare 与 execute 可以重叠。
告警与导出:在 Worker 链里的位置
MetricsStreamProcessor.minutePersistentWorker(:317)创建 AlarmNotifyWorker,作为 MetricsPersistentMinWorker 的下游(:320-323)。即分钟级指标持久化后,会调用 AlarmNotifyWorker.in(metrics)。

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-集群协调与动态配置。
完整数据流总图
把前面所有要素汇总:

关键文件速查
| 子系统 | 关键类 | 路径 |
|---|---|---|
| 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 是统一发货的铃铛,铃一响所有分厂一起发货。降采样是把成品再打包成"日报/周报"存到更长期的货架。