☰
Spinnaker Orca Peering 跨集群执行数据同步机制深度解析:多区域部署与数据库迁移场景实战
2026/9/25 3:17:51 网站建设 项目流程
  • 后端
  • DevOps
  • 云原生
  • 微服务

【免费下载链接】spinnaker

Spinnaker is an open source, multi-cloud continuous delivery platform for releasing software changes with high velocity and confidence.

项目地址:https://gitcode.com/gh_mirrors/sp/spinnaker
点击查看免费下载

本文以 Spinnaker 编排引擎 orca 的orca-peering模块(orca/orca-peering/README.md)为核心,系统讲解多个 orca 集群(各自持有独立数据库)之间如何同步执行(execution)数据,覆盖peer/partition/foreign execution三大核心概念、基于PeeringAgent的增量复制算法、完整的orca.ymlpeering profile 配置与参数调优,以及配套的监控指标、动态开关和orca-interlink远程操作能力。读完本文,你将能够为多区域 Spinnaker 部署或数据库迁移场景配置并运维一套可用的 orca 执行数据对等(peering)机制。

背景与问题:为什么 orca 需要 Peering

在标准拓扑中,单个 orca 集群对应一个独立的 [SQL] 数据库,该数据库既保存全部执行历史(execution history),也保存执行队列(execution queue)。一旦出现以下两类场景,问题就随之而来:

  • 多区域(multi-region)部署:不同区域各有一个 orca 集群和各自的数据库,用户期望某个区域的执行记录在其他区域也能被看到、被查询,甚至被远程操作。
  • 数据库迁移(database migration):把 orca 从旧数据库迁往新数据库的过渡期,新旧两套数据库需要保持执行数据一致。

orca-peering正是为了解决"多个 orca 安装(每个都带自己的数据库)之间相互通告变更"这一半实验性(semi-experimental)问题而生的模块。它是一个轮询式的数据同步代理:从 peer 数据库中把执行数据复制到本地数据库,并保证复制过来的数据以只读方式存在,避免两个集群对同一执行互相操作造成冲突。

核心概念:peer、partition 与 foreign execution

理解 peering 机制前,必须先建立三个概念:

  • peer(对等集群):一个我们从中复制数据的 orca 集群(其数据库可以是只读副本)。每个 orca 集群拥有一个 ID,例如us-east-1、us-west-2。在 yaml 配置中,peer 由"数据库连接 + ID"两者共同定义。典型场景下,ID 为us-west-2的 orca 集群可以与 ID 为us-east-1的集群互相对等,反之亦然。

  • partition(分区):数据库中的执行记录都带有 partition 标记,partition 与上述 peer ID 同义。当一条执行从 ID 为us-east-1的 peer 被"peered(复制)"过来时,它会被以partition = us-east-1的方式持久化到本地数据库。由于历史原因,早期执行记录可能没有 partition 字段,因此一个 orca 集群会把partition = NULL或partition = 本集群 ID的执行视为自己所有。这一点在源码中同样成立——MySqlRawAccess.kt 在查询执行 ID 时,partition 约束就是partition IS NULL或partition = 指定值两种分支。

  • foreign execution(外部执行):出现在本地数据库中、但 partition 被标记为某个 peer 的执行。这些执行本质上是只读的,当前 orca 集群无法对它们执行任何变更操作。源码中,复制时会强制把每条执行的 partition 字段改写为 peer ID(见下文 ExecutionCopier 的分析),从而保证"属于谁"的判定一致。

Peering 机制能做什么

从 README 的定义看,peering 机制完成三件事:

  1. 同步(复制)执行数据:把 pipelines 和 orchestrations 两类执行从 peer 的数据库复制到本地集群数据库;
  2. 允许对 foreign execution 执行操作:例如操作正在 peer 上运行的执行(依托orca-interlink模块,通过消息转发给实际拥有者执行);
  3. (尚未实现)接管执行所有权:取得之前由 peer 操作的执行并继续运行,此能力在 README 中标注为 "still to come"。

执行同步(Execution Peering)的原理

执行同步本质上就是把执行从一份数据库复制到另一份数据库,但有一个关键取舍:执行历史需要被 peered,执行队列则不能。

为什么队列不能复制?因为:

  • 队列被复制会导致同一执行在多个集群被重复调度(duplicate executions);
  • 队列的变更频率和带宽极高,复制它会为数据库带来难以承受的开销。

因此 peering 只针对执行数据本身(及其 stages),而队列维持在各集群本地。所有 peering 逻辑都位于 PeeringAgent.kt(算法细节见其代码注释),高层思路如下:

  1. 给定一个 peer ID 及其数据库连接(可指向只读副本);
  2. 将所有带该 peer ID 的 foreign execution 镜像到本地数据库;
  3. 复制过程中,所有执行被标注为来自该 peer(写入partition列);
  4. 任何对 foreign execution(即partition != 本集群 ID的执行)的操作尝试都会失败。

PeeringAgent 的单轮执行流程

从源码看,PeeringAgent继承自AbstractPollingNotificationAgent,每个轮询周期(由intervalMs控制)执行一次tick()。tick()的流程(PeeringAgent.kt)为:

  1. 先通过DynamicConfigService检查全局开关pollers.peering与针对本 peer 的开关pollers.peering.<PEERID>(两者均默认开启);
  2. 依次对PIPELINE与ORCHESTRATION两种执行类型执行复制;
  3. 传播删除操作(peerDeletedExecutions);
  4. 若配置了自定义 peerer,调用invokeCustomPeerer();
  5. 整个过程被pollers.peering.lag计时器包裹,用于度量单轮总耗时。

对每种执行类型,复制又分为两个阶段(PeeringAgent.kt):

  • 首轮运行(isFirstRun):只复制已完成的执行(completed)。原因很直接:首次批量复制可能要耗时 20 分钟以上,此时复制进行中的执行会立刻过时,没有意义。
  • 后续运行:先做已完成执行的增量复制,再复制活动中的执行(active)。进行中执行的数量少、变化快,因此每轮全量抓取当前活动执行 ID 并复制。

peerCompletedExecutions内部通过doMigrate完成"增量差量"计算(PeeringAgent.kt),算法要点:

  • 从源库取回updated_at大于游标updatedAfter的已完成执行 ID(同时覆盖partition = peeredId与partition = NULL两批);
  • 从本地库取回已被复制的对应执行 ID;
  • 需要复制的 = 源库有、且updated_at比本地更新的执行;
  • 需要删除的 = 本地有、但源库已不存在的执行;
  • 删除前校验maxAllowedDeleteCount阈值,超过则整轮不做删除并记错误指标(防止误删全部执行);
  • 复制完成后,把"最新 updated_at − clockDriftMs"作为新的游标,从而容忍集群间时钟漂移。

删除的传播(peerDeletedExecutions,PeeringAgent.kt)依赖deleted_executions表:该表的主键是自增 int,PeeringAgent用它作为"游标"记录已传播的删除位置;只有整批删除全部成功后游标才会推进,失败则下轮重试。对不存在的执行执行"删除"是无害的,只是浪费少量数据库 CPU。

并行复制与分块细节(ExecutionCopier)

实际的批量复制由 ExecutionCopier.kt 完成:

  • copyInParallel把待复制 ID 按chunkSize分块放入并发队列,起min(threadCount, 分块数)个工作线程并行消费(ExecutionCopier.kt);
  • 单个分块的复制顺序非常讲究(ExecutionCopier.kt):
    1. 先抓执行行、再抓 stage 行——因为差量计算以执行的updated_at为准,绝不能出现"stage 抓完、执行又被更新"的时序颠倒;stage 比执行新没问题,下一轮 agent 运行会自行修正;
    2. 先复制 stages、再复制 executions——若先保存执行,用户可能看到一条还没有任何 stage 的执行;
    3. 源库 stage 列表可能已变化(例如重启 deploy stage 会删除其合成 stages 重新开始),因此先把本地"源已不存在的 stage"删掉,再对剩下的做增量更新/复制;
    4. 复制执行行时,强制把partition字段改写为peeredId,然后通过loadRecords(基于INSERT ... ON DUPLICATE KEY UPDATE)写入本地。

测试用例 PeeringAgentSpec.groovy 对上述逻辑做了行为验证:包括"全局/单 peer 动态开关是否生效"、"已完成执行的差量计算是否正确(含删除集合、待复制集合、游标推进)",其中clockDrift在测试中被设为 10ms 以验证时钟漂移对游标的影响。

配置实战:orca.yml 中的 peering profile

README 给出了一个可直接参考的peeringprofile 配置片段(位于orca.yml),完整保留如下:

spring: profiles: peering pollers: peering: enabled: true poolName: foreign id: us-west-2 intervalMs: 5000 # This is the default value threadCount: 30 # This is the default value chunkSize: 100 # This is the default value clockDriftMs: 5000 # This is the default value queue: redis: enabled: false keiko: queue: enabled: false sql: enabled: true foreignBaseUrl: URL_OF_MYSQL_DB_TO_PEER_FROM:3306 partitionName: LOCAL_PARTITION_NAME connectionPools: foreign: jdbcUrl: jdbc:mysql://${sql.foreignBaseUrl}/orca?ADD_YOUR_PREFFERED_CONNECTION_STRING_PARAMS_HERE user: orca_service password: ${sql.passwords.orca_service} connectionTimeoutMs: 5000 validationTimeoutMs: 5000 maxPoolSize: ${pollers.peering.threadCount}

这份配置同时传达了几个关键设计:

  • 开启 peering 的同时必须关闭队列复制相关组件(queue.redis.enabled、keiko.queue.enabled均为false),与"队列不复制"的原则一致;
  • sql.connectionPools.foreign定义了指向 peer 数据库的连接池,maxPoolSize直接取自pollers.peering.threadCount,保证并行复制线程数有足够的连接可用;
  • 源码 PeeringAgentConfiguration.kt 以@ConditionalOnExpression("${pollers.peering.enabled:false}")条件装配整个 peering 代理,即默认关闭、只有显式开启该 profile 时才注册PeeringAgentBean;同时它要求peerId(peer 的 ID)与poolName两个参数必须指定,否则直接抛出ConfigurationException。
  • 补充说明:在源码的配置属性类 PeeringAgentConfigurationProperties.kt 中,该"peer ID"字段被声明为peerId(README 的配置键写作id),两者的默认值依次对应intervalMs=5000、threadCount=30、chunkSize=100、clockDriftMs=5000、enabled=false。

参数说明表

参数默认值说明
pollers.peering.enabledfalse用于开启或关闭 peering
pollers.peering.poolName[REQUIRED]访问 foreign 数据库所使用的连接池名称,对应上文sql.connectionPools.foreign
pollers.peering.id[REQUIRED]peer 的 ID,每个数据库必须唯一
pollers.peering.intervalMs5000执行迁移的间隔(每一轮执行一次增量复制)。间隔越短延迟越低,但 CPU 与数据库负载越高
pollers.peering.threadCount30用于批量迁移的线程数。大数值只在最初的批量导入阶段有明显帮助;此后增量通常很小,超过 2 基本没有区别
pollers.peering.chunkSize100复制数据时的分块大小(单次最多修改的行数)
pollers.peering.clockDriftMs5000允许操作同一数据库的多个 orca 实例之间存在这么大的时钟漂移
pollers.peering.maxAllowedDeleteCount100单次最多删除的执行数。若删除量 Δ 超过该值则本轮不执行删除并累加错误指标。用于防止误删全部执行;可通过DynamicConfigService动态调整

三个关键参数的深入解读

  • pollers.peering.intervalMs:这是 peering agent 两次运行之间的间隔(即上一轮 agent 运行结束到下一轮开始的时间)。加上单轮运行本身的耗时,共同决定了"一条执行从一个数据库复制到另一个数据库"的延迟。README 给出了一个参考基准:在约有 600 万历史执行、任意时刻约 200 个活动执行的 MySQL orca 安装上,单轮 agent 运行约耗时 8 秒。数值越小 peering 延迟越低,但数据库负载越高。

  • pollers.peering.clockDriftMs:定义比较执行updated_at时间戳时的"容差因子"。由于时间戳由各实例(而非数据库本身)写入,实例之间无法保证时钟完全同步。例如抓取源库快照时最新updated_at是 1000,但另一个未完全同步的实例可能在快照之后又以时钟 998 修改了某条执行——如果没有容差,这条变更可能被漏掉。源码中,每轮完成后的游标会减去clockDriftMs(PeeringAgent.kt),正是这一容差的落地。

  • pollers.peering.maxAllowedDeleteCount:这是防止 peering agent 灾难性故障的安全阀——防止它(错误地)决定把本地数据库的全部执行删除。取值应大于常规操作中单次最大删除量,但远小于库中执行总数;README 建议的量级约为全部执行的 0.25%。该值既可作为静态配置,也可通过DynamicConfigService在运行时调整,源码中doMigrate每次都以pollers.peering.max-allowed-delete-count动态读取(默认 100),删除量超出即跳过删除并累加错误指标(PeeringAgent.kt)。

数据库访问层:MySqlRawAccess 的实现细节

SqlRawAccess是 peering 与数据库交互的抽象层(SqlRawAccess.kt),定义了getCompletedExecutionIds、getActiveExecutionIds、getDeletedExecutions、getStageIdsForExecutions、getExecutions、getStages、deleteStages、deleteExecutions、loadRecords等操作。当前唯一的实现是 MySqlRawAccess.kt,其中值得注意的工程细节:

  • 初始化时查询数据库max_allowed_packet并预留 8192 字节语句开销,用于在批量写入时按包大小自动拆分(MySqlRawAccess.kt);
  • 执行 ID 查询带 partition 约束(NULL或指定值,见上文);删除时由于 stages 表索引多、删除代价高,单次删除分块取min(chunkSize, 5)进一步收紧(MySqlRawAccess.kt);
  • 对"已完成执行 ID 查询"这类读操作做了 3 次、间隔 500ms 的重试封装(withRetry,MySqlRawAccess.kt);
  • 表名映射集中在 Utils.kt:pipeline 对应pipelines/pipeline_stages,orchestration 对应orchestrations/orchestration_stages。

自定义扩展:CustomPeerer

如果内置的"执行 + stages"复制不足以覆盖你的数据同步需求,可以实现 CustomPeerer.kt 接口并以 Spring Bean 暴露:

  • init(srcDb, destDb, peerId):让自定义 peerer 完成初始化,可通过srcDb.runQuery拿到 jooq 上下文对源库执行查询;
  • doPeer():在每一轮默认 peering 动作全部完成后被调用,返回true表示成功、false表示失败。

PeeringAgent初始化时会先调用init并捕获异常(失败则该自定义 peerer 不会被启用,同时累加pollers.peering.customPeerer.numErrors指标),每个周期最后执行doPeer(PeeringAgent.kt)。

对 foreign execution 执行操作:orca-interlink

README 明确列出了用户可以通过 UI/API 对一条执行发起的操作(orca 会基于这些操作变更执行状态):

  • cancel(取消一条执行)
  • pause(暂停一条执行)
  • resume(恢复一条执行)
  • pass judgement(对一条执行进行判定)
  • delete(删除一条执行)

这些操作必须发生在拥有该执行的集群/实例上。当本地集群面对一条 foreign execution 时,不能直接操作,而是通过orca-interlink模块把意图转达给实际拥有者。从仓库源码看,orca/orca-interlink 模块提供了对应的消息事件类型:CancelInterlinkEvent、DeleteInterlinkEvent、PauseInterlinkEvent、ResumeInterlinkEvent、PatchStageInterlinkEvent、RestartStageInterlinkEvent(统一继承InterlinkEvent),并带有Interlink消息发送抽象与 AWS 实现InterlinkAmazonMessageHandler(使用 SQS 等云消息服务承载跨区域通信的推断来自该实现类的命名与包结构)。同时,InterlinkConfigurationProperties.java 暴露了interlink.flagger相关配置(enabled=true、maxSize=32、threshold=8、lookbackSeconds=60),用于对消息做标记/限流控制。README 在该节标注 "TBD",说明远程操作部分仍在演进中。

监控:Emitted Metrics 与推荐告警

peering agent 通过 PeeringMetrics.kt 发射如下指标(均带有peerId标签,部分带有executionType/state标签),可用于监控 peering 系统的健康状况:

指标说明
pollers.peering.lagTimer(秒),度量执行单轮迁移循环的耗时;该值 + agent 的intervalMs即为有效延迟。应是一个相当平稳的数值
pollers.peering.numPeeredCounter,已复制执行的数量(应保持平稳——即与活动执行数量大致相当)
pollers.peering.numDeletedCounter,已删除执行的数量
pollers.peering.numStagesDeletedCounter,复制过程中删除的 stage 数量,纯信息性指标
pollers.peering.numErrorsCounter,执行复制过程中遇到的错误数(应对其设置告警)

源码中除上述外还发射pollers.peering.customPeerer.numErrors(自定义 peerer 的错误计数),并且lag同时记录单种执行类型(PIPELINE/ORCHESTRATION)与整体(OVER_ALL)两种维度(PeeringMetrics.kt)。

如果使用了 peering 功能,README 建议为以下指标配置告警:

  • pollers.peering.numErrors > 0
  • pollers.peering.numPeered == 0持续一段时间(取决于你的活动执行稳态规模)
  • pollers.peering.lag > 60持续一段时间(约 3 分钟)

运行时动态控制:Dynamic properties

以下动态属性可通过DynamicConfigService在运行时控制,无需重启集群:

属性默认值说明
pollers.peering.enabledtrue设为false时关闭全部 peering
pollers.peering.<PEERID>.enabledtrue设为false时关闭指定 peer ID 的全部 peering
pollers.peering.max-allowed-delete-count100单轮 agent 运行允许删除的最大执行数

这三项开关在PeeringAgent.tick()中逐轮生效:前两项经dynamicConfigService.isEnabled(...)判断(默认true,见 PeeringAgent.kt),删除阈值则每次动态读取(见 PeeringAgent.kt)。测试 PeeringAgentSpec.groovy 专门验证了"全局禁用时不触发任何查询、单 peer 禁用时不触发该 peer 查询"的动态开关行为。

Caveats:使用前提与注意事项

  • 仅支持 MySQL:目前 peering 只支持 MySQL。要扩展到其他数据库引擎,只需为对应引擎实现一个新的SqlRawAccess(参见 SqlRawAccess.kt),装配逻辑中PeeringAgentConfiguration目前也仅在 jooq dialect 为MYSQL时构建MySqlRawAccess,否则抛出UnsupportedOperationException(PeeringAgentConfiguration.kt)。
  • 建议只用一个实例运行 peering agent/profile:目前还没有跨实例锁(cross instance locking),多个实例同时跑 peering 可能造成重复复制或竞争。该限制未来有望改进。值得说明的是,源码层面PeeringAgent继承自AbstractPollingNotificationAgent并接收NotificationClusterLock,锁基础设施已接入,但 README 明确提示当前阶段仍以单实例运行为准。

尚未实现的能力

README 明确标注了两个待办方向,属演进中的规划而非当前可用能力:

  • 接管执行所有权(Taking ownership):从 peer 手中取得某条执行的所有权并继续运行,即"failover/接管"能力;
  • 对 foreign execution 执行操作:README 的 "TBD" 说明该能力仍在开发中,需要与orca-interlink的远程消息通道协同落地。

从代码结构看,orca-peering模块目前的稳定能力集中在"执行数据的单向镜像 + 删除传播 + 监控与动态开关",多区域场景下的完整读写闭环仍需与orca-interlink一起演进。

  • 后端
  • DevOps
  • 云原生
  • 微服务

【免费下载链接】spinnaker

Spinnaker is an open source, multi-cloud continuous delivery platform for releasing software changes with high velocity and confidence.

项目地址:https://gitcode.com/gh_mirrors/sp/spinnaker
点击查看免费下载

相关推荐

上一篇:5分钟搞定Notion免费版PDF导出:告别复制粘贴的高效工具
下一篇:高效PDF文献翻译工具:Zotero PDF Translate功能解析与实用指南

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

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

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

立即咨询