第 11 章 · 数据迁移、CDC 与数据治理
本章覆盖 不停机分片扩容、CDC、Schema 演进与数据治理,含 时间窗/存储手算例题 与小紫演练清单。
前置:ch10 分库分表;PaaS Kafka;ch08 综合案例。
11.1 分片扩容六阶段(4库→8库)
评估 → 双写 → 全量 → 增量追赶 → 切读 → 停双写 → 清理
| 阶段 | 耗时例题 | 回滚 |
|---|---|---|
| 双写 | 1天上线 | 关新写 |
| 全量 | 见下 | 删新分片 |
| 增量 lag<1s | 数小时 | 暂停切读 |
| 切读 1→100% | 4~8h | 切回旧路由 |
| 停双写 | 1h | 数据修复 |
全量时间例题
| 项 | 值 |
|---|---|
| 历史订单 | 1.2亿行 |
| 含索引约 1KB/行 | 120GB |
| 带宽 100MB/s 有效 | 纯传 ≈ 20min |
| 导入+建索引 ×3~5 | 1~2h/库,4库并行墙钟 2~4h |
增量例题:双写写TPS=250,单条binlog≈2KB → 500KB/s;全量3h期间增量 ≈ 1.6GB。
封网:大促前 72h 禁结构级迁移。
11.2 双写与对账
def write_order(user_id, row):
old_db, old_tbl = route(user_id, OLD_SHARDS=4)
new_db, new_tbl = route(user_id, NEW_SHARDS=8)
insert(old_db, old_tbl, row) # 主路径
try:
insert(new_db, new_tbl, row)
except Exception:
log.error("dual_write_new_failed"); # 补偿 job 修复
| 不一致 | 发现 | 修复 |
|---|---|---|
| 新分片缺行 | 每小时 COUNT 对账 | 从旧补 INSERT |
| 字段不同 | checksum 抽样 | 以旧为准 |
抽样:全量1.2亿不可扫;按 user_id%1000 抽 0.1% 桶 ≈ 12万行/小时可接受。
11.3 CDC 全链路
MySQL binlog(ROW) → Debezium/Canal → Kafka → ES/CH/Redis失效/数据湖
| 用途 | 延迟 | 幂等键 |
|---|---|---|
| 搜索 | <5s | (table,pk,version) |
| 报表 | 分钟 | event_id |
| 缓存失效 | <1s | product_id |
{
"connector.class": "io.debezium.connector.mysql.MySqlConnector",
"database.hostname": "mysql-order.prod-middleware.svc.cluster.local",
"topic.prefix": "shop",
"table.include.list": "shop_order_.*\\.t_order_.*",
"snapshot.mode": "initial"
}
def handle_event(ev):
key = (ev["table"], ev["pk"], ev["ts_ms"])
if redis.setnx(f"cdc:seen:{key}", 1):
upsert_es(ev)
11.4 Kafka 容量精算
| 项 | 值 |
|---|---|
| 日订单 | 80万 |
| 每单CDC 3事件 | 240万条/天 |
| 平均2KB | 4.8GB/天 |
| +商品变更 | +0.75GB |
| 合计 | 5.55GB/天 |
磁盘 = 日增量 × 保留天 × 副本 = 5.55 × 7 × 3 ≈ 117GB → 150GB PVC
| 峰值 | 计算 |
|---|---|
| 写 250TPS×3事件 | 750 msg/s |
| 750×2KB | 1.5MB/s → 3分区够 |
| 教学推荐 | 6分区(max(并行,6)) |