集群协调与动态配置

这一篇解决什么
OAP 多节点怎么互相发现、怎么不重启改配置、怎么监控自己、怎么把数据外发、AI 基线怎么接、管理面怎么用。这篇收尾所有外围支撑子系统:集群协调、动态配置、Telemetry 自监控、Exporter、AI Pipeline、server-admin(含运行时规则热更新)、health-checker、server-tools、server-testing。读完能回答:OAP 集群怎么搭、配置怎么热更新、规则怎么不停机改。


集群协调 · server-cluster-plugin

解决什么问题

多个 OAP 节点要按 hash 分片转发数据,所以得"互相发现 + 维护 gRPC 连接"。集群协调插件就干这件事——底层可以是 zk/etcd/consul/nacos 注册中心,也可以是 k8s 直接查 pod。

什么是集群协调?一句话:让一组对等的 OAP 节点各自知道"集群里还有谁、地址多少",这样数据才能跨节点转发。OAP 的流处理是 hash 分片的,同一份数据要稳定落到同一个节点处理,节点之间互相建 gRPC 长连接按 hash 转发。

ClusterModule 定义了什么 Service

ClusterModule(oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/cluster/ClusterModule.java:23)只暴露三个 Service,见 services()(:32):

Service 方法 职责
ClusterRegister registerRemote(RemoteInstance) 节点启动时把自己注册到集群
ClusterNodesQuery queryRemoteNodes() 返回当前集群所有 RemoteInstance 列表
ClusterCoordinator(ClusterCoordinator.java:27) 抽象类,同时实现 ClusterRegister + ClusterNodesQuery + ClusterWatcherRegister,三合一基类,维护 List<ClusterWatcher>,节点列表变化时 notifyWatchers 回调(:40)

为什么把三个接口揉成一个抽象类?因为任意一种注册中心实现都必然要同时回答这三件事:我要注册、我要查列表、我要通知别人列表变了。拆成三个独立 provider 没意义,合在一起让插件作者一次实现全部,减少样板代码。不这么做的话,每个插件要写三个 provider 类,且三者共享的 List<ClusterWatcher> 状态没地方放。

RemoteInstance 是什么
RemoteInstance(RemoteInstance.java:25)本质是 Address(host+port+self 标志)的薄包装,Comparable 按 address 排序(:39)。Address.isSelf() 区分"自己"和"其他 OAP 节点",决定建本地 SelfRemoteClient 还是远端 GRPCRemoteClient。为什么要排序?后面 RemoteClientManager 要靠稳定顺序做 hash 分片——顺序一变,数据落点就变,流处理会乱。

各插件实现

所有插件继承 ClusterCoordinator,注册一个 XxxCoordinator 同时作为三个 Service 的实现。各 Provider 在 prepare()registerServiceImplementation(ClusterRegister.class, coordinator) 等。

图1

  • standalone(StandaloneManager.java:30):纯内存,registerRemote(:35)只存自己,queryRemoteNodes(:41)只返回自己。单机模式,不连任何注册中心。
  • zookeeper(ZookeeperCoordinator.java:48):用 Curator 的 ServiceDiscovery<RemoteInstance>,节点注册为 remote 路径下一个 ServiceInstance,payload 是 RemoteInstance(registerRemote :68)。通过 ServiceCache 监听节点变化,cacheChanged(:149)事件触发 notifyWatchers(:152)。
  • etcd(EtcdCoordinator.java:55):节点信息以 JSON 写入 etcd,key 为 serviceName/host_port,带 30s lease + keepAlive 续约(registerRemote :144 起 lease,:156 起 keepAlive)。queryRemoteNodes(:93)用前缀查询读所有 key。Watch 客户端监听 PUT/DELETE 事件,EtcdEventListener.onNext(:212)触发 notifyWatchers(:221)。
  • consul:用 ImmutableRegistration 注册服务,ServiceHealthCache 监听健康实例变化。
  • nacos:通过 Nacos naming 服务注册 + 订阅。
  • kubernetes(KubernetesCoordinator.java:52):不主动注册,而是通过 KubernetesLabelSelectorEndpointGroup(基于 label selector)从 Kubernetes API 列出 collector pod,构造 RemoteInstance 列表。port 取自 ConfigService.getGRPCPort()(createEndpointGroup :73)。这是"无注册中心"的发现方式——k8s 本身就是注册中心。

etcd 为什么用 lease + keepAlive 而不是直接写 key?因为节点崩溃时要让注册信息自动失效。lease 30s 到期没人续约,etcd 自动删 key,其他节点下次查询就看不到这个死节点。不这么做的话,死掉的 OAP 会永久留在集群列表里,数据一直往一个不存在的地址转发。

每个 Coordinator 都通过 TelemetryModuleMetricsCreator 创建 HealthCheckMetrics,在查询/注册失败时标记 unHealth(如 EtcdCoordinator.java:198createHealthCheckerGauge("cluster_etcd", ...)),并经 OAPNodeChecker.isHealth() 做节点健康校验(EtcdCoordinator.java:114)。这些 health_check_* 指标后面会被 health-checker 聚合成"健康分"。

与 RemoteClientManager 的关系

不要直接 find ClusterModule
业务代码不该直接 moduleManager.find(ClusterModule.NAME) 拿集群服务。应该走 CoreModuleRemoteClientManager——它才是稳定的对外门面,内部再去找 ClusterModule。这样换集群实现(zookeeper/etcd/k8s)对业务透明。

RemoteClientManager(oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/client/RemoteClientManager.java:59)实现 ServiceClusterWatcher 两个接口(:59 implements Service, ClusterWatcher):

  1. start()(:95)起定时任务,每 10 秒 refresh()(:99scheduleWithFixedDelay,参数 1, 10, SECONDS)。
  2. refresh()(:105)通过 moduleDefineHolder.find(ClusterModule.NAME).provider().getService(ClusterNodesQuery.class) 拿集群节点列表(:107)——这正是上面提示的关键:业务代码通过 CoreModule 拿 RemoteClientManager,RemoteClientManager 内部才 find ClusterModule。
  3. refresh(List<RemoteInstance>)(:118,synchronized)去重(distinct :179)、排序(Collections.sort :134)、与现有 client 比对(compare :142);若变化则 reBuildRemoteClients(:146)。
  4. reBuildRemoteClients(:202):对每个地址,self 用 SelfRemoteClient(:239),否则建 GRPCRemoteClientconnect()(:243-244);消失的 client 关闭(:261)。最后排序保证稳定顺序(:252 注释 //for stable ordering for rolling selector)。
  5. 同时实现 ClusterWatcher.onClusterNodesChanged(:278),所以支持两种刷新:定时轮询 + 事件推送(watcher 模式)。

什么是 Watcher 机制?就是"被观察者状态一变,主动通知所有观察者"。ClusterCoordinator 维护一个 List<ClusterWatcher>(:28),节点列表变化时 notifyWatchers 遍历回调每个 watcher 的 onClusterNodesChanged(:45)。RemoteClientManager 注册自己为 watcher,所以注册中心一变,它立刻被通知重建连接,不用等 10 秒轮询。两条腿走路,事件快、轮询兜底。

refresh 还会调 ServerStatusService.rebalancedCluster()(:147-150)通知集群拓扑变化,让 OAP 进入 rebalanced 状态(此期间某些写入行为会调整)。

详见 04-流处理引擎-Worker链 的 RemoteClientManager 节。


动态配置 · server-configuration

解决什么问题

让 OAP 不重启就能从外部配置中心(Apollo/Nacos/ZK/etcd/Consul/K8s ConfigMap/gRPC)拉取并热更新某些配置项,比如采样率、告警阈值、log analyzer 规则开关。

什么是动态配置中心?一个独立于 OAP 的存储(Apollo/Nacos 等),运维改了里面的值,OAP 运行时感知到并生效,不用重启。传统配置写在 application.yml,改完要重启进程;动态配置把"易变的运营参数"抽出来,让 OAP 边跑边读。

抽象

  • ConfigurationModule(oap-server/server-configuration/configuration-api/src/main/java/org/apache/skywalking/oap/server/configuration/api/ConfigurationModule.java:30)只暴露一个 Service:DynamicConfigurationService(:39)。
  • DynamicConfigurationService(DynamicConfigurationService.java:24)只有一个方法 registerConfigChangeWatcher(ConfigChangeWatcher)——模块启动时把自己的 watcher 注册进去。
  • AbstractConfigurationProvider(:25 起)是所有 provider 基类。prepare() 调子类 initConfigReader() 得到一个 ConfigWatcherRegister(:40),注册为 DynamicConfigurationService 实现(:41);notifyAfterCompleted()configWatcherRegister.start()(:53)。

ConfigChangeWatcher / ConfigSync 机制

ConfigChangeWatcher(ConfigChangeWatcher.java:30)抽象类,字段 module / provider / itemName / watchType(SINGLE|GROUP)(:31-34),抽象方法 notify(ConfigChangeEvent)(:48)和 value()(:53)。事件类型 ADD/MODIFY/DELETE(:72)。GroupConfigChangeWatcher 用于一组相关配置项(如多条告警规则)。

ConfigWatcherRegister(ConfigWatcherRegister.java:31)抽象基类,notifySingleValue(:35)是核心逻辑:比较新值与 watcher.value(),仅在变化时回调 notify,DELETE 时传 null(:41)。WatcherHolder 的 key 格式为 module.provider.itemName(:125-128)——这就是配置项命名约定,各模块按这个三元组注册自己的 watcher。

为什么要"比较再回调"?因为配置中心的 watch 事件可能重复触发(同一个值也会回调),直接通知 watcher 会让他反复重载配置,代价大且可能闪断。比一下 value(),没变就跳过。

两种同步风格:

图2

  • FetchingConfigWatcherRegister(FetchingConfigWatcherRegister.java:36):定时轮询。start()(:79)用虚拟线程定时器每 syncPeriod(默认 60s,构造器 :46)调 configSync()(:99),子类实现 readConfig(keys)(:155)/readGroupConfig(keys)(:157)拉配置,再逐项 notifySingleValue(:120)。isStarted 标志(:42)保证启动后不能再注册新 watcher(:55-56IllegalStateException)——启动后再加 watcher 会被轮询漏掉,行为不确定,干脆禁止。
  • ListeningConfigWatcherRegister(ListeningConfigWatcherRegister.java:27):事件监听。registerConfigChangeWatcher(:30)直接调子类 startListening(holder, callback)(:31),配置中心变化时回调 onSingleValueChanged(:33)/onGroupValuesChanged(:38)。start() 是空的(:46)——没有定时器要起。

各 provider 实现

Provider 基类 机制
apollo ConfigWatcherRegister Apollo 监听
consul Consul KV watch
etcd etcd KV watch
nacos Nacos config listener
zookeeper ZK node watch
k8s-configmap(ConfigurationConfigmapInformer) Fetching fabric8 informer 列 ConfigMap(按 labelSelector 过滤),定时读 configMapData()
grpc(GRPCConfigWatcherRegister) Fetching 见下

grpc-configuration-sync · OAP 作为配置客户端

重要澄清
grpc-configuration-sync 模块只是客户端(OAP 作为 gRPC client 去拉配置),不是服务端。

GRPCConfigWatcherRegister.java:35 继承 FetchingConfigWatcherRegister,构造器(:41)用 ConfigurationServiceGrpc.newBlockingStub 建到远端的 stub。readConfig(:52)调远端 call()(:60),用 UUID 做变更检测:若返回 UUID 与本地一致则返回 Optional.empty() 表示无变化(:62-65)。UUID 一致就跳过,省掉反序列化和逐项比对的代价。

proto 定义在 oap-server/server-configuration/grpc-configuration-sync/src/main/proto/configuration-service.proto,定义了 ConfigurationService.call / callGroup

服务端实现不在本仓库。全仓搜 ConfigurationServiceImplBase 只命中测试代码 oap-server/server-configuration/grpc-configuration-sync/src/test/java/org/apache/skywalking/oap/server/configuration/grpc/MockGRPCConfigService.java,生产代码无命中。因此 gRPC 配置服务端需要单独部署(可能是 SkyWalking 社区的独立 config-server 项目)。未确认该服务端是否在 SkyWalking 其它仓库——本仓库源码无法证实。


Telemetry · OAP 自监控

解决什么问题

让 OAP 把自己的运行时指标(内存、线程、各模块处理量、集群健康)以 Prometheus 格式暴露,供外部 Prometheus 抓取,实现"监控监控系统本身"。

抽象与实现

  • TelemetryModule(oap-server/server-telemetry/telemetry-api/src/main/java/org/apache/skywalking/oap/server/telemetry/TelemetryModule.java:28)暴露 MetricsCreatorMetricsCollector 两个 Service(:37-39)。
  • MetricsCreator 是工厂接口,创建 createCounter/createGauge/createHistogramMetric/createHealthCheckerGaugeHEALTH_METRIC_PREFIX = "health_check_"(MetricsCreator.java:31),健康检查指标名都加这个前缀。createHealthCheckerGauge(:53)默认实现就是把名字拼上前缀再建 gauge(:55),isHealthCheckerMetrics(:65)用前缀判断指标是否属于健康检查,extractModuleName(:73)去前缀还原模块名。
  • none 实现(NoneTelemetryProvider/MetricsCreatorNoop)做无操作占位——不想要自监控时配 none,所有 create 方法返回空实现,业务代码不用改。
  • PrometheusTelemetryProvider(oap-server/server-telemetry/telemetry-prometheus/src/main/java/org/apache/skywalking/oap/server/telemetry/prometheus/PrometheusTelemetryProvider.java:34):prepare()(:63)注册 PrometheusMetricsCreator/PrometheusMetricsCollector(:64-65),启动 HttpServer(:67),DefaultExports.initialize()(:72)暴露 JVM hotspot 指标。
  • HTTP 暴露:HttpServerHandler(oap-server/server-telemetry/telemetry-prometheus/src/main/java/org/apache/skywalking/oap/server/telemetry/prometheus/httpserver/HttpServerHandler.java:51 持有 CollectorRegistry.defaultRegistry)基于 Netty,请求到来时用 io.prometheus.client.exporter.common.TextFormat.write004(:68)把 CollectorRegistry.defaultRegistry 的所有指标序列化成 Prometheus text format 返回。

集群插件的 HealthCheckMetrics
集群插件里的 HealthCheckMetrics(如 cluster_zookeeper)也通过 Telemetry 的 MetricsCreator.createHealthCheckerGauge 创建(见 ZookeeperCoordinator.java:132),这些 health_check_* 指标会被 health-checker 子系统消费。


Exporter · 数据外发

解决什么问题

把 OAP 内部计算/聚合后的指标和 trace/log,按订阅关系重新导出到外部系统(如 OTel collector、Kafka),让 SkyWalking 数据流出到第三方链路。

Exporter vs Telemetry
Telemetry 暴露的是"OAP 自己有多忙"(自监控,面向 Prometheus),Exporter 暴露的是"OAP 算出来的业务指标/trace"(业务数据外发,面向 OTel collector/Kafka)。前者自监控,后者数据外发。且指标导出是订阅驱动——外部系统告诉 OAP"我要 metric X 的增量",OAP 才发。

模块与实现

  • ExporterModule(oap-server/server-core/.../exporter/ExporterModule.java:20)暴露三个 Service:MetricValuesExportServiceTraceExportServiceLogExportService。这是"消费方契约定义在 core、实现在 exporter 模块"的典型模式——core 定接口,业务代码调接口,具体发到 gRPC 还是 Kafka 由 provider 决定。
  • ExporterProvider(oap-server/exporter/src/main/java/org/apache/skywalking/oap/server/exporter/provider/ExporterProvider.java:34):prepare()(:66)注册三个导出器:GRPCMetricsExporter(实现 MetricValuesExportService,:70)、KafkaTraceExporter(实现 TraceExportService,:71)、KafkaLogExporter(实现 LogExportService,:72)。start()(:76)按开关分别 start(),notifyAfterCompleted()(:89)先拉一次订阅列表。

GRPCMetricsExporter · 订阅模式

GRPCMetricsExporter(oap-server/exporter/src/main/java/org/apache/skywalking/oap/server/exporter/provider/grpc/GRPCMetricsExporter.java:60 左右):

  • 连接远端 MetricExportService gRPC 服务,本地起 BatchQueue 缓冲(:82BatchQueueManager.create)。
  • 订阅模式:导出不是无脑全发。export(event)(:99)里,订阅列表为空且事件是 INCREMENT 时兜底全发(:103);否则只发匹配订阅的指标(:106-111 比对 metricName + eventType)。fetchSubscriptionList()(:126)硬编码 30s 周期(FETCH_SUBSCRIPTION_PERIOD = 30_000,:65)拉订阅——用 lastFetchTimestamp + ReentrantLock 做节流,避免每次 export 都打 RPC。
  • Kafka trace/log 导出器走 Kafka producer 把 trace/log 推到 Kafka topic。

为什么订阅驱动?如果 OAP 把所有指标无脑推给外部,外部不想要的就全是噪声流量。订阅让外部系统声明"只要这些",OAP 只发这些,带宽和外部处理压力都降下来。


AI Pipeline · AI 能力接入层

解决什么问题

把 OAP 与外部 AI 服务对接,提供"指标预测基线(用于异常检测告警)"和"HTTP URI 模式归一化识别"两类 AI 能力的接入层。

模块定义

  • AIPipelineModule(oap-server/ai-pipeline/src/main/java/org/apache/skywalking/oap/server/ai/pipeline/AIPipelineModule.java:24)只暴露一个 Service:BaselineQueryService(:33)。
  • AIPipelineProvider(AIPipelineProvider.java:32):prepare()(:61)注册 BaselineQueryServiceImpl(连接外部 baseline server);start()(:69)构造 HttpUriRecognitionService 并注入 EndpointNameGroupService.startHttpUriRecognitionSvr(:74-76)。
  • requiredModules()(:84)只依赖 CoreModule

BaselineQueryService · 预测基线

  • BaselineQueryService(.../services/BaselineQueryService.java:18)接口方法:querySupportedMetricsqueryPredictMetrics(serviceMetrics, start, end)queryPredictMetricsFromCache(serviceName, timeBucketHour)
  • BaselineQueryServiceImpl(BaselineQueryServiceImpl.java:51)通过 gRPC stub AlarmBaselineServiceGrpc(:52)调用外部 baseline 服务(包名 org.apache.skywalking.apm.baseline.v3)。构造器(:55)用 Guava Cache 缓存预测结果,1 小时过期(expireAfterAccess(1, HOURS),:57)。未配置地址则 stub 为 null(:59-61),方法返回空并告警(:70-72 等多处 stub == null 判空)。
  • PredictServiceMetrics 预测值结构,含 PredictSingleValue(value, upperValue, lowerValue)(BaselineQueryServiceImpl.java:173-177):预测基线值 + 上下界区间,用于判断实际值是否越界(异常检测)。
  • queryPredictMetricsFromCache(:95)的逻辑:缓存没有就按"过去 24h ~ 未来 24h"窗口(:114-117)向外部服务批量查预测,塞回缓存(:135-141),再按 (timeBucket, serviceName) 取(:144)。批量预取而不是逐次查,降低对外部 AI 服务的 QPS。

HttpUriRecognitionService · URI 模式识别

HttpUriRecognitionService(HttpUriRecognitionService.java:37 实现 HttpUriRecognition)通过 gRPC 调用外部 URI 识别服务,fetchAllPatterns(service)(:57)拉取归一化后的 URI 模式。这是把 /api/user/12345/api/user/99999 归并成 /api/user/:id 的能力,由外部 AI 模型完成,OAP 调用并把结果用于 endpoint 分组。

为什么要归一化 URI?不归一化的话,每个动态 ID 都是一个独立 endpoint,指标维度爆炸,没法做有意义的聚合分析。归一化后同一类接口合并成一条曲线,异常检测才有意义。SkyWalking 自己不做归一化(需要训练模型),外包给 AI 服务。

与 MQE baseline 操作的关系

已从源码确认的关键链路:

图3

  • MQEVisitorBase.java:42-43 import 了 AIPipelineModuleBaselineQueryService
  • MQEVisitorBase.java:78-91 持有 BaselineQueryService,通过 moduleManager.find(AIPipelineModule.NAME).provider().getService(BaselineQueryService.class) 懒加载(getBaselineQueryService :87)。懒加载是因为 MQE 模块不一定依赖 AI Pipeline,没装 AI Pipeline 时不能在构造期就 find,否则启动报错。
  • MQEVisitorBase.java:502 visitBaselineOP —— MQE 表达式中的 baseline() 操作符调用此方法,进而调 queryBaseline(:557)/queryLabeledBaseline(:582),从 BaselineQueryService.queryPredictMetricsFromCache 拿预测基线。

baseline 与异常检测
MQE 支持在 UI 查询里写 metric + baseline 之类运算。baseline() 操作符向 AI Pipeline 背后的外部预测服务请求"这个指标在此时此刻应该是什么值、正常波动区间多大",然后与实际值比较——超出上下界就可能是异常。AI Pipeline 模块本身不做预测,只做"OAP ↔ 外部 AI 服务"的 gRPC 桥接和缓存。

未确认
baseline server 和 URI recognition server 的具体实现(AI 模型)不在本仓库,本仓库只含客户端。仓内 ai-pipeline/src/test/java/BaselineQueryServer.java 是测试用的 mock server。


server-admin · 管理面

解决什么问题

提供一组默认关闭、需手动开启的管理 REST/gRPC 端口(默认 HTTP 17128 / gRPC 17129,区别于面向 agent 的 11800),承载 DSL 调试、运行时规则热更新、状态检查、UI 模版管理等运维操作。

admin-server · 共享管理主机

AdminServerModule(oap-server/server-admin/admin-server/src/main/java/org/apache/skywalking/oap/server/admin/server/module/AdminServerModule.java:40)通过 SW_ADMIN_SERVER=default 开启。services()(:48)暴露 HTTPHandlerRegisterGRPCHandlerRegisterAdminClusterChannelManager(admin-internal gRPC 基础设施,17129 端口)。注释(:51-55)明确把 admin-internal gRPC 与 agent/cluster gRPC 11800 分开,避免特权管理 RPC 与 agent 遥测共享爆炸半径——一个 agent 打爆 gRPC 端口不会拖垮 runtime-rule 的热更新通道。

AdminServerModuleProvider宿主,本身不提供业务端点,各 feature 模块在 start() 阶段 find(AdminServerModule.NAME).provider().getService(HTTPHandlerRegister.class).addHandler(...) 注册自己的 handler。

安全提示
admin 端口无内置鉴权(源码注释 :34-38 反复强调),运维必须用 IP 白名单 + 鉴权反代保护,绝不暴露公网。注释直接指向 docs/en/setup/backend/admin-api/readme.md 的安全说明。

dsl-debugging · DSL 调试

在线调试 MAL/LAL/OAL DSL 的 REST 面。RuntimeOalRestHandler(OAL)和 DSLDebuggingRestHandler(MAL/LAL 通用):提交一段 DSL 脚本,OAP 编译并试跑,返回结果/错误,便于规则上线前验证。

runtime-rule · 运行时规则热更新

这是 server-admin 中最重的子模块。模块 RuntimeRuleModule(oap-server/server-admin/runtime-rule/.../RuntimeRuleModule.java:30,NAME=receiver-runtime-rule,:31),默认关闭,需 SW_RECEIVER_RUNTIME_RULE=default 且 admin-server 开启才加载。

REST 路由(RuntimeRuleRestHandler):/runtime/rule/addOrUpdate/inactivate/delete/list/status/dump,以及按 DSL 分的 /runtime/mal/otel/*/runtime/mal/log/*/runtime/lal/*

图4

三层架构(RuntimeRuleModuleProvider 注释详述):

  • Scheduler(DSLManager + REST handler):DSL 无关。负责锁获取、集群 Suspend/Resume RPC 扇出、DAO 持久化、跨文件所有权校验��tick 调度(默认 30s + 启动同步 tick)、自愈、classloader 回收、alarm-reset 分发。
  • Orchestrators:DSLRuntimeApply(apply 流程:NEW/FILTER_ONLY/STRUCTURAL 分类 → 编译 → fireSchemaChanges → verify → commit/rollback);DSLRuntimeUnregister(下线流程)。
  • Engines(RuleEngine SPI):MalRuleEngineLalRuleEngine(未来 OalRuleEngine)。每个引擎实现 classify/compile/fireSchemaChanges/verify/commit/rollback/unregister,拥有 DSL 相关的 Javassist 类生成、DDL 探测、classloader 退役等。新增 DSL 只需实现 SPI + RuleEngineRegistry.register,不动 scheduler/orchestrator——这是 SPI 分层的意义:把"DSL 怎么编译"这种易变的逻辑隔离在引擎里,scheduler 只管调度和一致性,两者解耦。

集群协调:RuntimeRuleClusterServiceImpl(集群 gRPC 总线,收 Suspend/Resume/Forward RPC)+ RuntimeRuleClusterClient(广播/转发到 main 节点)+ MainRouter(main/peer 路由)。状态机:DSLRuntimeStateApplyStatusApplyPhaseSchemaApplyCoordinatorAppliedRuleScriptLockreconcile 包下 RuleSyncSuspendResumeCoordinatorStructuralCommitCoordinatorPendingApplyCommit 做最终一致收敛。layer 包做运行时分层与冲突检测。extension/DbOverrideRuntimeRuleResolver 让 DB 覆盖规则。

热更新规则
MAL/LAL 是 SkyWalking 用 DSL 写的"指标加工规则/日志分析规则"。传统方式改规则要重启 OAP。runtime-rule 让你通过 admin REST 提交新规则文件,OAP 在线编译、生成新的处理类、平滑切换数据流、必要时挂起旧流再恢复新流,并跨集群节点同步状态,做到不停机改规则。复杂度集中在"如何在不丢数据、不破坏 schema 的前提下切换运行中的处理流水线",所以有锁、状态机、两阶段提交、集群扇出这一大套机制。

为什么不直接热替换?因为运行中的流处理正在消费 Kafka/接收数据,直接切类会丢在途数据。挂起旧流(停止拉新数据但处理完在途的)→ 切换 → 恢复新流,这套 Suspend/Resume 才不丢。两阶段提交保证集群所有节点要么都切成功、要么都回滚,不会出现一半节点跑新规则一半跑旧规则的不一致。

inspect · 数据检视

admin-only。InspectRestHandler 暴露 GET /inspect/metrics(指标目录)和 GET /inspect/entities(某指标在时间范围内的实体列表,解码成 MQE-ready 形式)。entity 枚举扫描会与用户查询竞争存储资源,故放在私有 admin 端口。

status · 状态查询

替代旧 status-query-plugin。暴露 AlarmStatusQueryServiceStatusModuleProvider 注册:DebuggingHTTPHandler(per-query 调试 trace)、TTLConfigQueryHandler(/status/config/ttl,也挂公共 REST 端口供生态工具发现 TTL)、ClusterStatusQueryHandlerAlarmStatusQueryHandler

ui-management · UI 模板管理

admin-only。UIManagementRestHandler 提供 dashboard 模板 CRUD:GET/POST/PUT /ui-management/templatesPOST /ui-management/templates/{id}/disable。转发到 CoreModule 的 UITemplateManagementService。注释明确:11.0.0+ 起 sidebar 菜单不再由 OAP 提供,改由 Horizon UI 自带 bundle。


health-checker · 健康聚合

解决什么问题

把 Telemetry 暴露的 health_check_* 指标聚合成一个"健康分数"并通过 HTTP 暴露,供 LB/外部探针判断 OAP 是否该接流量。

实现

  • HealthCheckerModule(oap-server/server-health-checker/src/main/java/org/apache/skywalking/oap/server/health/checker/module/HealthCheckerModule.java)暴露 HealthQueryService
  • HealthCheckerProvider.java:50:prepare()(:82)把 score 初始化为 -1(:83,表示还没算过),建虚拟线程定时器(:84)。notifyAfterCompleted()(:106)起定时任务(scheduleAtFixedRate,初始延迟 2s,周期 config.getCheckIntervalSeconds(),:127):从 MetricsCollector.collect()(:110)收集所有样本,metricsCreator.isHealthCheckerMetrics(:112)过滤出 health_check_* 指标,值 <1 视为不健康(:114),累加计数并拼出不健康模块名列表(:115),写入 score(不健康模块数,0 表示全健康,-1 初始态,:120-124)和 details(:125)。
  • start()(:95)把 HealthCheckerHttpService 挂到 CoreModule 的 HTTPHandlerRegister(:100-103,HEAD + GET)。requiredModules()(:130)依赖 TelemetryModule

为什么是值 <1 而不是 ==0?因为 HealthCheckMetrics 是 gauge,健康时上报 1,不健康上报 0,中间态或上报失败可能给个介于 0 和 1 的值。用 <1 当不健康门槛,留出余量,避免边界抖动。


server-tools · profile-exporter

解决什么问题

一个独立的 CLI 工具,从已落库的 profile 数据里把"线程快照 + segment 基础信息"导出成文件,用于离线调试 Continuous Profiling。

实现

  • ProfileSnapshotExporterBootstrap(oap-server/server-tools/profile-exporter/tool-profile-snapshot-bootstrap/.../ProfileSnapshotExporterBootstrap.java:42):export(args) 解析 ExporterConfig、用 ApplicationConfigLoader 加载 OAP 配置、初始化 ModuleManager(复用 OAP 启动框架,但用 mock 模块MockCoreModuleProviderMockRemoteClientManager 避免起真实服务)、查 profiled segment + span、写基础信息文件、ProfileSnapshotDumper.dump 查并写快照文件,最后 System.exit
  • tool-profile-snapshot-server-mock 提供 mock 实现(MockSourceReceiverMockGRPCHandlerRegister 等),让导出器在"只读存储 + 不开端口"的精简模式下跑起来。
  • 这是工具而非常驻服务,按需手动运行。

MockCoreModuleProvider 的注意事项
CLAUDE.md 提到:给 CoreModule.services() 加 Service 时,要更新所有 CoreModuleProvider 实现,包括 server-tools/profile-exporter/ 里的 MockCoreModuleProvider,否则 profile exporter e2e 测试启动会失败。


server-testing · 测试基础设施

解决什么问题

给 OAP 其它模块的单元/集成测试提供 mock 框架和 DSL 测试工具。

实现

  • MockModuleManagerMockModuleProviderModuleManagerTesting:提供 mock 的 ModuleManager/ModuleProvider,便于不启动完整 OAP 测单个模块。
  • dsl/HierarchyRuleLoaderDslClassOutput 等支持 DSL 测试。
  • 该目录还有 org/junit/rules/TestRule.javaorg/junit/runners/model/Statement.java,用于在受 JPMS/类加载限制时提供 JUnit 兼容桩。这是纯测试支撑 jar,不参与生产运行。

外围子系统总览

图5


关键概念新手答疑小结

  • 集群协调:多 OAP 节点要按 hash 分片转发数据,因此需要"互相发现 + 维护 gRPC 连接"。ClusterModule(ClusterModule.java:23)抽象注册/查询接口,插件对接不同注册中心,RemoteClientManager(RemoteClientManager.java:59)消费节点列表建连。业务代码应通过 CoreModuleRemoteClientManager 间接使用,不直接 find ClusterModule
  • 动态配置 watch:OAP 不重启改配置。各模块注册 ConfigChangeWatcher(ConfigChangeWatcher.java:30,带 module.provider.item key),配置中心 provider 轮询或监听外部,值变化回调 watcher.notify(ConfigWatcherRegister.java:35 比较后再回调)。分 Fetching(轮询,FetchingConfigWatcherRegister.java:36)和 Listening(事件,ListeningConfigWatcherRegister.java:27)两派。
  • 自监控 vs 数据导出:Telemetry 是"OAP 自己有多忙"(Prometheus /metrics),Exporter 是"OAP 算出的业务数据外发"(gRPC/Kafka,订阅驱动)。两者都用 prometheus client 但服务对象不同。
  • AI baseline:MQE 的 baseline() 运算符经 AI Pipeline 模块向外部预测服务取基线 + 上下界(BaselineQueryServiceImpl.java:51),用于异常检测,不是 OAP 自己算的。
  • 热更新规则:runtime-rule 通过 admin REST 在线编译/切换 MAL/LAL 规则,配锁、状态机、集群扇出做最终一致,不停机改处理流水线。所有管理面端点默认关闭且无鉴权,必须网关保护。

关键文件速查

子系统 关键类 路径
集群 ClusterModule oap-server/server-core/.../cluster/ClusterModule.java
RemoteClientManager oap-server/server-core/.../remote/client/RemoteClientManager.java
ZookeeperCoordinator oap-server/server-cluster-plugin/cluster-zookeeper-plugin/.../ZookeeperCoordinator.java
KubernetesCoordinator oap-server/server-cluster-plugin/cluster-kubernetes-plugin/.../KubernetesCoordinator.java
动态配置 ConfigurationModule oap-server/server-configuration/configuration-api/.../ConfigurationModule.java
ConfigChangeWatcher 同目录 api/ConfigChangeWatcher.java
FetchingConfigWatcherRegister oap-server/server-configuration/configuration-api/.../FetchingConfigWatcherRegister.java
GRPCConfigWatcherRegister oap-server/server-configuration/grpc-configuration-sync/.../GRPCConfigWatcherRegister.java
Telemetry TelemetryModule / PrometheusTelemetryProvider oap-server/server-telemetry/telemetry-api/...telemetry-prometheus/.../PrometheusTelemetryProvider.java
Exporter ExporterModule / ExporterProvider / GRPCMetricsExporter oap-server/exporter/...
AI Pipeline AIPipelineModule / BaselineQueryServiceImpl oap-server/ai-pipeline/...
server-admin AdminServerModuleProvider oap-server/server-admin/admin-server/...
runtime-rule 三层 oap-server/server-admin/runtime-rule/...
health-checker HealthCheckerProvider oap-server/server-health-checker/.../HealthCheckerProvider.java
server-tools ProfileSnapshotExporterBootstrap oap-server/server-tools/profile-exporter/tool-profile-snapshot-bootstrap/.../ProfileSnapshotExporterBootstrap.java
server-testing MockModuleManager oap-server/server-testing/.../mock/MockModuleManager.java

全系列回顾

到这里,SkyWalking OAP 的所有核心子系统都拆解完了。回到 README 看整体地图,或回看任意一篇:

  • 00-整体架构与端到端数据流
  • 01-模块系统与启动流程
  • 02-数据采集层-Receiver
  • 03-分析引擎与四大DSL
  • 04-流处理引擎-Worker链
  • 05-存储层
  • 06-查询层与告警

心智模型
把外围子系统想成 OAP 工厂的"后勤部":集群协调是"分厂通讯录"(谁在哪、怎么联系),动态配置是"工艺参数远程调整"(不停线改设置),Telemetry 是"工厂自身的仪表盘"(机器转得快不快),Exporter 是"对外供货窗口"(按订单把成品发给外部),AI Pipeline 是"聘请的外部顾问"(预测基线、归一化 URI 都外包给它),server-admin 是"厂长办公室"(调试、改工艺单、查状态,默认锁门要授权进),health-checker 是"门卫"(看工厂健康分决定要不要放行流量)。