全量快照接 CDC:从水位追平到灰度切流与可回切退出
凌晨切流后,旧库为什么还在接单
一次订单库迁移已经进行了数小时:目标库的表都在,监控显示 CDC 延迟只有几秒,团队便把一半读请求切过去。很快,客服看到同一个订单在两个页面呈现不同状态;更糟的是,旧版本应用仍把取消请求写进旧库,新版本应用把付款结果写进新库。两边都在“成功提交”,却没有任何一边能证明自己是最终事实。
事故并不是复制工具少传了一行,而是控制面缺失。全量任务完成时没有保存与快照一致的日志起点,CDC 的“当前位点”也没有对应到源端提交水位;切写时只改了流量比例,没有指定唯一权威侧;发现异常后想回切,却从未把新库写入反向送回旧库。于是“回滚开关”只是把客户端送回一个已经落后的库。
在线迁移真正搬运的不是一份静态表,而是一条随事务持续前进的历史。团队必须始终回答四个问题:快照代表源库哪个时刻,快照之后的每个提交由谁接续,当前允许谁受理写入,以及回切目标是否仍拥有完整历史。只要其中一个答案含糊,低延迟和高行数都可能是假安全感。
把快照、水位和权威侧放进同一模型
先建立几个不会因产品名变化而失效的对象。
全量快照是在一个一致性视图里读取的基线。它不是“脚本开始时间到结束时间之间的数据平均值”。MySQL 8.4 的 InnoDB 一致性非锁定读依赖 MVCC;START TRANSACTION WITH CONSISTENT SNAPSHOT 让一致性视图在事务开始时建立,但其适用条件仍受隔离级别和存储引擎影响,细节应以 MySQL 8.4 一致性非锁定读和事务语句为准。非事务表、外部对象存储和异步派生表不能自动共享这个视图。
低水位是快照必须接续的日志位置,常见表达是 MySQL binlog 文件/位置或 GTID 集合、PostgreSQL LSN。它必须与快照建立动作有原子或可证明的对应关系。先查位点、过几秒再导表,会留下“位点之后、快照之前”的重复窗口;先导表、结束后再查位点,则可能留下静默漏数窗口。MySQL 官方把传统文件位置与 GTID 复制区分开来,GTID 让事务身份跨日志文件更稳定,见 MySQL 8.4 复制方法。
高水位是准备切换时源端已经确认提交的边界。CDC 消费到该边界,只说明传输器读到了这里;还要证明目标端对应事务已经提交、失败重试已经收敛,才能称为“已应用水位”。因此控制表至少要分开保存 source_committed_watermark、transport_seen_watermark 和 target_applied_watermark。
权威侧是当冲突发生时唯一获胜的数据源。它不是“主要流量所在的一侧”,而是业务写入契约。镜像、影子、回放和反向同步都不能自行晋升为权威。双写若没有事务身份、版本比较与对账协议,只是把一个失败点变成两个成功但互相矛盾的失败点。
一条可审计的切换状态机
把迁移写成状态机,目的不是增加流程表单,而是让自动化拒绝非法跃迁。每个状态都要同时保存进入证据、失败原因、操作者和可恢复动作。
Prepared 不能靠一句“检查完成”进入。源库日志保留必须覆盖全量耗时、最大追平耗时、人工决策时间和安全余量;目标库空间要容纳基线、索引重建、临时写放大和回切日志;迁移账号只能读取指定对象和复制位点,切流账号与校验账号分离。MySQL 8.4 默认启用 binary logging,默认 row format,但生产仍应显式检查 log_bin、binlog_format、GTID 与保留策略,而不是依赖默认值;官方边界见 binary log format。
Snapshotting 的持久状态不是百分比,而是 snapshot_id + low_watermark + object_manifest + schema_fingerprint。进程重启后若只剩“复制到第 73 张表”,却不知道这 73 张表属于哪个一致性视图,就不能续跑,只能废弃该基线或按工具的增量快照协议恢复。
CatchingUp 要观察“距离”而不只观察“时间”。时钟漂移会让秒级延迟失真,大事务也可能让日志长时间无已应用位点推进。优先用待应用事务数、GTID/LSN 差距、最老未完成事务年龄、sink 重试和死信对象共同判断。
TargetAuthoritative 是风险最大的跃迁。路由配置已发布不代表旧实例已经停止写入;要等连接池耗尽、后台任务暂停、定时器和人工脚本退出,再用源端审计证明冻结高水位之后没有未知 writer。
用配置把口头约定变成可执行门禁
下面的 YAML 不是某个迁移产品的专属语法,而是可以进入仓库、由控制器读取的切换清单。字段名可以改,语义不能丢。
migration_id: orders-v2
authority: source
snapshot:
snapshot_id: snap-orders-001
low_watermark: "gtid:<captured-with-consistent-view>"
schema_fingerprint: "sha256:<canonical-ddl>"
cdc:
applied_watermark_store: migration_control
max_unapplied_transactions: 0
max_oldest_unapplied_seconds: 5
retry_backlog_limit: 0
cutover:
mode: freeze
freeze_timeout_seconds: 60
read_canary_percent: 5
write_canary_percent: 0
required_stable_windows: 3
rollback:
reverse_sync_required: true
window_close_condition: "source_rebuilt_and_restore_drill_passed"
evidence_retention: "per-data-governance-policy"
objectives:
cutover_rpo: "zero acknowledged writes"
cutover_rto_seconds: 60
rollback_rpo: "zero acknowledged writes while rollback-ready"
rollback_rto_seconds: 300
evidence_clock: monotonicmax_unapplied_transactions: 0 不是要求日常永远零延迟,而是要求切换门槛时不存在已知未提交事务。freeze_timeout_seconds 越长,追平成功概率越高,但用户可用性代价越大;超过预算必须自动恢复旧库写入,不允许值班人员边讨论边无限冻结。write_canary_percent 在冻结方案中保持零,因为同一业务键同时向两侧随机分流会制造分裂权威。写灰度应按租户或稳定键隔离,并且每个键只有一个权威侧。
required_stable_windows 防止指标瞬时过线。每个窗口都应重新读取源高水位、目标已应用水位、差异计数和错误队列;连续通过后才允许跃迁。window_close_condition 不使用一个随意时长,因为业务周期、延迟任务和退款链路可能远长于技术观察窗口。关闭回切能力必须由恢复演练和业务证据共同决定。
objectives 中的数字只是演示预算,生产值必须来自业务 SLO、峰值写入速率和演练基线。切换 RPO 为零,表示冻结高水位以前所有已确认写入都出现在新权威侧,冻结期间被拒绝或排队的请求有可判定结果;它不等于“CDC 延迟为零”。切换 RTO 从冻结或故障决策开始,直到目标侧读写达到约定服务等级。回切 RPO 为零还要求切换后的每笔已确认新写已经反向应用,或可靠留存在可重放日志;回切 RTO 则包含阻断新写、旧库追平、重校验、路由生效和实例排空。证据应保存单调时钟耗时、起止状态、两端水位和合成事务结果,不能只写会议纪要中的开始结束时间。
从零运行切换正反模型
实验只需要 Python 3.12 或更高版本,使用标准库 sqlite3,不会连接任何外部数据库。先用 python --version 确认解释器,用 Python sqlite3 官方文档理解事务接口。把下面保存为临时目录中的 cutover_lab.py;它创建合成订单库、变更日志和控制状态,不含真实账号或数据。
import argparse
import json
import sqlite3
from pathlib import Path
SCHEMA = """
CREATE TABLE IF NOT EXISTS orders(
id INTEGER PRIMARY KEY, status TEXT NOT NULL, version INTEGER NOT NULL
);
CREATE TABLE IF NOT EXISTS changes(
seq INTEGER PRIMARY KEY AUTOINCREMENT,
order_id INTEGER NOT NULL, status TEXT NOT NULL, version INTEGER NOT NULL
);
"""
def open_db(path):
db = sqlite3.connect(path)
db.executescript(SCHEMA)
return db
def write(db, order_id, status):
old = db.execute("SELECT version FROM orders WHERE id=?", (order_id,)).fetchone()
version = 1 if old is None else old[0] + 1
db.execute(
"INSERT INTO orders VALUES(?,?,?) ON CONFLICT(id) DO UPDATE "
"SET status=excluded.status, version=excluded.version",
(order_id, status, version),
)
db.execute("INSERT INTO changes(order_id,status,version) VALUES(?,?,?)",
(order_id, status, version))
db.commit()
return version
def install_snapshot(target, snapshot_rows):
target.execute("BEGIN IMMEDIATE")
target.execute("DELETE FROM orders")
target.executemany("INSERT INTO orders VALUES(?,?,?)", snapshot_rows)
target.commit()
def consistent_snapshot(source, target):
# 位点和行集来自同一个源端读事务,模拟工具提供的原子快照握手。
source.execute("BEGIN")
low = source.execute("SELECT COALESCE(MAX(seq),0) FROM changes").fetchone()[0]
snapshot_rows = source.execute(
"SELECT id,status,version FROM orders ORDER BY id"
).fetchall()
source.commit()
install_snapshot(target, snapshot_rows)
return low
def broken_snapshot(source, target):
# 反例:先导数据,随后才取 CDC 起点;两步之间的提交会被永久跳过。
snapshot_rows = source.execute(
"SELECT id,status,version FROM orders ORDER BY id"
).fetchall()
install_snapshot(target, snapshot_rows)
write(source, 3, "CREATED")
return source.execute("SELECT MAX(seq) FROM changes").fetchone()[0]
def apply_until(source, target, after, high=None):
sql = "SELECT seq,order_id,status,version FROM changes WHERE seq>?"
args = [after]
if high is not None:
sql += " AND seq<=?"
args.append(high)
applied = after
for seq, oid, status, version in source.execute(sql + " ORDER BY seq", args):
target.execute(
"INSERT INTO orders VALUES(?,?,?) ON CONFLICT(id) DO UPDATE SET "
"status=excluded.status,version=excluded.version "
"WHERE excluded.version > orders.version",
(oid, status, version),
)
applied = seq
target.commit()
return applied
def rows(db):
return db.execute("SELECT id,status,version FROM orders ORDER BY id").fetchall()
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 = open_db(root / "source.db"), open_db(root / "target.db")
write(source, 1, "CREATED")
write(source, 2, "CREATED")
if mode == "bad-origin":
low = broken_snapshot(source, target)
applied = apply_until(source, target, low)
missing = sorted({r[0] for r in rows(source)} - {r[0] for r in rows(target)})
print(json.dumps({"state":"BLOCKED_NON_ATOMIC_ORIGIN", "low":low,
"applied":applied, "missing_on_target":missing},
ensure_ascii=False))
return
low = consistent_snapshot(source, target)
print(json.dumps({"state":"SNAPSHOTTED", "low":low}, ensure_ascii=False))
write(source, 1, "PAID")
write(source, 3, "CREATED")
if mode == "good":
applied = apply_until(source, target, low)
print(json.dumps({"state":"CATCHING_UP", "applied":applied}, ensure_ascii=False))
high = source.execute("SELECT MAX(seq) FROM changes").fetchone()[0]
applied = apply_until(source, target, applied, high)
assert applied == high and rows(source) == rows(target)
authority = "target"
write(target, 2, "CANCELLED")
# 这里只做一次反向应用;升级为可回切状态还需要持久 offset 等完整证据。
reverse_from = 0
reverse_from = apply_until(target, source, reverse_from)
assert rows(source) == rows(target)
print(json.dumps({"state":"FACT_EQUAL_AT_CHECKPOINT", "authority":authority,
"forward_high":high, "reverse_applied":reverse_from,
"equal":True}, ensure_ascii=False))
else:
# 错误模型:不冻结、无稳定分片规则,两侧同时成为 writer。
write(source, 1, "CANCELLED")
write(target, 1, "SHIPPED")
source_rows = {r[0]: r for r in rows(source)}
target_rows = {r[0]: r for r in rows(target)}
conflicts = [
{"id": oid, "source": source_rows.get(oid), "target": target_rows.get(oid)}
for oid in sorted(source_rows.keys() | target_rows.keys())
if source_rows.get(oid) != target_rows.get(oid)
]
print(json.dumps({"state":"BLOCKED_SPLIT_AUTHORITY",
"source":rows(source), "target":rows(target),
"conflicts":conflicts}, ensure_ascii=False))
if __name__ == "__main__":
p = argparse.ArgumentParser()
p.add_argument("mode", choices=["good", "bad-origin", "bad-split"])
p.add_argument("--workdir", default="cutover-lab-data")
args = p.parse_args()
run(Path(args.workdir), args.mode)在同一临时目录运行:
python cutover_lab.py good --workdir cutover-good
python cutover_lab.py bad-origin --workdir cutover-bad-origin
python cutover_lab.py bad-split --workdir cutover-bad-split正向命令依次输出 SNAPSHOTTED、CATCHING_UP 和 FACT_EQUAL_AT_CHECKPOINT,最后的 equal 为 true。这个状态只表示该检查点两库事实相等。升级为 ROLLBACK_READY 还必须持久化反向 offset,并证明 origin 防环、删除传播和重启续传均成立。
第一个反向命令输出 BLOCKED_NON_ATOMIC_ORIGIN,missing_on_target 为 [3]。订单 3 在导出行集之后、读取 CDC 起点之前提交;快照没有它,而 CDC 又从包含该提交的水位之后开始,所以即使 applied == low,缺口也不会自行消失。第二个反向命令输出 BLOCKED_SPLIT_AUTHORITY:订单 1 在源端已沿 CREATED -> PAID -> CANCELLED 走到版本 3,目标端却独立从 CREATED -> SHIPPED 走到版本 2;订单 3 还只存在于源端。版本较大只能说明某一侧局部更新次数更多,不能自动裁决两个独立 writer 的业务意图,必须有唯一发号源、全序事务身份或明确的冲突规则。
SQLite 中的 changes.seq 只提供单进程顺序,不具备 MySQL MVCC、GTID、网络重试或连接器重启语义。生产链路必须改用连接器确认过的 GTID/LSN,把低水位、已应用水位和 authority 持久化到迁移控制表,并通过阻断式 API 拒绝非法跃迁。
全量与增量怎样无缝接上
最稳妥的接续由迁移工具在同一协议内完成:建立一致性快照时捕获低水位,快照完成后从该水位读取日志,并把 offset 与 schema history 一起持久化。Debezium 的快照、offset 和 schema history 共同组成恢复链;其稳定版文档列出 initial 等快照模式及恢复条件,见 Debezium 快照机制。删除 offset 后直接复用旧目标数据,会让连接器重新发出历史事件;删除 schema history 则可能使旧日志无法按当时表结构解释。
自研导出与 CDC 拼接时,需要一次可审计握手。MySQL 8.4 已用 SHOW BINARY LOG STATUS 取代不再支持的 SHOW MASTER STATUS,但单独保存这条命令的输出仍不是一致性证明;应采用工具支持的锁、一致性事务、GTID 或备份元数据协议,把日志坐标与快照视图建立可证明的对应关系,命令用法见 MySQL 8.4 官方说明。PostgreSQL 可借助逻辑复制槽、导出快照和 LSN 建立基线,但复制槽会阻止所需 WAL 被清理,长时间全量会造成磁盘增长;接口与限制见 PostgreSQL 逻辑解码概念和复制槽函数。
日志起点一旦过期,正确动作通常是阻断并重建基线,而不是猜一个“最近位置”。MySQL 官方指出损坏或短缺的 binary log 会破坏复制同步,相关恢复行为见 MySQL binary log。对业务而言,重新全量很慢;对数据而言,带着未知缺口继续才是真正不可控。
冻结、受控双写与灰度切流
短时冻结适合写入可排队、可重试且追平时间可预测的系统。流程是:先停止入口写,再停后台 writer;记录冻结高水位;CDC 应用到高水位;执行差异和业务不变量校验;原子切换写入口;用新库完成一笔合成业务事务;最后恢复流量。冻结应返回明确的可重试错误或排队凭证,不能让客户端超时后猜测是否成功。
业务无法冻结时,双写必须升级为“受控迁移协议”。每个业务键由路由表确定权威侧,写请求携带稳定 operation_id 与单调业务版本;非权威侧接收的是可重放复制,不是第二次独立业务决定。跨库不具备原子提交时,要有 outbox/inbox、幂等键、重试上限和差异队列。禁止使用“先写 A,再尽力写 B,失败打日志”作为切换方案,因为日志并不能恢复客户端已观察到的成功。
读灰度可以较早开始,但应采用影子读:真实响应仍来自权威侧,异步读取目标并比较规范化结果,不把比较延迟加入用户请求。随后按稳定租户或主键哈希逐步让目标承担读流量,排除报表缓存、只读副本延迟和排序差异。写灰度则必须按数据所有权分片,不能对同一订单随机 5% 写旧库、95% 写新库。
项目接入至少需要三个独立开关:read_authority、write_authority 和 shadow_compare。一个 use_new_db=true 无法表达读已切、写未切、影子比较开启的中间状态,也无法在故障时只回退受影响维度。开关变更应审计到迁移 ID、配置版本、审批人和生效实例集合,并能证明旧配置实例已经排空。
回切窗口不是一个开关
切到新库后,如果新库继续产生写入,而旧库停止接收变化,回切能力会随第一笔新写立刻消失。真正的回切窗口只有三种建立方式:持续把新权威侧变更反向同步到旧库;把新写可靠记录为可重放日志并验证恢复速度;或者明确接受窗口内新写丢弃,并取得业务授权。第三种通常只适用于可重建缓存或一次性演练数据。
反向同步不能与正向 CDC 形成环。事件要带 origin、事务身份和版本,连接器过滤自身产生的回流;删除事件也必须保留,否则旧库会复活已经删除的数据。回切前重新执行追平与校验门槛,确认旧库的已应用水位覆盖新库高水位,再原子切回唯一 writer。发现分歧时先冻结两侧,不要让“修复脚本”和在线写继续竞争。
回切数据路径必须在切流前演练,而不是故障后临时设计。完整路径是“新权威侧提交 -> 新侧日志/outbox -> 防环反向传输 -> 旧侧幂等条件写 -> 旧侧已应用水位 -> 受影响分块与业务不变量重校验”。演练至少注入一次更新、一次删除和一次连接器重启,预期证据是三个操作保持事务身份和同键顺序、反向 offset 重启后连续、旧侧水位覆盖新侧演练高水位、重校验无差异。反例是只回放更新而丢弃删除:行数可能接近,回切后已删除对象却会复活,此时必须阻断 ROLLBACK_READY。
回切窗口的关闭条件应包含:新库跨业务周期稳定、备份恢复演练通过、旧版本应用已无法写旧库、所有延迟任务已迁移、回滚日志已归档、业务 owner 接受退出。此后回退应走“从新权威侧重新迁移”的灾备流程,而不是继续保留一套无人维护的旧双写链。
故障证据怎样从现象反推原因
看到目标行数停止增长,先区分源端没有提交、传输没有读取、sink 没有提交。源高水位前进而 transport 水位不动,检查日志权限、日志过期和连接器错误;transport 前进而 applied 不动,检查目标约束、死锁、重试与批次事务;三者都前进但差异增加,检查过滤规则、主键、DDL 和删除传播。
看到延迟瞬间归零也不能立即切换。连接器重启后 offset 被重置,可能从“当前”开始消费并把历史缺口伪装成零延迟。证据必须包含持久化 offset 身份、最近已提交事务、重启前后连续性和目标端落库记录。至少一次演练应主动暂停 sink、让源继续写、恢复后确认从旧 applied watermark 追平,而不是从新 source watermark 跳过。
看到旧库冻结后仍有写入,要按数据库账号、客户端地址、应用版本和 SQL 标签定位 writer。共享高权账号会让定位失败,因此迁移前就应为在线服务、定时任务、人工修复和 CDC 分配不同身份。切换门禁应查询冻结高水位后的提交者,不满足白名单就自动退回 CatchingUp。
权限、敏感数据与容量成本
迁移账号常同时接触全库数据、日志流和 schema 元数据,影响面比普通应用账号更大。源端读取、复制日志、目标写入、校验读取和切流控制应使用不同服务身份;凭证放入密钥管理系统,通过短期注入交付,不写进 YAML、命令历史或错误截图。证书私钥、连接串、binlog/WAL 样本、差异报告和死信事件都按生产数据分级处理。
日志和指标默认只记录迁移 ID、对象名哈希、位点、计数与错误码。为了排障临时采集行内容时,应限制列、脱敏、加密、设置访问审计和销毁期限。合规删除在迁移期间尤其容易失效:若目标快照重新带回已删除主体,或 tombstone 被过滤,迁移会构成数据复活。
容量预算不能只算目标表大小。源端全量扫描会增加 IO、buffer churn 和长事务历史版本;日志保留要覆盖最坏追平时间;传输层要容纳峰值变更与重试;目标端会同时承担导入、索引、校验和线上灰度读。控制器应设置源端负载上限、CDC backlog 上限和目标磁盘水位,越线后降速或暂停全量,而不是牺牲在线业务。
成本治理关注资源持续时间而非产品价格:临时实例、跨网络出口、重复存储、长日志保留和影子查询都会随迁移拖延增长。迁移 owner 每个阶段都应记录资源清单和销毁条件,避免项目“完成”后复制槽、连接器、临时账号和高规格目标继续运行。
退出时清理什么,保留什么
退出顺序应先撤销写能力,再清理传输。确认新库成为长期权威后,禁用旧应用凭证和旧写路由,停止正向/反向连接器,记录最终位点与配置摘要,删除临时复制槽或连接器任务,最后按保留策略处理旧库。直接先删连接器会丢失最后恢复证据,直接先删旧库会失去回切和差异定位入口。
本地实验可安全删除生成目录:
rm -rf cutover-good cutover-bad-origin cutover-bad-split
# PowerShell: Remove-Item -Recurse -Force cutover-good, cutover-bad-origin, cutover-bad-split生产中需要长期保留的是状态迁移记录、快照清单、低/高/已应用水位、配置哈希、校验报告、差异处置、切流审计和恢复演练结果;不应长期保留明文凭证、完整行样本和无访问控制的日志副本。团队把迁移平台 owner、数据库 owner、业务 owner、安全 owner 和切流指挥分别写入运行手册,并约定谁有权冻结、谁能改变 authority、谁判断回切、谁批准销毁。
一场迁移只有进入 Exited 才算结束:唯一权威侧已经稳定,旧 writer 已被技术性阻断,回切窗口已经按证据关闭,临时权限与资源已经回收,任何新故障都能从新库备份与变更历史恢复。目标库“看起来有数据”只是这条路最早的一站。
