1500万条质检数据从MySQL到MongoDB:多线程迁移方案设计与压测实录
为什么写这篇文章:两年多前,领导安排我做过一次千万级数据的迁移。面试中发现面试官对此兴趣很大,所以重新整理思路,并用多线程模拟复现了当时的方案。
任务背景:
当时领导负责另一个项目,需要做一个数据的迁移,但是他自己没时间就安排我做,从MySQL迁移1500万条数据进入Mongodb,至于为什么要这么做,由于我没参与那个项目,就不太了解。
约束条件:
不能修改源数据库:库是生产环境的,只能做读取操作,不能随便改。
数据不能丢失、不能重复。
写入顺序要与业务 ID 一致:下游系统会实时消费数据,顺序乱了会导致业务逻辑错误
技术选型
为什么不用DataX?
我首先调研了 DataX。DataX 是阿里开源的离线数据同步工具,功能很强大。但我发现了两个问题:
整批重试机制:DataX 写入失败时,会整批重试。如果网络抖动导致超时,第一次其实有部分成功了,第二次重试就会产生重复数据。
数据转换能力有限:我们的数据需要做数据清洗和格式转换并记录日志。DataX 的简单转换功能无法满足。所以,我决定自研迁移方案。
整体思路
架构:1个写入线程 + 3 个读取线程 + 阻塞队列
流程:
3 个线程并行,各自用主键游标分页从 MySQL 读取 2 万条数据。
每个线程对数据进行字段映射、数据清洗、日志打印后的 Document 列表放入优先级阻塞队列(按Document列表的最小业务id排序)。
单线程从队列取数据,批量写入 MongoDB。
写入失败时,获取失败下标,裁剪掉成功的前半部分,只重试失败的后半部分。
关键参数:
批量大小:2 万条/批
处理线程数:3
写入线程数:1
重试机制:裁剪重试(断点续传)
1、写入线程数量确定
首先为了满足写入顺序与业务id顺序完全一致。只能采用单线程向MongoDB插入数据,为了尽可能提高插入效率,可以从两个维度考量:
- 排除索引的影响,MongoDB的索引在大数据量的情况下,会极大影响写入性能,所以需要提前删除所有索引。
- 设置合理的插入批次,MongoDB支持最多10万/批的批量插入,但是并非批次越大越好,需要根据测试最终结果而定。我这里选择的2万/批,具体原因见下文。
2、重试机制
经过预处理清洗后的一批数据,虽然不会使MongoDB报错而无法插入,但是偶尔也会发生网络抖动导致失败。因此需要自己写一个重试机制
设计思路
批量写入失败时,MongoDB 会返回失败的下标。我把成功的前半部分裁掉,只重试失败的后半部分。
这样:成功的数据不会重复处理,失败的数据继续重试直到成功,不需要解析具体的错误原因,按下标裁剪就行,这就是我的“裁剪重试”机制。本质上是“断点续传”。
具体方法如下(伪代码)
try{mongoTemplate.insert(documents,"quality_detail_flat");}catch(BulkWriteExceptione){intfirstErrorIndex=e.getWriteErrors().get(0).getIndex();List<Document>retryList=documents.subList(firstErrorIndex,documents.size());insertWithRetry(retryList);// 递归重试}由于迁移过程中,全程有日志写入,所以极端情况下发现如果网络问题导致递归一直失败,可以人工暂停迁移。
3、读取线程数量确定
我提前分别对每批2000、10000、20000、40000进行过测试,数据分别如下:
| 批量大小 | 读取+数据处理 | MongoDB插入 | 单条平均时间(ms/条) |
|---|---|---|---|
| 2000条 | ~500ms | ~300ms | 0.15 |
| 10000条 | ~2050ms | ~800ms | 0.08 |
| 20000条 | ~3600ms | ~1300ms | 0.065 |
| 40000条 | ~6000ms | ~2300ms | 0.058 |
综合确定读取线程数量和写入批次
- 批次2000,单条速度跟其他批次差距过大,排除。
- 实测发现,两读一写总耗时22.5分钟,三读一写16.1分钟,多一个读线程能节省6分钟,且队列积压可控,所以选三读一写。
- 采用三读一写时,各种批次总导入时间和内存占用峰值对比如下:
| 批量大小 | 读+处理并行/批 | 写入/批 | 总耗时 | 队列积压速度 | 积压内存峰值 |
|---|---|---|---|---|---|
| 10000条 | 685ms | 800ms | 20分钟 | 读快115ms/批 | ~1.76GB |
| 20000条 | 1200ms | 1288ms | 16.1分钟 | 写快88ms/批 | ~770MB |
| 40000条 | 2000ms | 2300ms | 14.4分钟 | 写快300ms/批 | ~1.57GB |
可以看到:2万条每批,既能兼顾导入速度,又能大幅降低爆内存风险。
4、分段抓取与优先级队列
由于需要保持MongoDB数据的有序性,所以它每次插入的批次,最小的业务id必须等于已插入数据的最大id+1。这种批次数据我称之为:可写入数据(我的精准重试能保证数据完整性,所以可以简单粗暴的判断)
也就是说:
- 队列头必须是当前所有读取线程取到的最小业务id的数据。
- 要确保队头被取走之后,读取线程要在最短时间内,放入下一批可写入数据。
基于上述判断,我先用优先队列,按照每个批次数据的最小业务id排序,从而确保可以在队头直接取到最小数据。当然抓取的数据如果不是最小数据,会将该数据放回队列重新抓取。
再使用分段抓取,例如最开始时,线程1主键游标是1,偏移2万;线程2主键游标是20001,偏移2万。读取完一批之后,主键游标后移动6万位。确保尽快在队列放入可写入数据。
代码核心思路
1、分段抓取策略
3个读取线程错开起始位置,每个线程每次抓取2万条,读完一批后ID偏移6万位(3线程 × 2万条)。确保线程间数据不重叠。
for(longi=begin;i<maxId;i+=60000){//每轮查询20000条数据,查完sleep 3.6秒,这里默认id自增为1List<Long>selectList=selectFromMySQL(i,20000L);try{//模拟数据库读取和预处理清洗数据的总时间,设置为测试值/100,便于测试Thread.sleep(36);}catch(InterruptedExceptione){e.printStackTrace();}queue.add(selectList);}2、优先级队列与顺序控制
队列按批次最小ID排序。写入线程只有当前批次ID = 已写入最大ID + 1时才写入,否则放重回队列。保证写入顺序与源库一致。
- 取数据并判断
// 阻塞取数据,没数据就会休眠,不消耗CPUList<Long>batch=queue.poll(100,TimeUnit.MILLISECONDS);if(batch!=null&&batch.get(0)!=idx+1){queue.put(batch);continue;}- 定义优先队列排序规则
//优先阻塞队列,按照批次头部id排序publicstaticPriorityBlockingQueue<List<Long>>queue=newPriorityBlockingQueue<>(750,Comparator.comparing(batch->batch.get(0)));3、动态失败概率模型
批次越大,网络抖动导致失败的概率越高,采用指数模型模拟:
doublefailRate=0.05*Math.pow(batchIds.size()/40000.0,1.5);四、耗时模拟(按比例缩小100倍)
| 步骤 | 实际耗时 | 代码中sleep |
|---|---|---|
| 读取+处理 | 3600ms | 36ms |
| MongoDB插入 | 1300ms | 13ms |
模拟结果
这里大致可以看出队列堆积数量最大为(1500-1334)/2=84批,跟我当时测试环境跑的队列最大积压数50多批有差距,推测是模拟时时间缩小比例过大,造成的误差
这里总共执行时间*100倍之后大致时间为17.2分钟,跟实际时间非常接近。
完整代码见GitHub链接
https://github.com/jmingfu/Daily-Demo/blob/main/DataMigration