跳到正文
分布式定时任务重复执行故障复盘与幂等设计
分布式定时任务重复执行故障复盘与幂等设计

分布式定时任务重复执行故障复盘与幂等设计

从商户日结任务重复执行出发,分析故障转移、分布式锁、批次记录、租约状态机、唯一约束与跨系统幂等的职责边界。

故障概述#

一个商户日结任务部署在三台执行节点上。某天调度中心日志显示只产生了一次业务触发,但其中两台节点先后完成了全量结算,最终出现重复结算单、部分账户重复收款和资金流水无法对齐。这里最重要的现象不是“两台机器都收到任务”,而是同一个业务批次产生了两次可以落地的资金效果。

调度系统中的“触发一次”、Worker 的“执行一次”和业务系统的“生效一次”是三件不同的事。一次触发可能因为超时、心跳丢失、故障转移或人工重跑产生多个执行尝试;多个尝试本身并不一定构成事故,真正的事故是业务层没有把重复尝试压缩成同一个结算结果。

资金类故障发生后,第一步不是立即重跑,而是停止新的调度与自动重试,冻结仍可能继续写入的下游通道,按商户、结算日期和业务类型识别重复记录,并同时保存调度日志、Worker 日志、数据库事务记录和外部请求流水。只有先止住新的副作用,后续对账、冲正和根因分析才有稳定的证据基础。

先确认重复发生在哪一层#

“调度中心只触发了一次”只能说明业务触发记录可能只有一条,不能证明执行尝试只有一次。复盘时至少要区分下面三层标识:

层次建议标识要回答的问题
调度触发trigger_id计划任务在业务时间点被创建了几次
执行尝试attempt_idworker_id同一触发被派给了哪些节点,各自何时开始和结束
业务操作operation_id、业务幂等键数据库和外部系统最终接受了几次业务效果

根据“两台节点都跑完整批、第三台空闲”的现象,超时故障转移是一个合理假设:节点 A 开始执行后长时间没有完成或心跳异常,调度中心将同一触发重新派给节点 B;A 实际没有终止,恢复后继续写入,B 也从头执行,最终形成并发或前后重叠的两次全量处理。广播路由、人工补跑和错误的失败重试也可能产生相似结果,因此不能只凭现象宣布根因,还需要用 trigger_id、派发时间、心跳、超时阈值和 Worker 起止时间把链路对齐。

分片不平衡不是这类现象的首要解释。分片分成 1/1/97 会让某个节点负载过高、执行变慢,却不会自然让两个节点各自处理一份全量数据;只有分片条件未进入查询、广播被误当成分片,或者分片任务在重派时丢失了范围,才会进一步演变为重复全量执行。

调度层为什么不能单独保证业务只生效一次#

调度器面对节点失联时必须在“可能漏执行”和“可能重复执行”之间选择。节点可能在提交业务事务之前崩溃,也可能在提交之后、回报成功之前崩溃;如果调度器看不到业务提交结果,这两种状态在外部表现上几乎相同。不重试会留下漏处理风险,重试则会产生重复尝试,因此可靠调度通常更接近 At-Least-Once:允许再次执行,再由业务层识别重复。

这并不意味着任何范围内都无法实现 Exactly-Once。单个数据库事务、支持事务消息的受控系统,或者共享同一幂等协议的调用链,都可以在明确边界内提供一次性语义;问题在于调度器、业务数据库和银行通道无法自动共享一个原子提交。更准确的工程目标是:接受执行尝试可能重复,在已经声明的业务范围内实现 effectively-once,让重复尝试最终得到同一份业务结果。

各层职责应当明确分开:

机制主要作用不能替代什么
调度层超时、重试、故障转移、持久化触发记录保证任务最终有机会被执行业务幂等和资金正确性
选择层分布式锁、Leader、分片租约降低同一时间的重复计算崩溃后的重复尝试
批次层持久化批次记录识别整批状态、恢复与审计批次中途失败后的逐项正确性
记录层业务幂等键、唯一约束、状态机保证每个业务对象只产生一份有效结果外部系统的副作用
界外层幂等请求号、状态查询、对账控制账户或第三方通道的重复调用对方不支持幂等时的绝对保证
运营层告警、暂停、人工审核、冲正处理未知结果和异常数据正常路径中的程序正确性

正确性地基:逐项幂等与可恢复状态机#

先定义真正的业务幂等键#

结算任务不能只用“今天是否跑过”判断重复,因为一次批次可能在第 5000 个商户处崩溃。更稳妥的做法是为每个结算操作建立持久化记录,并用业务自然键约束它只能创建一次。下面以 PostgreSQL 为例:

CREATE UNIQUE INDEX uq_settlement_operation
ON settlement_operation (
merchant_id,
settle_date,
settlement_type,
currency
);
INSERT INTO settlement_operation (
merchant_id,
settle_date,
settlement_type,
currency,
status
)
VALUES (:merchant_id, :settle_date, :type, :currency, 'PENDING')
ON CONFLICT (
merchant_id,
settle_date,
settlement_type,
currency
) DO NOTHING
RETURNING id;

如果插入成功,当前 Worker 获得一条新的操作记录;如果没有返回行,说明相同业务键已经存在,Worker 应读取原记录并根据状态决定跳过、查询还是恢复,而不是再创建一个新的结算号。(merchant_id, settle_date) 只是最小示例,真实系统还要考虑结算类型、币种、账期、渠道和重算版本,业务键少一个维度会误伤合法操作,多一个不稳定字段又会让重复请求绕过约束。

唯一约束只保证操作记录不重复,不会自动保证余额更新、明细写入和外部打款都正确。属于同一数据库的账务明细、内部余额和操作状态,应尽量在同一个本地事务中提交,并继续为账务分录使用 operation_id 唯一约束。第二个 Worker 即使再次进入,也只能读到或复用同一条业务操作,不能创建第二份有效账务结果。

状态机必须能够从卡死状态恢复#

简单的 UNSETTLED -> SETTLING -> SETTLED 能阻止两个 Worker 同时推进,但如果进程在写入 SETTLING 后崩溃,记录会永久卡住。可恢复状态机需要租约、版本号和明确的终态,例如:

PENDING / RETRYABLE
-> PROCESSING(lease_until, worker_id, version)
-> SUCCEEDED
-> FAILED_PERMANENT
-> UNKNOWN_EXTERNAL

Worker 领取任务时使用带条件的更新,并把递增版本作为 fencing token:

UPDATE settlement_operation
SET status = 'PROCESSING',
worker_id = :worker_id,
lease_until = now() + interval '5 minutes',
version = version + 1
WHERE id = :id
AND (
status IN ('PENDING', 'RETRYABLE')
OR (status = 'PROCESSING' AND lease_until < now())
)
RETURNING id, version;

没有返回行表示当前记录已被其他 Worker 持有或已经进入终态。后续内部写入和完成更新必须携带领取时得到的 version,这样租约过期后的旧 Worker 即使恢复,也不能用过期版本覆盖新 Worker 的状态。租约时间需要结合单项处理耗时设置并允许续期,但续期仍然只是活性机制,最终正确性依赖唯一业务键、版本检查和下游幂等。

锁、批次记录和分片各自负责什么#

分布式锁可以减少多个实例同时扫描和处理同一批次,但它不是正确性的承重墙。Redis 或 Redisson 锁可能因为进程长暂停、网络分区、续期失败或 TTL 到期而失去所有权,旧持有者随后仍可能继续运行;如果下游支持 fencing token,应拒绝小于当前版本的旧请求,如果不支持,就必须依赖业务操作的唯一键和幂等结果挡住重复副作用。

批次记录也应放在可审计的持久化存储中,而不是只依赖一个短 TTL 的 Redis done-marker。例如为 (job_type, biz_date, scope) 建立唯一批次行,记录 CREATED、RUNNING、PARTIAL、COMPLETED、RECONCILING、FAILED 等状态、统计数量和最后检查时间。批次完成标记适合让整批重跑快速退出,但只有当所有逐项操作进入可接受终态并通过必要对账后,批次才应标记为完成。

分片解决吞吐,不解决幂等。每个分片必须把范围条件真正带入数据查询,并记录分片编号、总分片数和输入版本;调度中心在重派分片时仍可能让新旧 Worker 重叠,因此每个商户操作仍要经过同一套唯一键和状态机。是否使用单节点、锁或分片,最终应是成本和效率选择,而不是正确性选择。

跨系统调用:未知结果比明确失败更危险#

账户系统、银行通道和第三方支付不在本地数据库事务中。最危险的情况不是收到明确失败,而是请求已经被对方执行,本地却因为超时没有收到响应。如果此时生成新的请求号并重试,就可能发生第二次扣款或打款。

跨系统调用至少需要下面四层约束:

  1. 使用本地 operation_id 派生稳定的幂等请求号,同一次业务重试始终发送同一个键。
  2. 在本地事务中先保存待发送事件或 Outbox 记录,再由独立发送器调用外部系统,避免数据库已提交但请求完全丢失。
  3. 收到超时或连接中断时进入 UNKNOWN_EXTERNAL,优先用幂等键或对方流水号查询原请求状态,不能直接当作失败重发。
  4. 定期按双方流水、金额、商户和日期执行对账;对方不支持可靠幂等或状态查询时,对账和人工处理就是必须存在的最终控制面。

幂等键只有在对方明确承诺“相同键代表同一业务操作”时才有效,还要确认键的保留时间、参数不一致时的处理方式,以及处理中请求是否允许并发重入。仅仅在 HTTP Header 中放一个随机 UUID,而每次重试都生成新值,不构成幂等。

一套可落地的结算流程#

综合前面的边界,一次日结可以按照下面的顺序推进:

调度触发(trigger_id, biz_date)
-> 创建或读取唯一批次(job_type, biz_date, scope)
-> 可选:获取批次锁,减少重复扫描
-> 为每个商户创建唯一 settlement_operation
-> Worker 通过租约 + version 领取操作
-> 本地事务:写账务分录、内部余额和 Outbox
-> 发送器使用 operation_id 作为外部幂等键
-> 明确成功:记录外部流水,操作进入 SUCCEEDED
-> 明确永久失败:进入 FAILED_PERMANENT,停止自动重试
-> 超时或结果不明:进入 UNKNOWN_EXTERNAL,查询与对账
-> 所有操作完成并核对数量、金额后,批次进入 COMPLETED

下面的故障注入应在上线前实际演练,而不是只做正常路径单元测试:

故障点期望行为
同一批次被三个节点同时启动可以产生多个尝试,但每个业务键只有一条有效操作
写入操作记录后 Worker 崩溃租约到期后由其他 Worker 使用新版本恢复
本地账务事务提交后响应丢失重试命中同一操作与唯一分录,不重复记账
分布式锁过期,旧 Worker 恢复旧版本无法完成状态更新,下游幂等拒绝重复效果
外部通道已处理但本地超时使用同一幂等键查询,不创建新的打款请求
批次完成后被人工再次触发批次记录快速返回,逐项唯一键继续兜底

验收不能只检查“最终没有报错”,还要核对操作记录数、账务分录数、外部流水数、总金额和状态分布是否一致。对资金任务来说,能够在指定故障点杀死进程并安全恢复,比连续跑通十次正常流程更有价值。

重试、单节点与运维边界#

重试策略应该依据错误语义,而不是简单分成“技术异常重试、业务异常不重试”。明确的参数错误、规则不满足和代码缺陷通常属于永久失败,应停止批次并告警;连接超时、限流和临时不可用可以使用带抖动的指数退避和次数上限;并发冲突或租约竞争可以短暂重试;外部结果未知则必须先查询和对账,不能盲目再发。任何自动重试都应复用原 operation_id,并设置最大尝试、总时间和人工升级条件。

关键窗口切换为单节点可以降低并发和操作面,是合理的风险控制,但不能替代业务幂等。单节点同样会重启、超时和被人工补跑;如果关闭多节点后重复效果才消失,说明系统只是缩小了触发概率,正确性边界仍然没有补上。调度中心集群也主要改善可用性,只有配合 Leader、持久化触发记录和去重协议才能减少中心自身的重复派发,业务效果仍应由下游兜底。

一次可审计的执行日志至少应包含 trigger_id、attempt_id、batch_id、operation_id、worker_id、shard、lock/fencing version、开始与结束时间、输入版本、影响行数、外部幂等键和外部流水号。这些字段让团队能够回答“谁在什么时候处理了哪条业务、使用了哪一版输入、结果是否被外部接受”,而不是事故发生后只剩一句“调度显示成功”。

结论#

分布式定时任务不需要追求“永远只有一个 Worker 运行”,而要接受重复尝试是正常故障模型的一部分。锁、单节点和分片可以减少浪费,批次记录可以加快恢复,真正让资金效果只出现一次的,是定义清楚的业务幂等键、数据库唯一约束、带租约和版本的状态机,以及外部系统认可的幂等请求号。

这套设计也有边界:唯一约束只能保护约束所在的数据库,状态机必须处理卡死和过期持有者,外部通道不支持幂等时仍然需要查询、对账和人工控制。把这些限制写进方案,比承诺一个没有范围说明的 Exactly-Once 更可靠。

复盘这类事故时,最有价值的问题不是“为什么调度器执行了两次”,而是“为什么第二次执行仍然能够产生新的业务效果”。前一个问题决定怎样减少重复,后一个问题才决定系统是否正确。

参考资料#

版权许可

CC BY-NC-SA 4.0 本作品采用知识共享署名-非商业性使用-相同方式共享 4.0 国际许可协议进行许可。

相关文章

s1oopX

登录 s1oopX