原子计数器为何会漂移以及对账如何修复此问题
通过原子增量维护的计数器可能会与真实的行数慢慢产生偏差。一个定期的校准作业会重新计算真实值,并使其收敛。
通过原子增量维护的派生计数器会与其真实来源产生偏差。这并非因为增量操作本身是错误的,而是因为原子性只保证单个操作的完整性,而不保证与一个独立记录的一致性。崩溃、重试和部分失败都会让这个数字产生一点偏差,并且这些错误会累积起来。解决方法不是使用一个更大的锁,而是将该计数器视为一个快速缓存,并运行一个周期性作业,从真实来源重新计算真实值并将其写回。这样,即使该数字曾短暂地出错,它最终也会是正确的。
什么是“漂移”,以及为什么原子性无法阻止它
假设你将评论存储为行,并且你想在每篇帖子上显示评论数,而无需在每次加载页面时都扫描评论集合。显而易见的优化是使用派生计数器:在帖子上保留一个 comment_count 字段,并在每次创建评论时增加它。
// Fast path: bump the derived counter whenever a comment is created.
// This single update is atomic for this one document. Nothing here ties
// it to the actual number of comment rows that exist.
_, err := counters.UpdateOne(ctx,
bson.M{"_id": postID},
bson.M{"$inc": bson.M{"comment_count": 1}},
options.Update().SetUpsert(true),
)
$inc 是原子性的。两个并发的增量操作不会像“读取-修改-写入”那样丢失更新。因此,人们很容易得出计数器是正确的结论。事实并非如此,其原因在于对原子性所能带来的好处存在概念上的错误。
原子性是单个操作的属性。它指的是这一个增量操作要么完全发生,要么完全不发生,不会有任何交错操作来破坏其值。它并没有说明增量操作的序列是否与另一个集合中的行数相匹配。计数器和评论行是两个独立的状态片段,由两个独立的操作更新。它们之间的一致性并非这两个操作中任何一个所具有的属性。任何时候,只要这两个写入操作不属于同一个原子单元(在文档存储和不同集合之间,它们通常不属于同一个原子单元),它们之间就可能出现不一致。这种不一致就是漂移。
偏差从何而来
偏差并非单一的 bug。它是一系列在行写入和计数器写入之间产生的小间隙,每个间隙都会将计数推向一个可预测的方向。了解这些来源可以让你知道对账作业需要纠正什么,以及当真正的 bug 出现时,指标应该是什么样子。
| 原因 | 机制 | 方向 |
|---|---|---|
| 两次写入之间发生崩溃 | 行已插入,但进程在 $inc 操作完成前终止 |
计数偏低 |
| 以相反顺序崩溃 | 计数器已递增,但随后的行插入失败或回滚 | 计数偏高 |
| 至少一次重试 | 消息被传递了两次,同一事件使计数器递增了两次 | 计数偏高 |
| 缺少递减操作 | 行被删除或软删除,但没有运行相应的递减操作 | 计数偏高 |
| 回填或手动修复 | 直接在数据库中导入或修复行,绕过了递增路径 | 计数偏低 |
| 软删除和恢复的顺序错乱 | 删除操作递减了计数,但随后的恢复操作忘记递增,反之亦然 | 两者皆可能 |
有两点值得注意。首先,这些错误不会相互抵消。一个既会漏掉某些递增操作,又会重复应用另一些递增操作的系统,其结果并不会平均下来变得正确。它最终会落在一个错误的值上。其次,错误的方向具有诊断价值。一个只会出现偏高情况的计数器,指向的是重复处理或缺少递减操作。一个只会出现偏低情况的计数器,指向的是在行提交后失败的写入路径。修复计数的对账作业也能告诉你存在哪种问题,前提是你在弥合偏差前记录下它。
将计数器视为缓存,而非事实来源
要让这一切变得可控,其心智模型是不要再将计数器视为一个事实。它是对一个事实的缓存。这个事实就是评论行的集合。计数器是一个预先计算好的答案,其对应的问题是你不想在每次页面加载时都去运行的。
一旦将计数器视为缓存,就需要遵循两条规则。缓存允许存在过时数据,因此短暂的错误计数是可接受的,而非一场危机。并且,缓存必须有办法从其缓存的内容中重建,因为一个无法重新生成的缓存只是不可靠的主数据。事实来源保持其权威性。它就是行的集合,通过直接查询获得,并完全忽略计数器。任何时候你需要一个可以赖以决策的数字时,你就应该去计算行数。计数器用于廉价读取的场景,在这些场景下,微小的瞬时错误不会造成任何损失。
这种重新定义使得整个方法变得诚实。你不是在承诺一个永远正确的计数器,然后又悄悄地无法兑现。你承诺的是一个快速的近似计数器,外加一个你随时可以回退依赖的事实来源,以及一个能保持这两者数据接近的作业。
重新计算并覆盖的对账作业
对账是一个不在热路径上的定时作业,它会根据数据行重新计算真实计数值,并将其写入计数器。由于它基于计时器运行,而非针对每个请求运行,因此它可以执行热路径所避免的昂贵操作,即实际的计数。
朴素版本只有三行:统计行数,写入数值,完成。
// Naive reconciliation: recompute from the source of truth and overwrite.
// Correct in isolation, but see the race in the next section.
trueCount, err := comments.CountDocuments(ctx, bson.M{
"post_id": postID,
"deleted_at": bson.M{"$exists": false},
})
if err != nil {
return err
}
_, err = counters.UpdateOne(ctx,
bson.M{"_id": postID},
bson.M{"$set": bson.M{
"comment_count": trueCount,
"reconciled_at": time.Now(),
}},
)
这会一次性纠正所有类型的偏差,无论是过多还是过少,因为它完全不信任旧值。它会根据行重新派生出这个数字。无论旧计数器的值是多少,无论对错,都会被丢弃。
这里存在一个风险,这也是为什么朴素版本不是最终版本的原因。在 CountDocuments 返回和 $set 落地之间的片刻,可能会有新评论到达。热路径会为它们增加计数器的值。然后,$set 会用一个在这些增量操作存在之前计算出的数字覆盖计数器,导致这些增量丢失。本意是修复偏差的对账操作,反而制造了新的偏差。
根据水印进行协调,以保留正在进行的写入
避免覆盖并发写入的简洁方法是,只协调数据的已落定前缀,而不去动近期的写入。这需要对计数器的存储方式做一个小改动。将其拆分为两个字段:一个由协调过程拥有的 base_count,以及一个由热路径拥有的 live_delta。显示的数字是它们的和。
热路径不再触碰已协调的值。它只递增实时增量。
// Hot path now only touches live_delta. Reconciliation never overwrites
// this field, so a concurrent increment can never be clobbered.
_, err := counters.UpdateOne(ctx,
bson.M{"_id": postID},
bson.M{"$inc": bson.M{"live_delta": 1}},
options.Update().SetUpsert(true),
)
读取操作会将这两个字段相加:
displayed := doc.BaseCount + doc.LiveDelta
对账过程会选择一个水印,这是一个足够久远的过去时间戳,以至于不会有任何新行以更早的创建时间写入。在实践中,这意味着它要早于你的写入可见性延迟和最早的未结事务,所以几秒钟的时间通常就足够了。它会统计截至该水印的真实情况,将其设置为新的基准,并从实时增量中精确地移除新基准已包含的增量部分。所有这些写入操作都在一个原子更新中完成,因此字段之间永远不会出现不一致的情况。
// Settled boundary: no new row can appear with a timestamp older than this.
watermark := time.Now().Add(-30 * time.Second)
baseTrue, err := comments.CountDocuments(ctx, bson.M{
"post_id": postID,
"deleted_at": bson.M{"$exists": false},
"created_at": bson.M{"$lte": watermark},
})
if err != nil {
return err
}
// Atomic pipeline update. Install the recomputed base, and subtract from
// live_delta the number of rows the base has just absorbed (baseTrue minus
// the old base_count). Increments for rows after the watermark stay in
// live_delta untouched.
_, err = counters.UpdateOne(ctx,
bson.M{"_id": postID},
bson.A{
bson.M{"$set": bson.M{
"live_delta": bson.M{"$subtract": bson.A{
"$live_delta",
bson.M{"$subtract": bson.A{baseTrue, "$base_count"}},
}},
"base_count": baseTrue,
"reconciled_at": "$$NOW",
}},
},
)
其安全的原因是,协调 (reconciliation) 和热路径 (hot path) 现在写入的是不相交的字段。基准值 (base) 从行 (rows) 中重新派生,因此旧基准值中累积的任何漂移 (drift) 都会被清除。实时增量 (live delta) 仅持有最近的增量,一旦这些增量也超过了水位线 (watermark),下一个周期就会将它们合并到基准值中。协调期间的并发写入会进入实时增量,绝不会处于覆写的路径上。计数器无需锁也无需暂停写入即可收敛。
如果你想要最简单的版本,并且可以容忍稍大的瞬时误差,可以保留单字段覆写,并接受与协调 (reconcile) 竞争的写入可能会短暂丢失,然后在下一次运行时得到纠正。对于写入频率非常低的计数器来说,这是一个合理的选择。当写入速率高到足以让覆盖冲突 (clobbering) 变得重要时,你才需要采用水位线和拆分字段的方案。
测量偏差,从而发现真正的错误
对账所丢弃的数值比其写入的数值更有价值。在设置新的基准值之前,你知道旧值和真实值。这个差值就是偏差,它是你写入路径健康状况的直接证据。
drift := baseTrue - oldBaseCount
metrics.Observe("counter_drift", float64(drift),
"collection", "comments")
将其作为 gauge 或 histogram 发出,并按计数器类型进行标记。现在你就得到了一个信号,它会讲述一个故事。漂移值在零附近徘徊并保持稳定,意味着热路径基本正确,而对账过程只是在清理罕见的崩溃。在两次运行之间不断增长的漂移意味着增量正在以比你想象中更快的速度丢失或被重复应用,其正负号会告诉你具体是哪种情况。一次突然的尖峰与某次部署或事故相吻合,并直接指向破坏了写入路径的那个变更。
如果没有这个,对账过程会隐藏你的 bug。它每晚都会悄悄地掩盖损坏的增量路径,而你永远不会知道路径已损坏,因为显示的数字到早上总是看起来没问题。该指标将无声的修正转变为一个警报。这个任务仍然在修复症状,但现在它也报告了病因。如果漂移值超过了你关心的阈值,就呼叫相关人员,因为此时计数器的漂移速度已经超过了夜间任务可以安全掩盖的速度。
当近似计数器不够用时
这整个方法是用即时准确性来换取廉价的读取和自我修复。这种权衡对于一大类计数器是正确的,而对于某一特定类型则是错误的,而它们之间的分界线就是金钱。
对于统计数据、排名、显示计数、总浏览量、评论数、点赞总数、粉丝数而言,一个在几分钟内有几个误差的计数对用户是不可见的,并且不产生任何成本。在分布式系统中强制要求这些计数精确和实时,意味着在你最繁忙的写入操作的热路径上需要一个事务或一个锁,其成本会随着流量的增长而增长。一个带有定期对账功能的快速近似计数器是务实的答案,“最终正确”则是一个足够强的保证。
对于任何计数会影响到具有实际后果的决策的场景,情况就反过来了。钱包余额、你强制执行的付费配额、可能超卖的库存、绝不能为负的信用分类账:这些都不能是延迟对账的缓存。对于这些场景,真实源头应该位于读取路径上。你或者在做决策时进行计数或读取权威余额,或者让扣减操作本身针对约束条件是事务性的。一个短暂错误的计数对于点赞总数来说没问题,但对于金额来说是不可接受的。
还有一个需要权衡的运营成本。计算行数不是免费的,一个对巨大集合中的每个计数器都重新计算的对账作业,其本身就可能成为一个负载问题。通常的解决方案是:只对自上次运行以来发生变化的文档进行对账,按时间窗口进行扫描,错开执行以避免一次性计算所有内容,并且运行的频率要足够高以保持偏差较小,但又要足够低以使扫描保持廉价。这种调优是运行对账计数器的真正工作所在。设计很简单。随着数据增长,如何保持作业的成本可控,这才是工程技术的用武之地。