当前位置: 首页 > news >正文

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~300ms0.15
10000条~2050ms~800ms0.08
20000条~3600ms~1300ms0.065
40000条~6000ms~2300ms0.058

综合确定读取线程数量和写入批次

  • 批次2000,单条速度跟其他批次差距过大,排除。
  • 实测发现,两读一写总耗时22.5分钟,三读一写16.1分钟,多一个读线程能节省6分钟,且队列积压可控,所以选三读一写。
  • 采用三读一写时,各种批次总导入时间和内存占用峰值对比如下:
批量大小读+处理并行/批写入/批总耗时队列积压速度积压内存峰值
10000条685ms800ms20分钟读快115ms/批~1.76GB
20000条1200ms1288ms16.1分钟写快88ms/批~770MB
40000条2000ms2300ms14.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
读取+处理3600ms36ms
MongoDB插入1300ms13ms

模拟结果

这里大致可以看出队列堆积数量最大为(1500-1334)/2=84批,跟我当时测试环境跑的队列最大积压数50多批有差距,推测是模拟时时间缩小比例过大,造成的误差

这里总共执行时间*100倍之后大致时间为17.2分钟,跟实际时间非常接近。

完整代码见GitHub链接
https://github.com/jmingfu/Daily-Demo/blob/main/DataMigration

http://www.jsqmd.com/news/1256822/

相关文章:

  • 宁波旅游出行服务GEO城市合伙人选型推荐哪家靠谱:2026年代理合作如何锁定源头技术与长效收益? - 企业新闻快传
  • HSTracker:3大功能让你的炉石传说胜率飙升50%
  • 2026国内EMBA QS排名|头部院校中立择校测评 - 品牌2026推荐
  • 【Linux】二十四.线程篇一《一文吃透Linux线程全面解析:概念、内存管理、优缺点与用途》
  • 苏州工业场景下,沃锐智能全自动打包机方案如何破解行业痛点,全自动打包机厂家 - 品牌推荐师
  • 3分钟实现Figma中文界面:设计师必备的FigmaCN完整指南
  • LTE站点传输板更换及网管侧割接配合问题处理案例
  • 青龙面板自动签到工具:30+平台一键打卡完整指南
  • 快消行业SFA系统选型指南:2026年技术对比与实战分析
  • 上海周大福老凤祥旧金回收对比:专柜回收套路多,连锁实体计价更实在 - 日常比对手册
  • Adobe-GenP 3.0终极指南:如何免费激活Adobe全家桶
  • deepseek 怎么导出 pdf?借助 AI 导出鸭,轻松解决各类导出格式困扰
  • ComfyUI-Easy-Use:AI图像创作工作流优化的终极解决方案
  • 终极macOS炉石助手指南:3分钟掌握HSTracker卡组跟踪工具
  • c++ this 指针的用途
  • 告别网盘下载困扰:九大平台直链助手LinkSwift完全指南
  • 社区志愿服务管理系统的设计与实现
  • 2026口碑好的亚太EMBA中立择校测评 - 品牌2026推荐
  • 2026年重庆除甲醛公司选择指南:真实测评避坑 - 信息热点
  • 宝塔+雷池WAF部署
  • 2026 年湖南成人高考高升专报考条件、流程、专业全攻略,零基础在职人可报 - GrowthUME
  • AMD Ryzen硬件调试终极指南:免费开源工具SMUDebugTool完全掌控
  • 如何摆脱网盘下载限制?这款免费工具让你重获自由下载体验
  • 【提示词工程黄金法则】:20年AI架构师首曝角色设定5步法,90%工程师至今用错
  • 终极指南:3步快速备份QQ空间完整历史数据
  • windows原生安装hermes-agent(不使用WSL ,Docker)
  • 2026有孵化器的全球EMBA中立择校测评 - 品牌2026推荐
  • 终极指南:如何快速掌握N_m3u8DL-RE的5大创新技术解析与实战应用
  • 奔驰贴隔热窗膜选什么品牌?2026分档推荐+型号搭配(C/E/S/GLC/EQ全系适用) - 资讯快报
  • 2026AI论文写作工具推荐 常见疑问一站式解答 - 资讯快报