
后端工作流自动化任务调度【免费下载链接】temporalTemporal service项目地址https://gitcode.com/gh_mirrors/te/temporal点击查看免费下载Temporal Server 的 Worker 服务承担了集群内所有后台处理任务其中 Replicator 是跨数据中心Cross-DC简称 CDC复制的核心组件它负责消费远程 Temporal 集群产生的复制任务并将其应用到本地集群。本文以 service/worker/README.md 为主线结合仓库源码深入讲解 Worker 的角色定位、Replicator 的工作原理并给出完整的本地双集群复制环境搭建与故障转移Failover实操步骤。Worker 在 Temporal 集群中的角色按照 service/worker/README.md 的定义Temporal Worker 是 Temporal 服务中的一个角色role专门用于承载任何负责在 Temporal 集群上执行后台处理的组件。它不是一个暴露给业务开发者的 SDK Worker而是 Temporal 服务器自身内部的系统 Worker。从 service/worker/service.go 的注释可以看到 Worker 服务承载的核心职责Replicator处理由远程集群生成的复制任务replication tasksArchiver负责工作流历史的归档archival此外还包括 Scanner数据清理与一致性检查、Parent Close Policy Processor父工作流关闭策略执行、Batcher批量操作、Scheduler调度器等。service/worker/fx.go 中的依赖注入模块列出了实际注册进 Worker 的全部子模块migration、deletenamespace、chasmscheduler、callback、scheduler、batcher、workerdeployment、dlq、dummy等与 service/worker 目录下的子目录一一对应。Worker 服务的启动逻辑在 service/worker/service.go依次启动 membership monitor、Scanner然后在全局命名空间模式IsGlobalNamespaceEnabled()下才启动 Replicator见 service.go#L271-L273随后启动系统 SDK Worker 管理器worker.go它使用系统客户端在temporal-system默认任务队列上注册所有 Worker 组件的工作流与活动。Replicator跨集群复制的后台消费者Replicator 的职责README 中明确描述Replicator 是一个后台 Worker负责消费远程 Temporal 集群生成的复制任务replication tasks并将其传递给 processor处理器从而应用到本地 Temporal 集群。从源码看Replicator 的实现在 service/worker/replicator/replicator.go// Replicator is the processor for replication tasks Replicator struct { status int32 clusterMetadata cluster.Metadata namespaceReplicationTaskExecutor nsreplication.TaskExecutor customTaskHandler func(ctx context.Context, task *replicationspb.ReplicationTask) error clientBean client.Bean ... namespaceReplicationQueue persistence.NamespaceReplicationQueue namespaceProcessors map[string]*replicationMessageProcessor matchingClient matchingservice.MatchingServiceClient namespaceRegistry namespace.Registry }关键字段说明字段作用clusterMetadata集群元数据决定本地集群名称、远程集群列表及启用状态namespaceProcessors按源集群source cluster维度维护的处理器映射每个远程集群对应一个replicationMessageProcessornamespaceReplicationQueue持久化的命名空间复制队列保存待处理的复制消息与确认水位ack levelnamespaceReplicationTaskExecutor实际执行命名空间复制任务的执行器注册、更新、故障转移等customTaskHandler可选的扩展钩子用于处理自定义类型的复制任务启动与集群元数据监听Replicator 的Start()replicator.go#L103-L114做两件事调用listenToClusterMetadataChange()注册集群元数据变更回调启动后台 goroutine 周期性清理命名空间复制队列中已确认的消息。listenToClusterMetadataChangereplicator.go#L136-L185是理解 Replicator 动态性的关键当集群元数据发生变化例如通过管理命令 upsert 了远程集群回调会遍历新的集群信息跳过当前集群自身为每个Enabled的远程集群创建一个replicationMessageProcessor并启动若某个处理器对应的集群被禁用或移除则先Stop()再删除。for clusterName : range newClusterMetadata { if clusterName currentClusterName { continue } ... if clusterInfo : newClusterMetadata[clusterName]; clusterInfo ! nil clusterInfo.Enabled { remoteAdminClient, err : r.clientBean.GetRemoteAdminClient(clusterName) ... processor : newReplicationMessageProcessor(...) processor.Start() r.namespaceProcessors[clusterName] processor } }这解释了 README 第 3 步连接两个集群upsert-remote-cluster之后Replicator 无需重启即可自动感知新集群的原因。复制消息处理循环processor每个远程集群对应一个replicationMessageProcessorservice/worker/replicator/replication_message_processor.go其核心处理循环在processorLoop中周期性默认 1 秒带 0.2 抖动系数见 replication_message_processor.go#L31-L33调用handleReplicationTasks。handleReplicationTasksreplication_message_processor.go#L156-L269)的工作流程Ring 成员资格检查通过serviceResolver.Lookup(p.sourceCluster)查找负责该源集群的 Worker只有当自己hostInfo.Identity()就是负责人时才继续处理。这是保证每个源集群同一时刻只有一个 Worker 在处理的最佳努力机制源码注释指出 ring 重配置期间可能出现短暂多个 worker 同时执行但不会造成正确性问题因为命名空间复制任务受版本号校验保护拉取任务调用远程集群的 Admin 服务GetNamespaceReplicationMessages携带LastRetrievedMessageId与LastProcessedMessageId实现断点续传逐条处理对每个复制任务按类型分发handleReplicationTaskreplication_message_processor.go#L330-L366)REPLICATION_TASK_TYPE_NAMESPACE_TASK交给namespaceTaskExecutor.Execute处理命名空间注册、更新、故障转移REPLICATION_TASK_TYPE_TASK_QUEUE_USER_DATA调用 Matching 服务的ApplyTaskQueueUserDataReplicationEvent应用任务队列用户数据如版本化数据其他类型若有customTaskHandler则交给自定义处理器否则返回错误失败处理处理失败后按指数退避策略重试仍失败则写入 DLQ死信队列putNamespaceReplicationTaskToDLQ并记录ReplicatorFailures、ReplicatorDLQFailures指标推进水位成功处理后更新lastProcessedMessageID与lastRetrievedMessageID。值得注意的重试策略细节replication_message_processor.go#L62-L80)命名空间复制任务namespace task获得最多 30 次重试而其他类型任务默认最多 5 次。源码注释解释了原因命名空间任务受集群全局命名空间元数据 CASCompare-And-Swap保护在高并发下 CAS 可能频繁冲突或存储抖动更多重试次数可以避免任务被过早丢入 DLQ同时保持单线程循环对其他任务类型的响应性。已确认消息清理Replicator 还运行一个每 5 分钟触发一次的清理任务cleanupNamespaceReplicationQueuereplicator.go#L226-L247计算所有已连接集群确认水位的最低值connectedClustersLowWatermark调用DeleteMessagesBefore删除该水位之前的复制消息防止队列无限增长。本地 CDC 开发环境搭建QuickstartREADME 提供了 5 步搭建本地双集群复制环境的指南。以下命令均以 README 为准完整保留并结合当前仓库补充说明。前置条件本仓库代码已 clone且已安装 Go 与 Docker用于拉起依赖服务需要tctl命令行工具README 使用的管理命令为 tctl 风格。第 1 步启动 active主区域的开发服务器make start-cdc-active说明README 记录的 target 名为start-cdc-active。在当前仓库的 MakefileMakefile中跨集群开发环境相关的 target 为start-dependencies-cdc拉起依赖容器、start-xdc-cluster-a、start-xdc-cluster-b、start-xdc-cluster-c分别启动三个开发集群的 temporal-server对应 config/development-cluster-a.yaml、config/development-cluster-b.yaml、config/development-cluster-c.yaml。若当前 Makefile 中不存在start-cdc-active可执行make start-xdc-cluster-a或直接运行./temporal-server --config-file config/development-cluster-a.yaml --allow-no-auth start第 2 步启动 standby备区域的开发服务器make start-cdc-standby同理也可使用./temporal-server --config-file config/development-cluster-b.yaml --allow-no-auth start两个集群的配置差异集中在 clusterMetadata 段。以 cluster-a 为例clusterMetadata: enableGlobalNamespace: true failoverVersionIncrement: 100 masterClusterName: cluster-a currentClusterName: cluster-a clusterInformation: cluster-a: enabled: true initialFailoverVersion: 1 rpcName: frontend rpcAddress: localhost:7233配置项含义配置项说明enableGlobalNamespace: true开启全局命名空间模式这是 Replicator 启动的必要条件见上文IsGlobalNamespaceEnabled()判断failoverVersionIncrement: 100故障转移版本号的递增步长用于多集群间排序冲突解决masterClusterName主集群名称currentClusterName当前集群名称本集群身份clusterInformation.name.enabled是否启用该集群的复制连接clusterInformation.name.initialFailoverVersion该集群的初始故障转移版本号cluster-a 为 1cluster-b 为 2clusterInformation.name.rpcAddress该集群 frontend 的地址cluster-a 为localhost:7233cluster-b 为localhost:8233CDC 依赖服务额外的 standby/other UI 容器由 develop/docker-compose/docker-compose.cdc.yml 定义Linux 平台覆写见 develop/docker-compose/docker-compose.cdc.linux.yml。第 3 步连接两个 Temporal 集群在两个集群之间互相注册对方为远程集群tctl --ad 127.0.0.1:7233 adm cl upsert-remote-cluster --frontend_address localhost:8233 tctl --ad 127.0.0.1:8233 adm cl upsert-remote-cluster --frontend_address localhost:7233第一条在 active 集群frontend 端口 7233注册 standby 集群8233第二条在 standby 集群注册 active 集群7233。说明这是 README 记录的 tctl 命令。从 config/development-cluster-a.yaml 的配置注释可见新版 Temporal 建议使用temporalCLI 的 operator 子命令效果等价temporal --address 127.0.0.1:7233 operator cluster upsert --frontend-address 127.0.0.1:8233 --enable-connection --enable-replication该命令执行后目标集群的clusterMetadata更新本集群 Replicator 通过listenToClusterMetadataChange回调自动创建对应源集群的复制消息处理器见 replicator.go#L136-L185无需重启服务。第 4 步创建全局命名空间tctl --ns sample namespace register --gd true --ac active --cl active standby参数解析--ns sample命名空间名为sample--gd true注册为全局命名空间global namespace只有全局命名空间才会在集群间复制--ac active设置 active 集群为active--cl active standby将该命名空间复制到active与standby两个集群。注册全局命名空间会触发命名空间复制任务该任务通过集群间的复制队列传递给 standby 集群由 standby 的 Replicator 消费并应用对应REPLICATION_TASK_TYPE_NAMESPACE_TASK处理路径replication_message_processor.go#L342-L350。故障转移Failover操作全局命名空间任一时刻只有一个 active 集群负责读写通过切换 active 集群实现故障转移。切换到 standby故障转移tctl --ns sample namespace update --ac standby切回 active故障恢复/回切tctl --ns sample namespace update --ac active每次切换都会更新命名空间的 active 集群并递增故障转移版本号配置中的failoverVersionIncrement用于多集群间的版本冲突仲裁。切换动作会生成新的命名空间复制任务由 Replicator 传输到另一集群并应用在此过程中任务处理受集群全局命名空间元数据 CAS 保护因此副本消息处理采用更宽松的 30 次重试策略见 replication_message_processor.go#L38-L43避免瞬时冲突把任务丢入 DLQ。验证与测试仓库中提供了对应的测试与集成用例Replicator 单元测试service/worker/replicator/replicator_test.go 与 replication_message_processor_test.go覆盖启动/停止、集群元数据变更回调、消息拉取与 DLQ 等路径跨数据中心XDC集成测试tests/xdc 目录下共 25 个测试文件覆盖双集群/多集群场景下的复制与故障转移行为被动路径passive path测试tests/passivepath 针对 standby 集群执行路径的验证。小结Temporal Worker 是集群后台处理的统一宿主Replicator 则是其中负责跨集群复制的关键组件。理解它的工作方式——从集群元数据监听、按源集群创建处理器、ring 成员资格选主、周期拉取复制消息、按类型分发应用到重试与 DLQ 兜底、水位清理——有助于在生产环境排查复制延迟、DLQ 堆积与故障转移异常等问题。本文给出的 5 步本地 CDC 环境搭建与故障转移操作可直接用于功能验证与问题复现。赞分享后端工作流自动化任务调度【免费下载链接】temporalTemporal service项目地址https://gitcode.com/gh_mirrors/te/temporal点击查看免费下载相关推荐minikube 本地 Kubernetes 集群实战指南从快速启动到多集群、Addons 与源码级解析minikube 本地 Kubernetes 集群实战指南从快速启动到多集群、Addons 与源码级解析 minikube 是 Kubernetes 官方社区云原生容器编排CLI开发工具Milvus CDC 跨集群复制Replication技术总览主备灾备拓扑、角色约束与故障切换机制Milvus CDC 跨集群复制Replication技术总览主备灾备拓扑、角色约束与故障切换机制 本文以 docs/design docs/design数据库向量数据库分布式数据库后端Temporal多集群部署实战跨区域数据复制与灾备终极指南 Temporal多集群部署实战跨区域数据复制与灾备终极指南 Temporal作为现代化的 工作流编排引擎 其多集群部署能力为企业级应用提供了强大的 跨后端工作流自动化任务调度上一篇3分钟搞定U盘启动Rufus格式转换接口深度解析下一篇无需重启Atmosphere模块加载器让Switch功能实时升级创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考