Administrator
发布于 2021-04-27 / 1015 阅读
25

分库分表后的数据迁移与双写方案

8200 万行数据,怎么搬到 32 个分片里去

接上一篇。分片规则定好了、ShardingSphere 配好了,接下来的问题是:老库里那 8200 万行数据怎么搬过去,而且不能停服。

这篇文章写的是迁移本身,跟分片设计是两件事,但难度可能更大。

先算停机方案要多久

最省事的做法是挂维护页,导出导入。我先估算了一下时间:

1. 停服、停止写入                                5 分钟
2. mysqldump 导出 t_order(68 GB)              约 42 分钟
3. 传输到新集群                                  约 8 分钟
4. 按分片规则拆分成 32 份并导入                  约 2 小时 10 分
5. 建索引                                       约 1 小时 20 分
6. 数据校验                                     约 25 分钟
7. 切流、启动、回归                              约 30 分钟
------------------------------------------------------------
合计                                           约 5 小时 20 分

这些数字是我在一台配置相同的备机上真跑一遍测出来的,不是估算。5 小时 20 分,而且这是在一切顺利的前提下。

业务方给的答案是:最多停 30 分钟,且只能在凌晨 2 点到 4 点。 停机方案直接否掉。

不停机方案:五个阶段

阶段一  双写上线(新老库都写,读老库)
阶段二  存量数据迁移(后台跑批,限速)
阶段三  增量追平 + 一致性校验
阶段四  灰度切读(1% → 10% → 50% → 100%)
阶段五  停写老库、清理

整个周期三周,其中大部分时间花在"观察"而不是"操作"上。

阶段一:双写

双写的核心原则是:老库是主,新库是辅。新库写失败不能影响业务。

@Transactional
public void createOrder(Order order) {
    // 1. 写老库,这是主流程,失败就抛异常回滚
    orderMapper.insert(order);

    // 2. 写新库,异步、失败只记录
    if (dualWriteEnabled) {
        try {
            CompletableFuture.runAsync(() -> {
                shardingOrderMapper.insert(convert(order));
            }, migratePool).exceptionally(e -> {
                // 记到一张表里,后续补偿
                migrateFailMapper.insert(new MigrateFailLog("INSERT", order.getId(), e.getMessage()));
                meterRegistry.counter("migrate.write.fail").increment();
                return null;
            });
        } catch (Exception e) {
            log.error("dual write submit fail, orderId={}", order.getId(), e);
        }
    }
}

为什么异步?因为同步写新库会让接口多一次网络往返,我们实测 P99 会涨 8~12 ms。而新库写失败了我们有补偿机制,不需要强一致。

双写开关用配置中心(我们用的 Nacos),可以随时关:

@NacosValue(value = "${migrate.dual-write:false}", autoRefreshed = true)
private boolean dualWriteEnabled;

这个开关救过我们一次。双写上线第二天,新库有一次网络抖动,写入大量失败,虽然不影响业务,但失败日志表涨得很快。我们把开关关了十分钟,等网络恢复再打开。

双写最容易忽略的一点:UPDATE 和 DELETE 也要双写

我们第一版只做了 INSERT 的双写,结果迁移完成校验时发现差异。原因是双写上线后老库被 UPDATE 过的行,新库里还是迁移时的旧版本。

@Transactional
public void updateOrderStatus(Long orderId, Integer status) {
    int rows = orderMapper.updateStatus(orderId, status);
    if (dualWriteEnabled && rows > 0) {
        CompletableFuture.runAsync(() -> {
            // 注意:新库要用 ShardingSphere 的主键路由,
            // 而 orderId 自带 user_id 基因,ShardingSphere 能算出分片
            shardingOrderMapper.updateStatus(orderId, status);
        }, migratePool).exceptionally(...);
    }
}

阶段二:存量迁移

迁移程序是自己写的。没用 DataX 这类通用工具,因为要做分片路由(按 user_id % 32 决定写到哪个库哪张表),通用工具配起来反而麻烦。

@Component
public class OrderMigrator {

    private static final int BATCH = 2000;
    private final RateLimiter limiter = RateLimiter.create(3000);  // 3000 行/秒

    public void migrate(long startId, long endId) {
        long cursor = startId;
        while (cursor < endId) {
            limiter.acquire(BATCH);

            List<Order> page = oldOrderMapper.selectByIdRange(cursor, cursor + BATCH);
            if (page.isEmpty()) break;

            for (Order o : page) {
                // 用 INSERT IGNORE,避免重复迁移时主键冲突
                shardingOrderMapper.insertIgnore(o);
            }
            cursor = page.get(page.size() - 1).getId() + 1;

            // 记录进度,中断可续跑
            progressMapper.updateProgress("t_order", cursor);
        }
    }
}

几个关键设计:

  • 按主键区间分批,不用 LIMIT offsetOFFSET 越大越慢,扫 8000 万行时后面的批次会慢到无法接受。用 WHERE id > ? LIMIT 2000 走主键索引,每批都是恒定的快。
  • 限速 3000 行/秒。我们试过不限速,主库的 CPU 直接从 40% 涨到 85%,业务接口的 P99 涨了 3 倍。限速后 CPU 只涨 6 个百分点。8200 万行按 3000/秒算需要 7.6 小时,我们跑了三个晚上。
  • 记录进度。迁移程序中断是常态(网络、GC、机器重启),必须能续跑。
  • INSERT IGNORE。双写可能已经把某些行写进去了(那些是迁移开始后才创建的订单),重复插入会主键冲突。
<insert id="insertIgnore">
  INSERT IGNORE INTO t_order (id, user_id, order_id, ...)
  VALUES (#{id}, #{userId}, #{orderId}, ...)
</insert>

阶段三:一致性校验

迁移完不等于完成,必须验证。全量比对的思路:按主键区间分块,每块算一个校验和,只比对校验和,不一致的块再逐行比。

public ChecksumResult checksumBlock(long start, long end) {
    // 老库
    String oldCrc = oldOrderMapper.checksumByIdRange(start, end);
    // 新库:要按 user_id 拆到 32 个分片,分别算再合并
    long newCrc = 0;
    for (int i = 0; i < 32; i++) {
        newCrc ^= shardMapper(i).checksumByIdRange(start, end);
    }
    return new ChecksumResult(start, end, oldCrc, newCrc);
}
<select id="checksumByIdRange" resultType="string">
  SELECT BIT_XOR(CRC32(CONCAT_WS('#',
      id, user_id, order_id, status,
      ROUND(IFNULL(amount,0),2),
      DATE_FORMAT(create_time, '%Y-%m-%d %H:%i:%s')
  )))
  FROM t_order WHERE id >= #{start} AND id < #{end}
</select>

BIT_XOR(CRC32(...)) 这个写法有个好处:它与行的顺序无关,只与内容有关。老库和新库的行顺序可能不同(不同分片的数据物理顺序肯定不同),用普通的 MD5(GROUP_CONCAT(...)) 会误报。

几个必须注意的点:

  • 浮点和金额DECIMAL 要统一精度,用 ROUND(amount, 2)
  • 时间DATETIME 格式化成字符串再算 CRC,避免时区或精度差异。我们的 create_timedatetime(不带毫秒),没这个问题,但 update_timedatetime(3),必须处理。
  • NULL 值CONCAT_WS 遇到 NULL 会跳过,导致 NULL 和空串算出同样的结果。我们用了 IFNULL(col, '<NULL>')

第一轮校验结果:8200 万行,分成 8200 个块(每块 1 万行),差异块 341 个,涉及 12,847 行。逐行比对后发现全部是同一个原因:这些行在双写上线之后、存量迁移跑到它们之前被 UPDATE 了,而双写的 UPDATE 那会儿还没上线(我前面说的第一版只做了 INSERT)。

修复方法:对这 12847 行按主键重新做一次覆盖迁移(用 REPLACE INTO 而不是 INSERT IGNORE)。重新校验,差异归零。

阶段四:灰度切读

校验通过后开始切读流量。切读比切写安全,因为读错了不会污染数据。

user_id 做灰度维度,因为分片键就是它,同一用户的数据一致性不会被打破:

public List<Order> listOrders(Long userId, Integer status) {
    if (readFromNew(userId)) {
        return shardingOrderMapper.selectByUserAndStatus(userId, status);
    }
    return oldOrderMapper.selectByUserAndStatus(userId, status);
}

private boolean readFromNew(Long userId) {
    int percent = migrateConfig.getReadPercent();     // 配置中心,1 / 10 / 50 / 100
    return userId % 100 < percent;
}

灰度节奏(每个阶段观察 2 天):

阶段读新库比例观察重点结果
D1-D21%有无报错、耗时对比正常,P99 从 3840 ms 降到 14 ms
D3-D410%慢查询、连接池发现 2 条 SQL 未走分片键,全路由
D5-D650%数据库连接数、CPU正常
D7-D9100%全量业务回归正常

第二阶段发现的 2 条全路由 SQL 是重要收获。它们是运营后台的查询,没带 user_id,ShardingSphere 会路由到全部 32 个分片。这种 SQL 在测试环境数据量小时看不出问题,上生产就暴露了。我们后来把它们改到了 ES,这也是我在 ShardingSphere 那篇里强调"跨库分页要解决"的原因。

阶段五:停写老库

读流量 100% 切过去、稳定运行一周后,才敢停双写。

顺序很重要:

  1. 先把双写开关关掉,观察 24 小时,确认没有问题(如果新库缺数据,这时候还来得及用老库补)
  2. 再从代码里删掉双写逻辑,上线
  3. 老库保留一个月只读,不删
  4. 一个月后备份归档

回滚预案

这是整个方案里我最花心思的部分。每个阶段都要能退回去:

阶段出问题怎么退
双写期间关配置开关,新库的脏数据直接 truncate
存量迁移期间停止迁移程序,新库数据废弃重来
灰度读期间readPercent 调回 0,秒级回退
100% 读之后仍可调回 0,因为老库一直在双写,数据是新的
停双写之后退不回去了,所以这一步前必须观察满一周

关键洞察:只要双写还在,我们就有退路。所以双写是整个方案的保险绳,它应该一直开到最后一刻。

踩到的另外三个坑

1. 自增 ID 冲突

老库的 t_order.idAUTO_INCREMENT,新库用的是我们自定义的雪花 ID(带基因位)。迁移时直接把老的 ID 搬过去,新老 ID 会混在一起。

好消息是我们的雪花 ID 从 2021-01-01 的秒偏移开始(约 2^31 量级起跳),而老自增 ID 只到 8300 万左右,数值区间不重叠。我提前算过这个,如果重叠就得做 ID 映射,那会复杂十倍。

2. 迁移期间的从库延迟

迁移程序虽然限速了,但依然是大量写入。MySQL 主从延迟一度到 40 秒,导致依赖从库的报表任务数据不准。

解决办法:迁移只在凌晨 1 点到 6 点跑,且监控 Seconds_Behind_Master 超过 10 秒就自动降速(把 RateLimiter 的速率调低)。

@Scheduled(fixedDelay = 10000)
public void adjustRate() {
    long lag = getSlaveLagSeconds();
    if (lag > 20) {
        limiter.setRate(1000);
    } else if (lag > 10) {
        limiter.setRate(2000);
    } else {
        limiter.setRate(3000);
    }
}

3. 双写的幂等

异步双写失败了记进 migrate_fail_log 表,补偿任务会重试。重试必须幂等,否则重复写会覆盖更新的数据。

我们的补偿逻辑用了 update_time 做版本判断:

UPDATE t_order SET status = #{status}, update_time = #{updateTime}
WHERE id = #{id} AND update_time <= #{updateTime}

这样迟到的旧更新不会覆盖新值。

时间线

D1        双写上线(只写了 INSERT)
D3        补上 UPDATE / DELETE 双写
D3-D6     存量迁移,每晚 1:00-6:00,累计 14.5 小时
D7        第一轮校验:341 个差异块,12847 行
D8        修复差异,重新校验通过
D9        灰度读 1%
D11       灰度读 10%,发现全路由 SQL
D13       灰度读 50%
D15       灰度读 100%
D22       关闭双写
D23       删除双写代码
D+30      老库归档

三周,其中真正"干活"的时间大概 30 小时,其余都在等观察期。

小结

  • 停机方案要 5 小时 20 分(实测),业务只给 30 分钟,只能走不停机双写。
  • 双写的核心原则:老库为主、新库为辅,新库失败不阻断业务。异步写 + 失败记录 + 补偿。
  • UPDATE 和 DELETE 同样要双写。我们第一版漏了,导致 12847 行数据不一致。
  • 存量迁移按主键区间分批(不用 OFFSET)、限速 3000 行/秒、记录进度可续跑、用 INSERT IGNORE
  • 一致性校验用 BIT_XOR(CRC32(CONCAT_WS(...)))与行顺序无关,这点很关键。注意处理 NULL、金额精度、时间格式。
  • 灰度按分片键维度切,配置中心控制,秒级可回退。
  • 双写是保险绳,要一直开到最后一刻。关掉双写就退不回去了,那之前必须观察满一周。

这个项目做完后我记了一条经验:数据迁移的难点不在于把数据搬过去,在于证明搬过去的数据是对的。我们花了 14.5 小时搬数据,却花了 5 天做校验和灰度。这个比例我觉得是对的,也可以反过来作为评估迁移工作量的参考。

参考