Apache DolphinScheduler 注册中心(Registry)SPI 扩展指南:插件配置、源码原理与自定义实现
2026/9/23 14:30:29 网站建设 项目流程

Apache DolphinScheduler 注册中心(Registry)SPI 扩展指南:插件配置、源码原理与自定义实现

【免费下载链接】dolphinschedulerApache DolphinScheduler is the modern data orchestration platform. Agile to create high performance workflow with low-code项目地址: https://gitcode.com/gh_mirrors/do/dolphinscheduler

注册中心是 Apache DolphinScheduler 集群协调的基石,负责 Master/Worker 节点元数据存储、上下线感知、任务负载均衡以及故障转移时的全局锁。本指南以官方 SPI 扩展文档为核心,结合当前仓库源码,完整讲解注册中心插件的配置方式、参数语义,并深入剖析dolphinscheduler-registry-api的 SPI 接口设计与自定义插件扩展流程,帮助你从"会用"进阶到"能写"。

注册中心在 DolphinScheduler 中承担什么职责

在深入 SPI 之前,先明确注册中心的价值。根据 dolphinscheduler-registry/README.md,DolphinScheduler 使用注册中心完成三件事:

  1. 节点元数据存储与上下线感知:存储 Master/Worker 的元数据,当节点上下线时,其他节点能够及时收到通知;
  2. Worker 元数据与负载均衡:存储 Worker 元数据,供 Master 做任务分发时的负载均衡;
  3. 全局锁:执行故障转移(Failover)时获取全局锁,保证同一时刻只有一个节点执行主备切换逻辑。

因此,一个合格的注册中心实现必须满足三个基本能力:在订阅路径上数据被新增/删除/更新时通知服务端;提供创建/释放全局锁的机制;当服务端异常下线时自动删除其元数据(即临时节点语义)。

注册中心 SPI 插件架构总览

整个注册中心功能划分为三个模块,位于 dolphinscheduler-registry 目录下:

模块作用
dolphinscheduler-registry-api定义注册中心 SPI 标准接口与核心模型,是所有插件必须依赖的 API 层
dolphinscheduler-registry-plugins具体的注册中心实现插件,当前提供 zookeeper、etcd、jdbc 三种
dolphinscheduler-registry-all聚合导出模块,用于将选定的插件实现打包进发行版,新增插件时需在此模块的pom.xml中追加依赖

其中dolphinscheduler-registry-api是最关键的 SPI 契约层,其核心接口与模型如下(源码位于 dolphinscheduler-registry-api/src/main/java/org/apache/dolphinscheduler/registry/api):

  • Registry接口:注册中心 SPI 主接口,每个插件都必须实现它;
  • ConnectionListener接口:负责监听客户端与注册中心之间的连接状态变化,状态定义见ConnectionState
  • SubscribeListener接口:负责监听指定前缀下子节点的状态变化,事件内容定义见EventADD/REMOVE/UPDATE三种类型);
  • RegistryClient:对Registry的封装,为上层(Master/Worker/AlertServer)提供心跳注册、节点列表查询、分布式锁等业务级能力;
  • RegistryException:注册中心统一异常。

开箱即用的三种注册中心插件

dolphinscheduler-registry-plugins目录下目前提供三种官方插件,每种插件均以独立的 Maven 模块发布:

  • dolphinscheduler-registry-zookeeper:默认注册中心,基于 Apache Curator + ZooKeeper;
  • dolphinscheduler-registry-etcd:基于 jetcd 的 etcd 注册中心;
  • dolphinscheduler-registry-jdbc:基于数据库(MySQL/PostgreSQL)的注册中心,无需额外部署中间件。

从当前仓库的默认配置可以印证 ZooKeeper 的默认地位:dolphinscheduler-master/src/main/resources/application.yaml 与 dolphinscheduler-worker/src/main/resources/application.yaml 中默认配置均为registry.type: zookeeper,而 dolphinscheduler-standalone-server/src/main/resources/application.yaml(单机模式)默认采用registry.type: jdbc,直接复用内置数据库,无需额外中间件即可启动。

如何使用:以 ZooKeeper 为例的插件配置

官方 SPI 文档以 ZooKeeper 为例给出了最简配置。需要说明的是,配置方式经历了演进:

  • 早期版本通过registry.properties文件(位于dolphinscheduler-service/src/main/resources/registry.properties)配置,内容为:

    registry.plugin.name=zookeeper registry.servers=127.0.0.1:2181
  • 当前仓库已全面迁移到 Spring Boot 风格的application.yaml配置,每个服务(master/worker/api)维护自己的registry配置块。以 dolphinscheduler-worker/src/main/resources/application.yaml 为例:

    registry: type: zookeeper zookeeper: namespace: dolphinscheduler connect-string: localhost:2181 retry-policy: base-sleep-time: 1s max-sleep: 3s max-retries: 5 session-timeout: 60s connection-timeout: 15s block-until-connected: 15s digest: ~

配置完成后启动 DolphinScheduler 集群,集群即使用 ZooKeeper 作为注册中心存储服务端元数据。官方 dolphinscheduler-registry-zookeeper/README.md 提供了同样的完整示例。

配置前缀规则(重要)

官方 SPI 文档强调了一条贯穿始终的规则:所有配置信息的前缀都需要加上registry。例如参数base.sleep.time.ms,在 registry 中应配置为:

registry.base.sleep.time.ms=100

该规则在当前源码中同样成立:各插件属性类均以@ConfigurationProperties(prefix = "registry")绑定配置(见 ZookeeperRegistryProperties.java、EtcdRegistryProperties.java),因此 YAML 中表现为registry.zookeeper.session-timeoutregistry.zookeeper.retry-policy.base-sleep-time这样的层级结构。

ZooKeeper 插件参数详解

ZooKeeper 插件全部参数定义于 ZookeeperRegistryProperties.java,默认值与语义如下:

配置项默认值说明
zookeeper.namespacedolphinscheduler所有节点在 ZooKeeper 中挂载的根路径,不同集群可用不同 namespace 隔离
zookeeper.connect-string无(必填)ZooKeeper 连接串,如127.0.0.1:2181,多节点逗号分隔
zookeeper.retry-policy.base-sleep-time1s重试基础休眠时间,构建ExponentialBackoffRetry时使用
zookeeper.retry-policy.max-sleep3s重试最大休眠时间
zookeeper.retry-policy.max-retries3最大重试次数
zookeeper.session-timeout60s会话超时时间,必须为正数
zookeeper.connection-timeout15s连接超时时间,必须为正数
zookeeper.block-until-connected15s启动时阻塞等待连接成功的最大时长
zookeeper.digestZooKeeper ACL 认证信息,配置后启用digest鉴权并应用CREATOR_ALL_ACL权限

该属性类实现了Validator接口,在 validate() 中强制校验connectStringnamespace非空,sessionTimeoutconnectionTimeoutblockUntilConnected必须为正数,启动时还会将生效配置打印到日志,方便核对。

从配置到连接的源码链路

ZooKeeper 插件的接入完全由 Spring Boot 条件装配驱动。ZookeeperRegistryAutoConfiguration上标注了:

@ConditionalOnProperty(prefix = "registry", name = "type", havingValue = "zookeeper")

即只有registry.type=zookeeper时该自动配置类才会生效(见 ZookeeperRegistryAutoConfiguration.java),随后以@ConditionalOnMissingBean(value = Registry.class)的条件创建ZookeeperRegistry并立即调用start()

ZookeeperRegistry的构造过程(见 ZookeeperRegistry.java)做了三件事:

  1. retry-policy三项参数转换为 Curator 的ExponentialBackoffRetry
  2. 通过CuratorFrameworkFactory.builder()设置connectStringnamespacesessionTimeoutMsconnectionTimeoutMs
  3. 若配置了digest,则追加digest授权与CREATOR_ALL_ACL的 ACL Provider。

start()则调用client.blockUntilConnected(blockUntilConnected)同步等待连接建立,超时未连接上会抛出RegistryException并关闭客户端。

其他插件配置:etcd 与 jdbc

etcd 插件

在 master/worker/api 的application.yaml中配置(完整示例见 dolphinscheduler-registry-etcd/README.md):

registry: type: etcd endpoints: "http://etcd0:2379, http://etcd1:2379, http://etcd2:2379" # 以下均有默认值 namespace: dolphinscheduler connection-timeout: 9s retry-delay: 60ms # 单位毫秒 retry-max-delay: 300ms retry-max-duration: 1500ms # SSL 选项按需配置 cert-file: "deploy/kubernetes/dolphinscheduler/etcd-certs/ca.crt" key-cert-chain-file: "deploy/kubernetes/dolphinscheduler/etcd-certs/client.crt" key-file: "deploy/kubernetes/dolphinscheduler/etcd-certs/client.pem" # 认证选项按需配置 user: "" password: "" authority: "" load-balancer-policy: ""

etcd 插件参数定义于 EtcdRegistryProperties.java,包括endpoints(必填)、namespace(默认dolphinscheduler)、connectionTimeout(默认 9s)、ttl(默认 30s,临时节点租约时长)、重试策略三项、负载均衡策略以及 SSL/认证相关配置。启用逻辑与 ZooKeeper 一致,由EtcdRegistryAutoConfiguration上的@ConditionalOnProperty(prefix = "registry", name = "type", havingValue = "etcd")控制(见 EtcdRegistryAutoConfiguration.java)。

jdbc 插件

jdbc 插件复用 DolphinScheduler 自身的数据库作为注册中心,省去额外部署 ZooKeeper/etcd 的成本,单机模式(standalone)即默认使用它。使用步骤(详见 dolphinscheduler-registry-jdbc/README.md):

  1. 初始化数据库表:MySQL 执行src/main/resources/mysql_registry_init.sql,PostgreSQL 执行src/main/resources/postgresql_registry_init.sql(位于 dolphinscheduler-registry-jdbc 模块下);

  2. 修改配置:在 master/worker/api 的application.yaml中设置:

    registry: type: jdbc # 心跳刷新间隔,默认 3s,不能小于 1s heartbeat-refresh-interval: 3s # 心跳超时时间,默认 60s,必须大于 3 * heartbeat-refresh-interval session-timeout: 60s # Hikari 连接池配置,默认复用 DolphinScheduler 自身数据源 hikari-config: jdbc-url: jdbc:mysql://127.0.0.1:3306/dolphinscheduler username: root password: root maximum-pool-size: 5 connection-timeout: 9000 idle-timeout: 600000

jdbc 插件的参数校验逻辑位于 JdbcRegistryProperties.java:heartbeatRefreshInterval不得小于 1s,sessionTimeout必须大于3 * heartbeatRefreshInterval;未显式配置jdbcRegistryClientName时默认取主机名:server.port。该插件的自动装配类JdbcRegistryAutoConfiguration会以@MapperScan扫描注册中心专属的 Mapper(JdbcRegistryDataMapperJdbcRegistryLockMapper等),并基于hikari-config构建独立的SqlSessionFactory(见 JdbcRegistryAutoConfiguration.java)。

注意:若使用 MySQL 作为注册中心,需要将mysql-connector-java.jar加入 DolphinScheduler 的 classpath,发行版不会内置该驱动(参见 jdbc 插件 README 的说明)。

如何扩展:实现自定义注册中心插件

官方 SPI 文档指出:dolphinscheduler-registry-api定义了实现插件的标准,扩展插件时只需实现标准接口即可。文档记载的扩展入口是org.apache.dolphinscheduler.registry.api.RegistryFactory;从当前仓库源码看,该 SPI 已演进为直接实现Registry接口并通过 Spring Boot 条件装配注册为 Bean的方式,下面以当前代码为准确认完整扩展步骤。

第一步:实现RegistrySPI 接口

Registry接口定义在 Registry.java,共 12 个核心方法,自定义插件必须全部实现:

方法语义
start()启动注册中心客户端并建立连接
isConnected()当前是否已连接
connectUntilTimeout(Duration)在给定超时内阻塞等待连接成功,超时抛RegistryException
subscribe(path, SubscribeListener)订阅路径及其子路径,数据变化时回调监听器
addConnectionStateListener(ConnectionListener)注册连接状态监听器
get(key)读取 key 的值,key 不存在时抛异常
put(key, value, deleteOnDisconnect)写入键值;deleteOnDisconnect=true表示断连时自动删除(临时节点语义)
delete(key)删除 key
children(key)返回 key 的子节点集合
exists(key)key 是否存在
acquireLock(key)/acquireLock(key, timeout)获取指定前缀的全局锁,可带超时
releaseLock(key)释放全局锁

Event事件模型(见 Event.java)包含四个字段:被监听的前缀key、事件发生的完整路径path、路径对应的数据data、事件类型typeADD/REMOVE/UPDATE)。

第二步:参考既有插件的实现范式

ZooKeeper 插件的 ZookeeperRegistry.java 是最佳的参考范本,其关键实现技巧包括:

  • 临时节点put()中根据deleteOnDisconnect选择CreateMode.EPHEMERALCreateMode.PERSISTENT,并配合orSetData()实现"存在即更新";
  • 订阅:使用 Curator 的TreeCache监听整棵子树,通过内部类EventAdaptorTreeCacheEvent转换为统一的EventNODE_ADDED→ADDNODE_UPDATED→UPDATENODE_REMOVED→REMOVE);
  • 分布式锁:基于InterProcessMutex实现,并用ThreadLocal缓存每个线程已获取的锁;由于 etcd/jdbc 无法实现可重入锁,这里做了兼容处理——同一线程重复获取同一把锁时直接返回true,多次获取只需释放一次;
  • 断连清理close()时逐个关闭TreeCache与 Curator 客户端。

第三步:编写自动装配类并接入条件开关

参照 ZookeeperRegistryAutoConfiguration.java 的结构:

@Configuration(proxyBeanMethods = false) @ConditionalOnProperty(prefix = "registry", name = "type", havingValue = "your-plugin") public class YourRegistryAutoConfiguration { @Bean @ConditionalOnMissingBean(value = Registry.class) public YourRegistry yourRegistry(YourRegistryProperties properties) { YourRegistry registry = new YourRegistry(properties); registry.start(); return registry; } }

同时提供带@ConfigurationProperties(prefix = "registry")的属性类(可参考 ZookeeperRegistryProperties.java 或 EtcdRegistryProperties.java),这样所有配置项自然遵循"registry前缀"规则。装配完成后,RegistryConfiguration(见 RegistryConfiguration.java)会自动把该RegistryBean 包装成RegistryClient供上层使用。

第四步:在聚合模块中登记依赖

将新插件依赖添加到dolphinscheduler-registry-allpom.xml,使其随发行版打包;同时dolphinscheduler-registry-plugins的父pom.xml中加入新模块。官方 README 明确说明:新增注册中心时必须在dolphinscheduler-registry-all的 pom.xml 中追加依赖

第五步:复用统一测试基类验证

dolphinscheduler-registry-it模块提供了跨插件的统一测试基类RegistryTestCase(见 RegistryTestCase.java),覆盖isConnectedconnectUntilTimeoutsubscribe(ADD/UPDATE/REMOVE 全链路)、连接状态监听、get/put/deletechildrenexists、加锁解锁等全部 SPI 行为。ZooKeeper 插件的测试 ZookeeperRegistryTestCase.java 通过 Testcontainers 拉起zookeeper:3.8容器并将registry.zookeeper.connect-string指向映射端口后继承该基类即可跑通全部用例。自定义插件继承同一个基类,即可获得与其他插件一致的回归保障。

上层如何使用注册中心:RegistryClient 与节点路径

了解 SPI 后,再看上层如何消费它。RegistryClient(见 RegistryClient.java)在构造时即通过RegistryNodeType初始化三个常驻路径:/nodes/master/nodes/worker/nodes/alert-server(见 RegistryNodeType.java)。完整的节点路径规划如下:

用途路径
全量节点/nodes
Master 节点/nodes/master
Worker 节点/nodes/worker
AlertServer 节点/nodes/alert-server
Master 节点选举锁/lock/master-node
Master 故障转移锁/lock/master-failover
Master 任务组协调器锁/lock/master-task-group-coordinator
AlertServer 节点锁/lock/alert

RegistryClient提供的典型能力包括:persistEphemeral(key, value)注册临时心跳节点(对应Registry.put(key, value, true))、getServerList(RegistryNodeType)按心跳 JSON 解析出 Master/Worker/AlertServer 的存活列表(对应MasterHeartBeatWorkerHeartBeat等模型)、getLock/releaseLock获取全局锁、subscribe订阅节点变化。这也印证了 README 中"通知、负载均衡、全局锁"三大用途在代码层面的落地。

FAQ:注册中心连接超时怎么办

官方 SPI 文档给出最直接的建议:增大相关超时参数。结合上文参数表,可按现象分层处理:

  1. 启动即报连接失败(如zookeeper connect failed in ...):检查registry.zookeeper.connect-string是否正确可达,并适当调大block-until-connected(默认 15s),它是start()阶段阻塞等待连接建立的最大时长(见 ZookeeperRegistry.java);
  2. 运行期频繁断连/重连:调大session-timeout(默认 60s)与connection-timeout(默认 15s),并检查retry-policy三项(base-sleep-timemax-sleepmax-retries)是否足以应对瞬时网络抖动;
  3. 业务侧等待连接超时Cannot connect to registry in ... s):对应Registry.connectUntilTimeout(timeout)的调用点,可结合调用方传入的超时时间与上述连接参数综合评估;
  4. jdbc 插件场景:若使用 jdbc 注册中心出现"节点被误判离线",重点检查heartbeat-refresh-interval(默认 3s)与session-timeout(默认 60s,必须大于 3 倍心跳间隔)的匹配关系,以及hikari-config连接池参数(maximum-pool-sizeconnection-timeoutidle-timeout)是否过小。

所有超时类参数都遵循本文所述的registry前缀规则配置,改动后重启对应服务生效。

延伸阅读

  • 注册中心设计初衷与模块划分:dolphinscheduler-registry/README.md
  • SPI 契约源码:Registry.java、RegistryClient.java、RegistryNodeType.java
  • 各插件配置与用法:dolphinscheduler-registry-plugins下各插件模块的 README(zookeeper / etcd / jdbc)
  • 默认配置示例:dolphinscheduler-master/src/main/resources/application.yaml、dolphinscheduler-worker/src/main/resources/application.yaml
  • 统一 SPI 测试基类:RegistryTestCase.java

【免费下载链接】dolphinschedulerApache DolphinScheduler is the modern data orchestration platform. Agile to create high performance workflow with low-code项目地址: https://gitcode.com/gh_mirrors/do/dolphinscheduler

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询