用于“先退款后恢复”重试的幂等信用分类账
失败的任务会获得点数退款。然后重试成功。在不重复收费的情况下收回这些点数,需要一个幂等账本和原子性保障。
一个按使用量计费的系统会在作业失败时退还额度。令人头疼的情况是,同一个作业稍后在重试时恢复并成功。这样一来,之前的退款就错了,你必须在不重复计费的情况下收回这些额度。解决方案是一个由反向操作对组成的只追加账本,其中每一次退款和每一次收回都是一次基于当前状态守护的单一条件写入。本文将展示该账本的结构、阻止重复退款和重复计费的原子性守护机制,以及重试队列如何在不陷入无限循环的情况下为其提供数据。
先退款后恢复问题
在按量计费系统中,用户操作会消耗点数:渲染视频、运行转录、生成图像。常见的设计是预先收费,然后运行作业,并在失败时退款,这样用户就无需为他们从未收到的工作付费。麻烦就从这笔退款开始。
分布式作业的失败过程并不干净利落。一个 worker 在 30 秒时超时,将作业标记为失败,并退还 40 点数。但正在执行工作的 GPU 实际上在 32 秒时完成了任务,结果在超时触发两秒后存入了存储。该作业并未失败,它恢复了。现在,账本上显示了一笔为存在且将被交付的工作所做的退款。
朴素的实现会让情况变得更糟。如果失败处理和成功处理是两个没有共享状态的独立代码路径,你就会遇到两种 bug 之一。同一作业的两个失败信号(一次超时加上一个稍后的错误回调)各自触发退款,导致用户为一个 40 点数的作业获得了 80 点数。或者,一次退款后紧跟着一次恢复,然后第二次恢复信号又各自重新计费,导致用户为一个作业被收费两次。在一个以正确性为核心的领域,这两种情况都是正确性故障。
其根本原因与大多数分布式系统 bug 背后的原因相同:一个操作分两步完成,先读取状态,再写入状态,而在这两者之间的间隙中,另一个 worker 可能会基于同样的陈旧读取进行操作。
将每次变更都建模为一组逆向操作对
第一个决定是绝不修改或删除之前的账本条目。账本是只追加的。一个被收费、退款、然后又收回的任务会产生三个不可变的行:
- 收费:-40
- 退款:+40
- 收回:-40
用户的余额是这些行的总和,而不是一个你可以覆盖的可变计数器。每个条目都带有 job_id、type、amount 和 created_at,因此每笔信用的完整历史都是可审计的。当支持人员询问用户为何拥有当前余额时,你只需重放这些行。
退款和收回互为逆向操作,这使得恢复操作得以表达。只有在先发生了退款且该退款尚未被收回时,收回操作才有意义。这个条件不是你在应用程序代码中用 if 语句来检查的。它本身就是写入操作的保障条件。
防止双重退款的原子性防护
双重退款 bug 是一个经典的“检查时到使用时”(time-of-check-to-time-of-use)竞争条件。工作进程 A 读取到 is_refunded = false,并决定退款。在它写入之前,工作进程 B 也读取到同样的 is_refunded = false,也决定退款。两者都执行了写入操作。用户收到了两次退款。
弥合这个间隙意味着检查和写入必须是同一个不可分割的操作。在支持条件更新的文档存储中,这会是一个单独的 UpdateOne 操作,其筛选条件是检查,其更新内容是写入。这个标志位存在于作业文档上,并且只有将它从 false 翻转为 true 的那个工作进程才被允许追加分类账行。
// Refund only if this job has not already been refunded.
// The filter and the $set are one atomic operation, so two concurrent
// failure signals cannot both win.
now := time.Now()
res, err := jobs.UpdateOne(ctx,
bson.M{"_id": jobID, "is_refunded": false},
bson.M{"$set": bson.M{
"is_refunded": true,
"refunded_at": now,
}},
)
if err != nil {
return err
}
if res.ModifiedCount == 1 {
// We won the guard. Append the ledger row exactly once.
appendLedger(ctx, jobID, "refund", cost, now)
}
ModifiedCount 返回值为 1 的工作进程赢得了守卫,并追加了退款。所有其他工作进程都会看到 ModifiedCount 返回 0,因为一旦标志被设置,筛选器就不再匹配,所以它们不执行任何操作。数据库会为你序列化条件写入,因此任何程度的并发都不会产生两次退款。
当作业恢复时收回额度
恢复操作使用镜像守卫。当一个已退款的作业稍后成功时,你会收回额度。条件是作业已被退款(is_refunded = true)且尚未被收回(is_reclaimed = false),并且检查和写入同样是一个操作。
// Reclaim only if the job was refunded and not yet reclaimed.
now := time.Now()
res, err := jobs.UpdateOne(ctx,
bson.M{"_id": jobID, "is_refunded": true, "is_reclaimed": false},
bson.M{"$set": bson.M{
"is_reclaimed": true,
"reclaimed_at": now,
}},
)
if err != nil {
return err
}
if res.ModifiedCount == 1 {
appendLedger(ctx, jobID, "reclaim", cost, now)
}
其精妙之处在于普通成功路径上发生的情况。一个首次尝试即成功的任务从未被退款,因此 is_refunded 仍为 false,筛选器不匹配,ModifiedCount 为 0,并且不会写入任何回收行。这完全正确:因为没有任何东西被退款,所以也就没有任何东西需要回收。成功处理器无条件地运行相同的回收调用,而由防护条件来决定该调用是否适用。无论是否发生过退款,也无论是因为重复投递而运行一次还是五次,该成功路径都是幂等的。
请注意,回收金额来自任务上记录的 cost,而不是来自重试报告的任何内容。将金额固定为原始费用意味着一个恢复的任务回收的正是它所退款的金额,绝不会是其他数字。
计量任务的状态转换
两个布尔标志 is_refunded 和 is_reclaimed 定义了一个小型状态机。其中的每条有效路径都保持分类账平衡,上面的每个守卫都是此状态机的一条边。
| 状态 | is_refunded | is_reclaimed | 净分类账影响 | 有效的下一转换 |
|---|---|---|---|---|
| 已计费,运行中 | false | false | -cost | 成功(保持),或失败(退款) |
| 首次尝试成功 | false | false | -cost | 终止 |
| 失败,已退款 | true | false | 0 | 恢复(回收),或重试 |
| 已恢复,已回收 | true | true | -cost | 终止 |
两个终止状态都以 -cost 的净影响结算,这是正确的:用户为工作恰好支付了一次。已退款但未回收的状态以 0 结算,这对于确实失败且从未交付的工作是正确的。没有路径可以达到 -2 倍成本或在交付后达到 0,因为每条边都由一个只有一个写入者可以翻转的标志来控制。
使用退避策略重新排队失败的工作
一个退回的作业并非总是彻底失败。如果失败是由暂时性依赖问题(例如,一个掉线的 GPU 节点或一个超时的存储调用)引起的,那么该作业会回到重试队列,而不是进入“墓地”队列。使用 RPush 命令将作业推送到 Redis 列表的尾部,可以实现一个简单的先进先出(FIFO)重试通道,并且每次重新入队都会携带一个递增的尝试次数和一个“不早于”时间戳,这样出现故障的依赖项就不会在紧密循环中被反复冲击。
// Re-queue a failed job with exponential backoff, capped.
attempt := job.Attempts + 1
backoff := time.Duration(
math.Min(
float64(baseDelay)*math.Pow(2, float64(attempt)),
float64(maxDelay),
),
)
job.Attempts = attempt
job.NotBefore = time.Now().Add(backoff)
payload, _ := json.Marshal(job)
rdb.RPush(ctx, "jobs:retry", payload)
消费者会跳过任何 NotBefore 在未来的作业,不作处理将其重新入队,因此无需单独的延迟队列机制即可实现延迟。以 maxDelay 为上限的指数级增长意味着第一次重试很快,但一个持续失败的依赖项会退避到慢速轮询,而不是采用一个会消耗其正在等待的资源的忙碌循环。
死信守卫阻止无限重试
仅靠退避机制本身无法阻止一个永远不会成功的作业。如果没有上限,一个永久损坏的作业会永远在退款、重新入队、失败、退款的循环中,并且由于退款守卫是幂等的,它不会造成财务损失,但它会浪费工作进程并使队列变得混乱。尝试次数也是这个上限。
// Terminal guard: past the cap, dead-letter instead of retrying.
if job.Attempts >= maxAttempts {
rdb.RPush(ctx, "jobs:dead-letter", payload)
return
}
死信列表是用于人工检查的暂存区,而不是一个自动重试通道。没有任何东西会定时地消费它。同样重要的是,is_reclaimed 和 is_refunded 这两个标志也充当防止重复投递的幂等性密钥。如果一个作业从重试队列中被投递了两次(至少一次投递保证了这种情况会发生),那么第二次投递会命中一个其标志已反映其最终状态的作业,并且受保护的写入将匹配不到任何内容。重新处理一个已经结算的作业不会改变余额。队列可以任意次地投递一条消息,而账本依然保持正确,因为正确性存在于条件写入中,而非投递保证中。
权衡:账本复杂性与正确性
简单的模型只有一行逻辑:失败即退款,完成。状态更少,代码更少,无需推理。然而,一旦失败不是最终状态,这个模型就错了,而这在分布式系统中很常见。每一次与缓慢成功的任务竞争的工作进程超时,每一次重复的回调,每一次恢复成功的重试,都会悄无声息地打破这个简单的模型,而悄无声息的计费错误是代价最高昂的那种。
账本模型需要你付出两个标志位、一对反向操作、一个带退避机制的重试队列以及一条死信路径的代价。这是实实在在的复杂性,要将其牢记于心也并非易事。问题在于,业务领域是否值得这样做。对于一个免费套餐功能,少数积分的计算错误不会对任何人造成损害,那么就跳过所有这些,容忍偶尔的偏差。对于付费积分,其余额就是金钱形式,用户发现重复收费就会产生支持工单和信任问题,在这种情况下,只追加的审计日志和原子防护,在第一次遇到已退款任务又恢复执行时,就足以证明其价值了。让整套机制生效的规则很简单:让每个账本操作都幂等,并为每个状态转换设置一个只允许单个写入者通过的原子防护。