1500万条质检数据从MySQL到MongoDB:多线程迁移方案设计与压测实录

📅 2026/7/24 21:31:12
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){intfirstErrorIndexe.getWriteErrors().get(0).getIndex();ListDocumentretryListdocuments.subList(firstErrorIndex,documents.size());insertWithRetry(retryList);// 递归重试}由于迁移过程中全程有日志写入所以极端情况下发现如果网络问题导致递归一直失败可以人工暂停迁移。3、读取线程数量确定我提前分别对每批2000、10000、20000、40000进行过测试数据分别如下批量大小读取数据处理MongoDB插入单条平均时间ms/条2000条~500ms~300ms0.1510000条~2050ms~800ms0.0820000条~3600ms~1300ms0.06540000条~6000ms~2300ms0.058综合确定读取线程数量和写入批次批次2000单条速度跟其他批次差距过大排除。实测发现两读一写总耗时22.5分钟三读一写16.1分钟多一个读线程能节省6分钟且队列积压可控所以选三读一写。采用三读一写时各种批次总导入时间和内存占用峰值对比如下批量大小读处理并行/批写入/批总耗时队列积压速度积压内存峰值10000条685ms800ms20分钟读快115ms/批~1.76GB20000条1200ms1288ms16.1分钟写快88ms/批~770MB40000条2000ms2300ms14.4分钟写快300ms/批~1.57GB可以看到2万条每批既能兼顾导入速度又能大幅降低爆内存风险。4、分段抓取与优先级队列由于需要保持MongoDB数据的有序性所以它每次插入的批次最小的业务id必须等于已插入数据的最大id1。这种批次数据我称之为可写入数据我的精准重试能保证数据完整性所以可以简单粗暴的判断也就是说队列头必须是当前所有读取线程取到的最小业务id的数据。要确保队头被取走之后读取线程要在最短时间内放入下一批可写入数据。基于上述判断我先用优先队列按照每个批次数据的最小业务id排序从而确保可以在队头直接取到最小数据。当然抓取的数据如果不是最小数据会将该数据放回队列重新抓取。再使用分段抓取例如最开始时线程1主键游标是1,偏移2万线程2主键游标是20001偏移2万。读取完一批之后主键游标后移动6万位。确保尽快在队列放入可写入数据。代码核心思路1、分段抓取策略3个读取线程错开起始位置每个线程每次抓取2万条读完一批后ID偏移6万位3线程 × 2万条。确保线程间数据不重叠。for(longibegin;imaxId;i60000){//每轮查询20000条数据查完sleep 3.6秒,这里默认id自增为1ListLongselectListselectFromMySQL(i,20000L);try{//模拟数据库读取和预处理清洗数据的总时间设置为测试值/100便于测试Thread.sleep(36);}catch(InterruptedExceptione){e.printStackTrace();}queue.add(selectList);}2、优先级队列与顺序控制队列按批次最小ID排序。写入线程只有当前批次ID 已写入最大ID 1时才写入否则放重回队列。保证写入顺序与源库一致。取数据并判断// 阻塞取数据没数据就会休眠不消耗CPUListLongbatchqueue.poll(100,TimeUnit.MILLISECONDS);if(batch!nullbatch.get(0)!idx1){queue.put(batch);continue;}定义优先队列排序规则//优先阻塞队列,按照批次头部id排序publicstaticPriorityBlockingQueueListLongqueuenewPriorityBlockingQueue(750,Comparator.comparing(batch-batch.get(0)));3、动态失败概率模型批次越大网络抖动导致失败的概率越高采用指数模型模拟doublefailRate0.05*Math.pow(batchIds.size()/40000.0,1.5);四、耗时模拟按比例缩小100倍步骤实际耗时代码中sleep读取处理3600ms36msMongoDB插入1300ms13ms模拟结果这里大致可以看出队列堆积数量最大为1500-1334/284批跟我当时测试环境跑的队列最大积压数50多批有差距推测是模拟时时间缩小比例过大造成的误差这里总共执行时间*100倍之后大致时间为17.2分钟跟实际时间非常接近。完整代码见GitHub链接https://github.com/jmingfu/Daily-Demo/blob/main/DataMigration