分布式数据一致性对账

分布式数据一致性对账:对账任务设计、差异检测与修复、对账系统架构、与事件溯源结合,从对账模型到生产实践

无论我们采用 2PC、TCC、Saga 还是本地消息表,最终一致性方案都无法百分之百保证数据在任意时刻都相同。网络抖动、消息丢失、重复消费、宕机恢复都会让上下游数据库出现静默差异。数据对账(Reconciliation) 是兜底的"最后一道防线":定期或实时地比较两份数据,发现差异并自动修复,把最终一致性从"理论承诺"变成"可验证事实"。

1. 为什么需要数据对账

1.1 分布式数据不一致的来源

来源典型场景不一致类型
网络分区事务提交成功但响应丢失,发起方误判失败缺数据 / 多数据
消息乱序与重复MQ 的 at-least-once 语义,消费端重复处理重复数据
部分成功跨服务多步操作中途宕机,补偿未完成缺数据
缓存与库异步更新缓存写库成功但同步失败缓存漂移
存储复制延迟主从延迟、跨机房复制中断副本不一致
人工运维失误误删数据、错误回滚任意差异

这些不一致大多不会立刻暴露,而是在对账、结算、审计时才发现。对账的目标就是主动发现并修复它们,避免差异长期累积。

1.2 对账在一致性体系中的位置

一致性保障分层:
  第 1 层:事务协议(2PC/TCC/Saga/Seata)—— 保证正常路径下的正确性
  第 2 层:消息可靠性(重试/幂等/本地消息表)—— 保证异步路径的最终一致
  第 3 层:数据对账(比对/检测/修复)—— 兜底一切未覆盖的静默差异
              ↑ 本文重点

前两层解决的是"设计上可控"的差异,对账解决的是"事实上可能漏掉"的差异。对账不是替代事务方案,而是与它们配合形成闭环,正如 https://plumephp.com/distributed-transactions/ 中"最大努力通知"用于对账场景一样。

1.3 对账的适用场景

  • 金融结算:订单、支付、账单、清算,每一分钱都要对上
  • 跨系统数据同步:主库与数据仓库、搜索索引、缓存之间的同步校验
  • 消息链路验证:消息发出量与消费量的核对,判断是否有积压或丢失
  • 库存与账务:扣减流水与库存余额的对账,防止超卖或盘差
  • 合规审计:留存对账记录,满足监管对账要求

2. 对账核心模型

2.1 账本抽象:源头与目标

对账的本质是比较两套数据视图。我们把数据来源称为"源头账本(Source Ledger)",被核对方称为"目标账本(Target Ledger)"。两个账本都需要具备:

  • 可枚举:能够按主键或分片遍历全部记录
  • 可标识:每条记录有稳定唯一键(业务单号 + 序号)
  • 可比较:存在一个可计算的比对指纹(版本号、哈希、金额等)
source_ledger (订单库)          target_ledger (账务库)
┌─────────────────────┐        ┌─────────────────────┐
│ order_no | amount   │        │ order_no | amount   │
│ A001     | 100.00   │  ────► │ A001     | 100.00   │  ← 一致
│ A002     | 50.00    │  ────► │ (缺失)              │  ← 差异:缺数据
│ A003     | 20.00    │  ────► │ A003     | 25.00    │  ← 差异:金额不符
│          │          │  ────► │ X999     | 10.00    │  ← 差异:多数据
└─────────────────────┘        └─────────────────────┘

2.2 对账任务模型

一个对账任务由如下要素定义:

public class ReconciliationTask {
    private String taskId;            // 任务唯一 ID
    private String sourceLedger;      // 源头账本标识
    private String targetLedger;      // 目标账本标识
    private String compareKey;        // 比对键,如 order_no
    private List<String> compareFields; // 需要比对金额、状态等字段
    private String mode;              // FULL / INCREMENT / REAL_TIME
    private int shardCount;           // 分片数
    private long watermark;           // 对账水位(毫秒时间戳)
    private String diffHandler;       // 差异处理器(自动修复策略)
    private Schedule schedule;        // 调度策略:每日 / 每小时 / 持续
}

2.3 对账水位与周期

对账不能每次全量扫描全表。增量对账依赖水位(watermark):

  • 源头账本和目标账本都维护 updated_at / version 字段
  • 对账任务记住上次处理到的水位 W,本次只比较 updated_at > W 的记录
  • 完成后将水位推进到 W'
  • 周期性(如每日一次)跑一次全量对账兜底,防止增量阶段漏掉更新
时间轴:
  W0 ───────────────► W1 ───────────────► W2 ───────────────► W3
  增量对账区间1        增量对账区间2         增量对账区间3
       └── 每日 03:00 插入一次全量对账,扫描全表做最终兜底

3. 对账任务设计

3.1 分片扫描

两张表动辄上亿行,单任务顺序扫描太慢。按 compareKey 哈希分片,并行扫描:

func (t *ReconciliationTask) ShardKeys(total int) []Shard {
    shards := make([]Shard, 0, t.ShardCount)
    for i := 0; i < t.ShardCount; i++ {
        shards = append(shards, Shard{
            TaskID:   t.TaskID,
            ShardNo:  i,
            ShardTotal: t.ShardCount,
            // 每片扫描 compare_key % shardCount == i 的记录
            Range:    "WHERE MOD(CRC32(compare_key), shard_total) = shard_no",
        })
    }
    return shards
}

分片任务可以提交到任务队列并行执行,每个分片独立记录进度(shard_id + offset),失败的分片可以单独重跑,不影响其他分片。

3.2 批次与游标

每个分片内部按主键游标分批拉取:

SELECT order_no, amount, version
FROM orders
WHERE MOD(CRC32(order_no), 4) = 1
  AND id > :cursor
ORDER BY id
LIMIT 1000

拉取一批后与目标账本对应批次做指纹比对,记录游标继续。每批次完成即上报进度,任务可从中断处恢复。

3.3 幂等与去重

对账任务本身要具备幂等性:

  • 同一批数据可能被重复处理(任务重跑)
  • 差异记录需要以 (task_id, compare_key, diff_type) 建立唯一约束,重复上报直接忽略
  • 修复动作天然要幂等(见第 5 节)
CREATE TABLE IF NOT EXISTS reconciliation_diff (
    task_id      VARCHAR(64) NOT NULL,
    compare_key  VARCHAR(128) NOT NULL,
    diff_type    VARCHAR(16) NOT NULL,   -- MISSING / MISMATCH / EXTRA
    source_snapshot JSON,
    target_snapshot JSON,
    status       VARCHAR(16) DEFAULT 'PENDING',
    created_at   DATETIME DEFAULT CURRENT_TIMESTAMP,
    PRIMARY KEY (task_id, compare_key, diff_type)
);

4. 差异检测

4.1 全量比对与增量比对

维度全量比对增量比对实时比对
数据范围全部记录自上次水位后的变更单条变更事件
延迟小时级分钟级秒级
成本高中低(事件驱动)
用途每日兜底 / 首次初始化常规周期对账关键链路即时校验
典型触发定时任务定时任务MQ / CDC 事件

4.2 指纹比对

全量比对不能逐字段传输大字段,通常使用指纹:

指纹 = hash(主键 + 关键业务字段 + 版本号)

两边的指纹相同 → 记录一致;指纹不同 → 再拉取明细定位差异字段。常用算法:

// 一致性校验指纹(Murmur3 / xxHash / MD5 均可)
public class FingerprintUtil {
    public static String calcFingerprint(String key, String amount, String status) {
        // 将参与比对的字段拼接后哈希,注意字段顺序要一致
        return Hashing.murmur3_128()
            .newHasher()
            .putString(key, UTF_8)
            .putString(amount, UTF_8)
            .putString(status, UTF_8)
            .hash()
            .toString();
    }
}

4.3 差异分类

类型含义修复方向
MISSING源头有、目标没有(缺数据)从源头补写到目标
EXTRA目标有、源头没有(多数据)从目标清理或挂起待人工
MISMATCH双方都有但关键字段不一致按权威源覆盖,或人工裁定
DUPLICATE目标重复记录去重合并

关键原则:明确谁是权威源(Source of Truth)。绝大多数场景下"源头账本"是权威,修复方向为"源头 → 目标";但也存在双方都是视图、需以第三方为准的情况。

4.4 差异检测实现

func (r *Reconciler) DetectBatch(ctx context.Context, source, target []Record, shard Shard) []Diff {
    sourceIndex := indexByKey(source)
    targetIndex := indexByKey(target)
    diffs := make([]Diff, 0)

    // 1) 源头在、目标不在 → MISSING
    for key, s := range sourceIndex {
        t, ok := targetIndex[key]
        if !ok {
            diffs = append(diffs, Diff{Key: key, Type: MISSING, Source: s})
            continue
        }
        // 2) 指纹不一致 → MISMATCH
        if s.Fingerprint != t.Fingerprint {
            diffs = append(diffs, Diff{Key: key, Type: MISMATCH, Source: s, Target: t})
        }
        delete(targetIndex, key)
    }
    // 3) 目标在、源头不在 → EXTRA
    for key, t := range targetIndex {
        diffs = append(diffs, Diff{Key: key, Type: EXTRA, Target: t})
    }
    return diffs
}

5. 差异修复

5.1 修复策略

策略触发条件动作风险
自动修复差异规则明确、权威源明确从权威源重放补写低,但需审计
半自动常规差异自动修,特殊差异挂起自动 + 人工确认中
人工处理金额类、法律类、方向不明的差异生成工单人工裁定高,需全链路留痕

5.2 自动修复的幂等设计

修复动作必须幂等,因为可能被对账系统重试,也可能与业务并发执行。“补写"本质是一次 Upsert,而不是盲目覆盖:

@Transactional
public RepairResult repairDiff(Diff diff) {
    // 1) 从源头读取权威数据
    SourceRow row = sourceRepo.findById(diff.getKey()).orElse(null);
    if (row == null) {
        // 源头也没了:挂起,等待人工
        return RepairResult.suspend(diff, "source row disappeared");
    }
    // 2) 用乐观锁 + 版本号做条件更新,避免覆盖目标侧更新的数据
    int updated = targetRepo.updateIfVersionLessThan(
        diff.getKey(), row.toTargetRecord(), diff.getTargetVersion());
    if (updated == 0) {
        return RepairResult.suspend(diff, "target version moved, manual review");
    }
    // 3) 记录修复日志
    repairLogRepo.save(new RepairLog(diff, row, RepairStatus.SUCCESS));
    return RepairResult.success(diff);
}

经验法则:修复时优先使用"版本号条件更新”,只有目标版本未变时才覆盖,否则挂起人工。这样能避免对账修复与线上业务互相覆盖。

5.3 修复审计与回滚

  • 每次修复写入 repair_log,包含差异 ID、修复前后快照、执行人(系统或人)、时间
  • 提供"一键回滚"能力:如果自动修复判断错误,可以基于日志将目标回滚到修复前快照
  • 对账差异与修复记录本身也要纳入监控,防止"对账系统每天在打架"

6. 对账系统架构

6.1 整体架构

┌────────────────────────── 对账平台 ──────────────────────────┐
│                                                              │
│  任务编排层:对账任务管理 / 调度 / 分片 / 水位 / 重跑          │
│        │                                                    │
│  差异检测层:批量比对引擎 + 实时比对引擎(事件流)             │
│        │                                                    │
│  差异存储层:reconciliation_diff 表 + diff 事件流             │
│        │                                                    │
│  修复执行层:自动修复 / 工单系统 / 审计与回滚                  │
│        │                                                    │
│  可观测层:对账延迟、差异量趋势、修复成功率指标与告警          │
│                                                              │
└──────────────────────────────────────────────────────────────┘
        │ 读取源头账本                     │ 读取目标账本
   ┌────┴─────┐                      ┌────┴─────┐
   │ 订单库     │                      │ 账务库     │
   │ MySQL     │                      │ MySQL/HBase│
   └──────────┘                      └──────────┘

6.2 实时对账与离线对账

  • 离线对账:Hive/Spark 批任务,处理历史全量数据,成本低
  • 准实时对账:定时增量任务(每 5 分钟),覆盖在线变更
  • 实时对账:订阅 MQ/CDC 事件流(如 Debezium),对关键字段做即时比对,命中差异立即告警

实时对账适合"支付回调"、“余额变动"等核心链路;离线对账适合每日大扫除。两者共用同一套差异模型与修复通道。

6.3 与事件溯源结合

事件溯源(Event Sourcing)让对账变得更简单也更有力:既然事件的日志就是事实本身,那么从事件流重放即可重建任意账本。对账因此可以细化为"账本重建比对”:

事件流(orders.events)──► 投影重建账本 A
                       └──► 投影重建账本 B
对比 A、B 即可知道两个投影逻辑是否一致(投影 bug 检测)
同时对比 事件流计数 vs 各服务实际落库计数(事件丢失检测)

具体做法:

  • 每个业务事件携带 event_id + aggregate_id + version
  • 对账时统计事件流中每个 aggregate_id 的最大版本,与实际库中的版本号比对
  • 版本不一致 → 说明投影丢失了某次事件,触发"从事件流重放到最新版本"的修复
// 事件溯源场景下:比对"事件版本"与"库版本"
SELECT aggregate_id, MAX(version) FROM outbox_events GROUP BY aggregate_id;
-- 与
SELECT id, version FROM orders;
-- 两表 JOIN,找出 version 落后 / 缺失的行,从事件流重放补齐

这与 https://plumephp.com/distributed-event-driven-architecture/ 中的事件溯源模式天然互补:事件溯源提供"可重放的事实源",对账提供"验证重放是否正确"的机制。

6.4 对账可观测性

对账系统自己也要被观测:

  • 对账延迟:任务从启动到完成的时间,超时告警
  • 差异量趋势:日差异量突增往往预示业务代码 bug 或消息链路异常
  • 修复成功率:自动修复失败率高说明修复策略或权威源判定有问题
  • 积压水位:待处理差异积压过多说明修复通道堵塞
# Prometheus 指标示例
reconciliation_diff_total{task="order-vs-ledger", type="MISSING"}
reconciliation_diff_total{task="order-vs-ledger", type="MISMATCH"}
reconciliation_repair_success_rate{task="order-vs-ledger"}
reconciliation_watermark_lag_seconds{task="order-vs-ledger"}

7. 生产实践

7.1 对账 SLO 设计

场景对账频率可接受延迟修复时限
支付核心链路实时 + 每日全量秒级分钟级自动修复
一般业务表同步每 5 分钟增量 + 每日全量分钟级小时级
非关键历史数据每日全量小时级当天内
审计类归档每日 / 每周天级按合规要求

7.2 常见坑与经验

  1. 字段时区 / 精度不一致:金额用 decimal,时间统一 UTC,否则"看起来一样"的字段指纹永远不同
  2. 权威源不唯一:对账前必须明确谁是权威,否则两边互相覆盖形成震荡
  3. 修复与业务并发:无条件覆盖会丢业务更新,务必用版本号条件更新
  4. 水位遗漏:updated_at 未建立索引、时区混乱导致增量漏扫,定期全量兜底必不可少
  5. 对账本身成为热点:分片键要均匀,避免单个分片拖慢整个任务
  6. 只修不查:差异修复只是止血,必须回溯根因(是消息丢?补偿漏?代码 bug?),否则差异每天重现

总结

环节关键设计核心原则
对账模型源头/目标账本、比对键、水位明确权威源
任务设计分片扫描、批次游标、幂等去重可中断恢复
差异检测全量/增量/实时、指纹比对指纹先行、明细兜底
差异修复自动/半自动/人工、版本条件更新幂等 + 审计 + 可回滚
系统架构任务编排 + 检测 + 修复 + 可观测事件溯源可重放
生产实践SLO、水位兜底、根因回溯修 + 查并重

数据对账不是银弹,但它把"最终一致性"从口号变成可测量、可验证、可修复的工程能力。它和 https://plumephp.com/distributed-transactions/、https://plumephp.com/distributed-idempotency-reliability/ 一起,构成了分布式系统数据正确的完整闭环。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「distributed-systems」更多文章

  1. 分布式数据库前沿深度解析:TiDB、Spanner 与 CockroachDB 的共识与事务实现
  2. 异地多活与容灾架构深度解析:同城双活、两地三中心与多活设计
  3. 幂等设计与消息可靠性:不丢不重、防止重复消费的分布式基石