迁移校验与对账:分块 Checksum、不变量、Tombstone 与幂等修复
两边都是一千万行,为什么余额还是错了
一场会员库迁移的验收报告很漂亮:源库与目标库总行数完全相同,随机抽查的几百个会员也一致,CDC 延迟已经追平。切流后,财务发现一批账户余额变大,客服发现早已注销的会员重新出现。复盘时才看到两个变化恰好互相抵消:一条删除没有传到目标,另一条新记录没有写入目标,所以总行数仍然相等;抽样又没有抽到这两个主键。
校验失败常常不是没有运行 SQL,而是问题问得太弱。COUNT(*) 只能回答“当前有多少行”,不能回答“是不是同一批行”;抽样只能发现被抽中的错误,不能证明未抽中的数据正确;单个全表 checksum 能发现不同,却很难在变化中的大表上稳定定位,更不能说明差异来自漏数、重复、乱序、类型转换还是删除传播。
可靠对账是一条闭环:先定义同一比较时点和规范化表示,再用分层信号缩小差异;定位到业务键和变更历史后,生成带幂等键的修复计划;修复完成重新校验原分块与关联不变量;最后保存足以复核、又不泄露敏感行内容的证据。任何“发现差异后人工改一下”的流程都没有闭环,因为它无法证明修复没有被重放、覆盖新写或制造第二处错误。
先固定比较时点,再谈相等
源库持续写入时,先查源再查目标会天然读到两个时间点。即使迁移完全正确,结果也可能不同。校验任务必须携带 comparison_watermark:源端在该水位的一致性视图,与目标端“已应用到该水位”的视图比较。无法创建跨库一致性视图时,就短时冻结写,或在业务版本列上限制 version <= watermark。
MySQL 8.4 的一致性读由 InnoDB MVCC 提供,事务隔离级别会改变同一事务后续读是否继续使用旧快照,见 MySQL 8.4 一致性非锁定读。PostgreSQL 当前版允许事务导出快照供其他会话导入,但导出快照有事务生命周期和隔离级别约束,见 pg_export_snapshot。这些能力只能固定单个引擎内的视图;跨引擎比较仍需迁移位点把两个视图对齐。
对账清单至少记录:迁移 ID、对象、源快照/事务 ID、源提交水位、目标已应用水位、schema 指纹、规范化规则版本和分块计划版本。没有这些身份,同一个 checksum 即使连续两次相同,也无法证明比较的是同一历史边界。
从计数到业务不变量的分层证据
第一层是结构证据:表、列、类型、精度、默认值、nullable、主键、唯一约束、字符集、collation 和时区规则是否符合目标设计。schema 不同会让“同一 SQL”产生不同值,例如大小写排序、尾随空格、时间截断和 decimal 精度变化。
第二层是便宜的总体信号:总行数、非空计数、主键最小最大值、金额 SUM、版本最大值、删除标记数。它们适合发现大面积故障和安排深入顺序,但任何一项相等都不是完整性证明。无主键表上的重复行还会让 COUNT(DISTINCT ...) 与业务记录数产生不同含义。
第三层是分块 checksum。先由同一份不可变计划生成共同的 [lower, upper) 边界,再让源和目标都按这些边界读取稳定键;不能让两端根据各自现存主键分页,否则一条缺失记录就会把后续行推入不同块。两侧在相同规范化规则下排序并序列化每行,再对块内容计算密码学哈希,只有哈希不同的块需要继续二分或逐键比较。块太大会让重算慢、锁和缓存压力大;块太小会制造海量元数据与查询。生产通常根据行宽、索引命中、扫描时延和变更热度动态调整,而不是固定“一万行一块”。
第四层是业务不变量。订单头金额应等于有效明细汇总;账户可用余额与冻结余额满足账务方程;已注销主体不能仍有活跃授权;同一业务幂等键只能有一个成功结果。这些约束能发现源与目标“同样错误”,也能发现跨表迁移顺序造成的暂态破坏。
第五层是变更历史证据。某个键不同,要沿 source commit、CDC event、sink apply 和 repair journal 追踪。只有当前值,没有版本、操作类型和事务身份,就无法区分目标漏收最新更新,还是源端后来又发生了合法变化。
这条链上没有“抽几行后直接通过”的捷径。摘要负责缩小搜索面,不变量负责判断业务语义,变更历史负责解释原因,修复账本负责让失败重试可判定;少一层都可能把发现差异与真正恢复混为一谈。
规范化决定 Checksum 是否有意义
checksum 不是把数据库返回的字符串随手拼接。规范化必须无歧义:列顺序固定;每个值带类型和长度;NULL 与空串不同;整数不用本地千分位;decimal 按约定 scale;时间转为约定时区和精度;二进制按原始字节;文本采用明确 Unicode 编码;行按稳定键排序。字段之间只加逗号会发生 ['ab','c'] 与 ['a','bc'] 的拼接歧义。
跨产品迁移时,不要强求数据库内置 checksum 函数输出相同,因为函数、编码、类型转换和聚合顺序可能不同。更稳妥的是两侧各按主键流式读取,在同一个校验器里执行规范化,或让两侧生成统一的长度前缀格式。密码学哈希用于降低碰撞风险,不等于加密;哈希值仍可能成为低基数字段的推断线索,报告访问权限不可放开。
下面是一份可以进入项目仓库的校验计划,字段会直接改变资源消耗和失败语义:
validation_id: member-cutover-final
comparison_watermark: "source-commit:<frozen-high-watermark>"
source_schema_fingerprint: "sha256:<canonical-source-schema>"
target_schema_fingerprint: "sha256:<canonical-target-schema>"
canonicalization_version: v1
chunks:
key: member_id
strategy: range
boundary_semantics: "[lower, upper)"
boundary_plan_uri: "evidence://member-cutover-final/chunks-v1.json"
target_rows: 50000
max_query_seconds: 3
checks:
- row_count
- key_set
- canonical_sha256
- invariant: "account.balance = sum(active_ledger.amount)"
deletes:
require_delete_event_applied: true
retain_tombstone_for_compaction: true
repair:
require_expected_source_version: true
require_observed_target_version: true
require_affected_rows: 1
delete_requires_approved_plan: true
require_idempotency_key: true
max_rows_per_batch: 1000
evidence:
include_raw_values: false
retention: "per-data-governance-policy"max_query_seconds 让控制器在源库压力上升时生成新版本边界计划,而不是让某一侧临时改块后继续比较。修复计划同时保存 expected_source_version 与 observed_target_version:执行时前者必须仍等于源版本,后者必须仍等于目标版本;条件写影响行不是 1 就回滚并重新定位。删除同样需要审批计划和 journal,不能在修复末尾用一条不受控的 DELETE ... NOT IN (...) 清场。include_raw_values: false 要求默认只保存键的不可逆标识、差异类型、版本和摘要;确需原值的个案应走受控取证。
Tombstone 是删除事实,不是空值噪声
插入和更新把“存在的值”送到目标,删除必须传播“这个键不再存在”的事实。只复制当前快照而不传删除,会让目标永久保留旧记录;把软删除字段漏出投影,也会让目标把已注销对象当作活跃对象。
在 Debezium 的 Kafka 事件语义中,源端 delete 先产生 op: d 的删除事件:旧行在 before,源位点在 source,业务版本若由表列承载也从 before 读取。随后同键、值为 null 的 tombstone 只为 compacted topic 清理该键历史;它没有事件 value,因而不承载业务版本、op 或源位点,见 Debezium tombstone 官方说明。过滤 tombstone 可能不影响数据库 sink 已执行的 delete,却会改变重放和压缩后的最终状态,因此不能把它当作“节省消息”的无害优化。
校验删除至少有三种办法:比较完整键集合;维护源端删除台账并确认每个带版本的 delete event 已应用,同时确认 compacted topic 需要的 tombstone 未被误滤;在目标保留受控的删除审计表,记录源事务身份、键哈希和应用状态。仅比较源端现存行无法解释目标多出的键究竟是漏删、越界导入还是源端快照之后才删除。
乱序会让删除问题更隐蔽。若带版本 5 的 delete event 已应用,随后到达版本 4 的 update,按到达顺序盲写会复活记录;紧随 delete event 的 null tombstone 本身不能参与版本比较。sink 必须把最后接受的业务版本或源事务顺序保存在目标行或独立删除版本表中,据此拒绝旧事件;若跨分区无法提供同键顺序,就必须保证同键分区稳定,或在目标端做条件更新。
运行一个会暴露五类错误的正反模型
实验需要 Python 3.12 或更高版本,不安装第三方包。它使用两个本地 SQLite 文件模拟源与目标,以规范 JSON 序列化、SHA-256 分块哈希、业务不变量和修复账本形成闭环。Python 的事务与行工厂接口见 sqlite3 官方文档,哈希接口见 hashlib 官方文档。
把以下内容保存为临时目录中的 reconcile_lab.py:
import argparse
import hashlib
import json
import sqlite3
from decimal import Decimal
from pathlib import Path
SCHEMA = """
CREATE TABLE members(
id INTEGER PRIMARY KEY, name TEXT NOT NULL, balance TEXT NOT NULL,
version INTEGER NOT NULL, deleted INTEGER NOT NULL DEFAULT 0
);
CREATE TABLE ledger(
id INTEGER PRIMARY KEY, member_id INTEGER NOT NULL,
amount TEXT NOT NULL, active INTEGER NOT NULL DEFAULT 1
);
CREATE TABLE repair_journal(
repair_key TEXT PRIMARY KEY, entity TEXT NOT NULL, entity_id INTEGER NOT NULL,
expected_source_version INTEGER, observed_target_version INTEGER,
action TEXT NOT NULL, approved_by TEXT NOT NULL, status TEXT NOT NULL
);
"""
def connect(path):
db = sqlite3.connect(path)
db.executescript(SCHEMA)
return db
def seed_source(db):
db.executemany("INSERT INTO members VALUES(?,?,?,?,?)", [
(1, "Alice", "30.00", 3, 0),
(2, "Bob", "5.00", 2, 0),
(3, "Carol", "0.00", 4, 1),
(4, "Dana", "8.00", 1, 0),
])
db.executemany("INSERT INTO ledger VALUES(?,?,?,?)", [
(101, 1, "10.00", 1), (102, 1, "20.00", 1),
(201, 2, "5.00", 1), (401, 4, "8.00", 1),
])
db.commit()
def copy_all(source, target):
for table in ("members", "ledger"):
cols = len(source.execute(f"SELECT * FROM {table} LIMIT 1").description)
target.executemany(
f"INSERT INTO {table} VALUES({','.join('?' * cols)})",
source.execute(f"SELECT * FROM {table}"),
)
target.commit()
def corrupt(target):
# 漏数:删掉 Dana;越界导入一行,使总行数仍为 4;漏删让 Carol 复活。
target.execute("DELETE FROM members WHERE id=4")
target.execute("INSERT INTO members VALUES(5,'Eve','0.00',1,0)")
target.execute("UPDATE members SET deleted=0 WHERE id=3")
# 乱序:旧版本覆盖 Alice;重复:多一条账本;大小写值发生变化。
target.execute("UPDATE members SET balance='10.00',version=2 WHERE id=1")
target.execute("INSERT INTO ledger VALUES(999,1,'10.00',1)")
target.execute("UPDATE members SET balance='5',name='bob' WHERE id=2")
target.commit()
def normalized(row):
# 类型和列顺序固定;Decimal 统一两位,NULL 不与空串混淆。
return {
"id": int(row[0]), "name": row[1],
"balance": format(Decimal(row[2]), ".2f"),
"version": int(row[3]), "deleted": int(row[4]),
}
def canonical(row):
data = normalized(row)
return json.dumps(data, ensure_ascii=False, sort_keys=True,
separators=(",", ":")).encode("utf-8")
BOUNDARIES = ((1, 3), (3, 5), (5, 7))
def chunks(db, boundaries=BOUNDARIES):
result = {}
for lower, upper in boundaries:
rows = db.execute(
f"SELECT id,name,balance,version,deleted FROM members "
"WHERE id>=? AND id<? ORDER BY id", (lower, upper)
).fetchall()
digest = hashlib.sha256()
for row in rows:
payload = canonical(row)
digest.update(len(payload).to_bytes(4, "big"))
digest.update(payload)
result[f"[{lower},{upper})"] = {
"row_count": len(rows), "sha256": digest.hexdigest()
}
return result
def invariants(db):
failures = []
for mid, balance, deleted in db.execute(
"SELECT id,balance,deleted FROM members ORDER BY id"
):
ledger = sum(Decimal(x[0]) for x in db.execute(
"SELECT amount FROM ledger WHERE member_id=? AND active=1", (mid,)
))
if not deleted and Decimal(balance) != ledger:
failures.append({"member": mid, "balance": balance,
"active_ledger": str(ledger)})
return failures
def diff(source, target):
s = {r[0]: r for r in source.execute("SELECT * FROM members")}
t = {r[0]: r for r in target.execute("SELECT * FROM members")}
result = {
"missing_on_target": sorted(s.keys() - t.keys()),
"extra_on_target": sorted(t.keys() - s.keys()),
"version_regression": [], "delete_mismatch": [], "value_mismatch": [],
}
for key in sorted(s.keys() & t.keys()):
source_row, target_row = normalized(s[key]), normalized(t[key])
if source_row["deleted"] != target_row["deleted"]:
result["delete_mismatch"].append(key)
elif target_row["version"] < source_row["version"]:
result["version_regression"].append(key)
elif source_row != target_row:
result["value_mismatch"].append(key)
return result
def repair(source, target, plan):
for item in plan:
key = item["repair_key"]
done = target.execute(
"SELECT status FROM repair_journal WHERE repair_key=?", (key,)
).fetchone()
if done and done[0] == "DONE":
continue
entity, row_id = item["entity"], item["id"]
try:
target.execute("BEGIN IMMEDIATE")
target.execute(
"INSERT INTO repair_journal VALUES(?,?,?,?,?,?,?,?)",
(key, entity, row_id, item["expected_source_version"],
item["observed_target_version"], item["action"],
item["approved_by"], "STARTED"),
)
if item["action"] == "UPSERT_MEMBER":
src = source.execute(
"SELECT * FROM members WHERE id=? AND version=?",
(row_id, item["expected_source_version"]),
).fetchone()
if src is None:
raise RuntimeError(f"source version changed for member {row_id}")
if item["observed_target_version"] is None:
cursor = target.execute("INSERT INTO members VALUES(?,?,?,?,?)", src)
else:
cursor = target.execute(
"UPDATE members SET name=?,balance=?,version=?,deleted=? "
"WHERE id=? AND version=?",
(src[1], src[2], src[3], src[4], row_id,
item["observed_target_version"]),
)
elif entity == "members":
cursor = target.execute(
"DELETE FROM members WHERE id=? AND version=?",
(row_id, item["observed_target_version"]),
)
else:
observed = item["observed_row"]
cursor = target.execute(
"DELETE FROM ledger WHERE id=? AND member_id=? AND amount=? AND active=?",
observed,
)
if cursor.rowcount != 1:
raise RuntimeError(f"conditional repair affected {cursor.rowcount} rows: {key}")
target.execute(
"UPDATE repair_journal SET status='DONE' WHERE repair_key=?", (key,)
)
target.commit()
except Exception:
target.rollback()
raise
def run(root, mode):
root.mkdir(parents=True, exist_ok=True)
for name in ("source.db", "target.db"):
(root / name).unlink(missing_ok=True)
source, target = connect(root / "source.db"), connect(root / "target.db")
seed_source(source)
copy_all(source, target)
if mode == "bad":
corrupt(target)
print(json.dumps({
"counts": [source.execute("SELECT count(*) FROM members").fetchone()[0],
target.execute("SELECT count(*) FROM members").fetchone()[0]],
"chunk_mismatch": chunks(source) != chunks(target),
"diff": diff(source, target),
"invariant_failures": invariants(target),
}, ensure_ascii=False, sort_keys=True))
return
before = chunks(target)
assert before == chunks(source) and not invariants(target)
corrupt(target)
d = diff(source, target)
plan = []
repair_members = set(d["missing_on_target"])
repair_members.update(d["version_regression"])
repair_members.update(d["delete_mismatch"])
repair_members.update(d["value_mismatch"])
for mid in sorted(repair_members):
expected = source.execute(
"SELECT version FROM members WHERE id=?", (mid,)
).fetchone()[0]
observed = target.execute(
"SELECT version FROM members WHERE id=?", (mid,)
).fetchone()
plan.append({
"repair_key": f"member:{mid}:source-v{expected}", "entity": "members",
"id": mid, "expected_source_version": expected,
"observed_target_version": None if observed is None else observed[0],
"action": "UPSERT_MEMBER", "approved_by": "migration-reviewer",
})
for mid in d["extra_on_target"]:
observed = target.execute(
"SELECT version FROM members WHERE id=?", (mid,)
).fetchone()[0]
plan.append({
"repair_key": f"member:{mid}:approved-delete-v{observed}",
"entity": "members", "id": mid, "expected_source_version": None,
"observed_target_version": observed, "action": "DELETE",
"approved_by": "migration-reviewer",
})
source_ledger_ids = {r[0] for r in source.execute("SELECT id FROM ledger")}
for row in target.execute("SELECT * FROM ledger ORDER BY id").fetchall():
if row[0] not in source_ledger_ids:
plan.append({
"repair_key": f"ledger:{row[0]}:approved-delete", "entity": "ledger",
"id": row[0], "expected_source_version": None,
"observed_target_version": None, "observed_row": row,
"action": "DELETE", "approved_by": "migration-reviewer",
})
repair(source, target, plan)
repair(source, target, plan) # 同一计划重放,不会重复执行。
assert all(not values for values in diff(source, target).values())
assert chunks(source) == chunks(target) and not invariants(target)
print(json.dumps({"status":"RECONCILED", "diff":diff(source, target),
"journal_rows":target.execute(
"SELECT count(*) FROM repair_journal").fetchone()[0],
"recheck":True}, ensure_ascii=False, sort_keys=True))
if __name__ == "__main__":
p = argparse.ArgumentParser()
p.add_argument("mode", choices=["good", "bad"])
p.add_argument("--workdir", default="reconcile-lab-data")
args = p.parse_args()
run(Path(args.workdir), args.mode)运行正反两条路径:
python reconcile_lab.py bad --workdir reconcile-bad
python reconcile_lab.py good --workdir reconcile-good反向输出中的 counts 是 [4, 4],说明总行数检查没有报错;但共同 [lower, upper) 边界上的 chunk_mismatch 为 true,missing_on_target 包含会员 4,extra_on_target 包含会员 5,version_regression、value_mismatch、delete_mismatch 分别包含会员 1、2、3,业务不变量还会暴露 Alice 的余额与有效流水不相等。这里分别对应漏数、越界数据、旧更新乱序覆盖、值变化、delete event 漏应用和重复流水。逐键分类与块摘要调用同一 normalized(),因此 5 与 5.00 这类按规则等价的 decimal 表示不会仅因存储文本不同而误报;本例会员 2 仍因名称大小写变化进入 value_mismatch。两端使用相同 schema,因而这些差异不能归因为精度漂移;要识别精度漂移,还必须比较类型、精度、标度、字符集、排序规则和默认值组成的 schema fingerprint。
正向路径先检查初始副本,再注入相同错误,为 upsert 和 delete 分别生成审批计划,并连续执行两次。每个计划同时固定源版本与观察到的目标版本,条件 DML 影响行必须为 1;最终输出 RECONCILED、空差异、recheck: true,journal_rows 为 6 且第二次执行不再增长。SQLite 的 BEGIN IMMEDIATE 只模拟本地事务互斥;生产数据库仍需按自身隔离级别处理源版本读取与目标条件写之间的并发窗口。
抽样为什么只能排风险,不能签完整性
抽样适合做快速健康信号、人工语义检查和超大对象的优先级排序。它的发现概率取决于错误比例与抽样设计:错误集中在冷分区、特定租户、长尾字符或历史月份时,均匀随机样本很容易错过;按主键前缀取样更可能系统性遗漏高位键。
要提高抽样价值,可以按租户规模、时间分区、数据类型边界和风险标签分层,故意覆盖 NULL、最大精度、非 ASCII 文本、超长 LOB、软删除和复合主键。但报告必须把结论写成“这些样本未发现差异”,不能写成“数据一致”。真正的切流门槛仍需完整键集合、全分块摘要或能覆盖全部数据的等价证明。
LOB 和加密列可能让全量逐字节比较代价过高。可以先比较长度、内容哈希、对象版本和元数据,再对差异块读取原文;客户端随机加密会让同一明文产生不同密文,此时需要在受控可信边界内比较解密后的规范值或业务生成的稳定指纹,不能简单认定密文不同就是迁移错误。
从不同的块定位到一条错误历史
差异定位采用逐层收敛,而不是立即导出整表。先比较 schema 指纹和总体信号;再找 checksum 不同的一级块;对不同块二分或按更小范围重算;最后比较键集合与逐行规范摘要。输出应把差异分成 missing_on_target、extra_on_target、value_mismatch、version_regression、delete_mismatch 和 schema_semantic_mismatch。
定位到键后,沿四份证据拼接时间线:源端该键在比较水位前的最后提交;CDC 中的操作类型、事务身份和顺序;sink 的条件写结果与重试;目标当前版本。若源有版本 7、CDC 有版本 7、sink 显示约束拒绝,问题在目标映射或约束;若 CDC 只有版本 6,问题在捕获过滤、日志缺口或事务解码;若 sink 已成功写版本 7、当前却是版本 6,说明有旧事件、人工脚本或第二 writer 覆盖。
数据库快照 checksum 与在线 CDC 校验也要避免相互追逐。热键不断变化时,同一块反复失败可能只是比较视图不一致。先核对源 comparison watermark 与目标 applied watermark;必要时只对差异键做短时冻结或使用版本上界,而不是无限重跑整表。
幂等修复不是一条 UPSERT
修复计划是一份不可变输入,至少包含 repair_key、实体、键、差异类型、expected_source_version、observed_target_version、动作、计划哈希和审批身份。执行器在同一事务中插入修复账本、重读并核对源版本、按观察到的目标版本执行条件写、检查影响行恰好为 1,最后标记完成;同一 repair_key 重跑只返回既有结果。源或目标任一版本已经变化,或条件 DML 影响 0 行/多行,执行器都必须回滚并重新定位,不能让旧修复覆盖在线新写。
修复动作要按语义区分。漏行可以按版本条件 upsert;多余行需要确认是漏 delete event 还是越界导入,再生成独立审批计划并写入同一 repair journal 后条件删除;乱序应恢复最高合法版本并修复 sink 的版本比较;重复账务记录不能只删任意一条,必须依据业务幂等键和审计链选择;schema 漂移先比较两侧真实 schema fingerprint,确认列类型、precision/scale 或转换规则确有变化,再修映射并重放受影响范围。
大批修复应限速、分批提交并设置失败隔离。单个超大事务会制造锁、日志和复制延迟;逐行自动提交又会放大网络和日志开销。批大小应由行宽、索引、锁等待、日志增长和在线延迟共同调节。每个批次保存起止键、成功数、跳过数、冲突数和目标已提交事务身份,进程中断后从账本恢复。
修复后不能只重跑原来的逐行比较。还要重新计算受影响块 checksum、键集合、关联业务不变量与删除台账,并观察 CDC backlog 是否恢复。否则“把会员行修对”可能仍留下重复流水,“补一条订单”也可能破坏库存或唯一约束。
项目怎样接入持续校验
校验器不应直接耦合切流脚本。项目可以定义一个只读的验证 API 或批任务契约:输入迁移 ID、对象、comparison watermark 和计划版本;输出状态、覆盖范围、不同块、差异分类、资源消耗与证据 URI。切流控制器只接受签名后的 PASSED 结果,并检查水位和 schema 指纹仍与当前状态一致,过期结果自动失效。
在全量阶段,每完成一个分片就记录源行数、目标行数和摘要,尽早发现映射错误;在追平阶段,持续检查版本回退、tombstone backlog 和业务不变量趋势;冻结高水位后执行最终完整校验;切流后继续运行低影响巡检,捕获隐藏 writer 和延迟任务。这样校验从一次验收 SQL 变成迁移期间的持续反馈。
任务部署可以是靠近数据库的临时 worker,也可以是受控批处理平台。靠近数据能减少跨网传输,但会竞争数据库 CPU、IO 和连接;集中平台便于调度与审计,却可能把敏感数据带出原安全域。无论形态,都要限制并发连接、查询超时、块扫描速率和临时磁盘,避免“为了证明数据库正确而把数据库压垮”。
常见失败怎样留下可判断证据
校验任务持续超时,先查看是否走主键/分区索引、块是否跨热点、MVCC 长快照是否制造版本保留、目标是否同时建索引。不要无限提高超时;缩小块、降低并发、转向只读副本前先确认副本水位满足 comparison watermark。
所有块突然不同,优先检查规范化版本、列顺序、时区、collation、decimal scale 和 schema 指纹,而不是认定全量丢失。一个校验器版本变更足以让全部摘要变化,所以规范化代码和规则必须版本化,升级时用固定测试向量验证前后预期。
只有删除相关块反复多行,检查 source connector 是否产生 delete、转换链是否过滤 tombstone、sink 是否把 null 当作忽略、软删除列是否在投影中,以及重放压缩 topic 后键是否仍存在。证据要同时取源提交、原始事件、转换后事件和 sink 结果,不能只看最终 topic。
修复任务显示成功但重校验仍失败,检查事务是否真正提交、账本与业务写是否同事务、后续旧事件是否覆盖、计划是否引用错误水位。若账本先标 DONE 再提交业务写,崩溃窗口会产生“记录成功、实际未修”;正确顺序是同事务提交,或用可恢复中间态让重启继续。
权限、敏感报告与证据保留
校验账号通常只需要指定 schema/table 的只读权限,不应拥有 DDL、账号管理或全库访问。修复执行器使用单独身份,只允许受控存储过程或目标表上的条件 DML;修复审批者不直接持有数据库密码。生产连接采用 TLS、短期凭证和密钥注入,命令示例始终使用 localhost 与占位凭证。
差异报告可能比原表更危险,因为它把异常用户、删除记录、余额和身份字段集中在一起。默认保存迁移 ID、表、块、键哈希、差异类型、版本、摘要和计数,不保存完整行。需要原值取证时,单独加密、严格授权、审计下载并设置销毁任务;日志不得打印连接串、证书、明文主键或 LOB。
证据保留要兼顾复核与最小化。长期保留校验计划、代码版本、comparison watermark、schema 指纹、块边界、摘要、差异分类、修复计划哈希、账本状态和重校验结论;原始行样本与临时导出在复核结束后销毁。哈希算法和规范化版本必须一起保存,否则未来无法解释旧摘要。
容量、成本与团队门禁
全表 checksum 本质是全量读,可能驱逐业务缓存、拉高存储吞吐并增加副本延迟。容量模型至少估算总字节、平均行宽、索引扫描比例、并发块数、每秒读取上限、临时摘要存储和重算比例。热表优先按分区和水位增量校验,冷表安排在负载窗口;越过源端延迟或磁盘阈值时自动降并发。
跨区域校验若把原始行拉到中央,会产生网络出口与数据驻留风险。让两侧在本地计算规范摘要、只交换块摘要能降低传输,但差异块仍需受控下钻。成本评审记录 worker 运行时长、扫描字节、跨网字节、临时存储和重复重算,不记录产品价格,也不让资源在迁移结束后继续常驻。
团队门禁应明确责任:数据库 owner 维护快照与负载边界;迁移 owner 维护水位、分块和工具版本;业务 owner 定义不变量和差异裁决;安全 owner 审批敏感取证;切流指挥只接受未过期的验证结果。任何人都不能一边修改规范化规则,一边独自批准校验通过。
清理实验与结束一次对账
本地实验只生成两个目录,可直接清理:
rm -rf reconcile-good reconcile-bad
# PowerShell: Remove-Item -Recurse -Force reconcile-good, reconcile-bad生产清理先停止新任务,再等待运行中块提交或标记可恢复中断;撤销校验与修复凭证;删除临时明文导出;按策略保留摘要和账本;最后回收 worker、临时磁盘与网络规则。不要先删修复账本,否则重启时无法判断哪些补偿已经提交;也不要无限保留差异原文,把一次迁移变成永久敏感数据副本。
一次对账结束的判据不是“差异 SQL 返回零行”这一瞬间,而是:比较水位仍有效,schema 与规范化规则未变化,所有分块和键集合覆盖完整,业务不变量通过,tombstone 已收敛,修复账本无悬挂状态,受影响范围重校验通过,证据可复核且敏感原文已按策略清理。到这里,团队才拥有可以支撑切流或退出的工程证据。
