集群协调与动态配置
这一篇解决什么
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) 等。

- 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 都通过 TelemetryModule 的 MetricsCreator 创建 HealthCheckMetrics,在查询/注册失败时标记 unHealth(如 EtcdCoordinator.java:198 用 createHealthCheckerGauge("cluster_etcd", ...)),并经 OAPNodeChecker.isHealth() 做节点健康校验(EtcdCoordinator.java:114)。这些 health_check_* 指标后面会被 health-checker 聚合成"健康分"。
与 RemoteClientManager 的关系
不要直接 find ClusterModule
业务代码不该直接moduleManager.find(ClusterModule.NAME)拿集群服务。应该走CoreModule的RemoteClientManager——它才是稳定的对外门面,内部再去找ClusterModule。这样换集群实现(zookeeper/etcd/k8s)对业务透明。
RemoteClientManager(oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/client/RemoteClientManager.java:59)实现 Service 和 ClusterWatcher 两个接口(:59 implements Service, ClusterWatcher):
start()(:95)起定时任务,每 10 秒refresh()(:99的scheduleWithFixedDelay,参数1, 10, SECONDS)。refresh()(:105)通过moduleDefineHolder.find(ClusterModule.NAME).provider().getService(ClusterNodesQuery.class)拿集群节点列表(:107)——这正是上面提示的关键:业务代码通过 CoreModule 拿 RemoteClientManager,RemoteClientManager 内部才 find ClusterModule。refresh(List<RemoteInstance>)(:118,synchronized)去重(distinct:179)、排序(Collections.sort:134)、与现有 client 比对(compare:142);若变化则reBuildRemoteClients(:146)。reBuildRemoteClients(:202):对每个地址,self 用SelfRemoteClient(:239),否则建GRPCRemoteClient并connect()(:243-244);消失的 client 关闭(:261)。最后排序保证稳定顺序(:252注释//for stable ordering for rolling selector)。- 同时实现
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(),没变就跳过。
两种同步风格:

FetchingConfigWatcherRegister(FetchingConfigWatcherRegister.java:36):定时轮询。start()(:79)用虚拟线程定时器每syncPeriod(默认 60s,构造器:46)调configSync()(:99),子类实现readConfig(keys)(:155)/readGroupConfig(keys)(:157)拉配置,再逐项notifySingleValue(:120)。isStarted标志(:42)保证启动后不能再注册新 watcher(:55-56抛IllegalStateException)——启动后再加 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)暴露MetricsCreator和MetricsCollector两个 Service(:37-39)。MetricsCreator是工厂接口,创建createCounter/createGauge/createHistogramMetric/createHealthCheckerGauge。HEALTH_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:MetricValuesExportService、TraceExportService、LogExportService。这是"消费方契约定义在 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 左右):
- 连接远端
MetricExportServicegRPC 服务,本地起BatchQueue缓冲(:82起BatchQueueManager.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)接口方法:querySupportedMetrics、queryPredictMetrics(serviceMetrics, start, end)、queryPredictMetricsFromCache(serviceName, timeBucketHour)。BaselineQueryServiceImpl(BaselineQueryServiceImpl.java:51)通过 gRPC stubAlarmBaselineServiceGrpc(: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 操作的关系
已从源码确认的关键链路:

MQEVisitorBase.java:42-43import 了AIPipelineModule和BaselineQueryService。MQEVisitorBase.java:78-91持有BaselineQueryService,通过moduleManager.find(AIPipelineModule.NAME).provider().getService(BaselineQueryService.class)懒加载(getBaselineQueryService:87)。懒加载是因为 MQE 模块不一定依赖 AI Pipeline,没装 AI Pipeline 时不能在构造期就 find,否则启动报错。MQEVisitorBase.java:502visitBaselineOP—— 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)暴露 HTTPHandlerRegister、GRPCHandlerRegister、AdminClusterChannelManager(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/*。

三层架构(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(
RuleEngineSPI):MalRuleEngine、LalRuleEngine(未来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 路由)。状态机:DSLRuntimeState、ApplyStatus、ApplyPhase、SchemaApplyCoordinator、AppliedRuleScriptLock。reconcile 包下 RuleSync、SuspendResumeCoordinator、StructuralCommitCoordinator、PendingApplyCommit 做最终一致收敛。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。暴露 AlarmStatusQueryService。StatusModuleProvider 注册:DebuggingHTTPHandler(per-query 调试 trace)、TTLConfigQueryHandler(/status/config/ttl,也挂公共 REST 端口供生态工具发现 TTL)、ClusterStatusQueryHandler、AlarmStatusQueryHandler。
ui-management · UI 模板管理
admin-only。UIManagementRestHandler 提供 dashboard 模板 CRUD:GET/POST/PUT /ui-management/templates、POST /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 模块如MockCoreModuleProvider、MockRemoteClientManager避免起真实服务)、查 profiled segment + span、写基础信息文件、ProfileSnapshotDumper.dump查并写快照文件,最后System.exit。tool-profile-snapshot-server-mock提供 mock 实现(MockSourceReceiver、MockGRPCHandlerRegister等),让导出器在"只读存储 + 不开端口"的精简模式下跑起来。- 这是工具而非常驻服务,按需手动运行。
MockCoreModuleProvider 的注意事项
CLAUDE.md 提到:给CoreModule.services()加 Service 时,要更新所有CoreModuleProvider实现,包括server-tools/profile-exporter/里的MockCoreModuleProvider,否则 profile exporter e2e 测试启动会失败。
server-testing · 测试基础设施
解决什么问题
给 OAP 其它模块的单元/集成测试提供 mock 框架和 DSL 测试工具。
实现
MockModuleManager、MockModuleProvider、ModuleManagerTesting:提供 mock 的ModuleManager/ModuleProvider,便于不启动完整 OAP 测单个模块。dsl/下HierarchyRuleLoader、DslClassOutput等支持 DSL 测试。- 该目录还有
org/junit/rules/TestRule.java、org/junit/runners/model/Statement.java,用于在受 JPMS/类加载限制时提供 JUnit 兼容桩。这是纯测试支撑 jar,不参与生产运行。
外围子系统总览

关键概念新手答疑小结
- 集群协调:多 OAP 节点要按 hash 分片转发数据,因此需要"互相发现 + 维护 gRPC 连接"。
ClusterModule(ClusterModule.java:23)抽象注册/查询接口,插件对接不同注册中心,RemoteClientManager(RemoteClientManager.java:59)消费节点列表建连。业务代码应通过CoreModule的RemoteClientManager间接使用,不直接 findClusterModule。 - 动态配置 watch:OAP 不重启改配置。各模块注册
ConfigChangeWatcher(ConfigChangeWatcher.java:30,带module.provider.itemkey),配置中心 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 是"门卫"(看工厂健康分决定要不要放行流量)。