1. 项目背景
业务场景:本地生活电商产品经理在需求评审会上宣布:“订单的收货地址要从字符串改成结构化对象——省/市/区/详细地址各自独立字段,方便后期做区域配送成本分析。” 开发团队瞬间陷入焦虑——2000 万历史订单中的 address 字段还是老格式"广东省深圳市南山区科技园路1号",直接把旧数据全部转换成新格式?线上不能停服务;用 updateMany 一次全量更新?2000 万行要跑 6 个小时,期间锁冲突导致正常下单超时;双写新旧字段同时维护?代码会越来越乱且数据不一致风险高。
痛点:Schema 演进是 MongoDB 项目中不可避免的痛——字段改名、类型变换、嵌套变平铺、冗余字段的引入和淘汰。没有系统化的数据迁移方案,团队要么用updateMany暴力迁移导致线上故障,要么在应用层兼容 5 种历史 Schema 版本导致代码腐烂,要么因为害怕迁移而不敢改字段设计,导致技术债务越积越重。
2. 项目设计
小胖(满脸愁容):大师,PM 要把订单地址从字符串改成结构化对象。2000 万条历史订单怎么处理?updateMany 一跑估计到明天都跑不完!
大师:数据迁移不能一把梭。线上不停机的 Schema 迁移有四部曲——双写、灰度读、后台迁移、清理旧格式。这不是技术选型不同,而是唯一可行的零停机路径。
小胖:双写?灰度读?说人话!
大师:分四步走:
双写(Dual Write):修改写入代码,下单时间时写入新字段
addressStruct: {province, city, district, detail}和旧字段address: "省市区地址"。新订单两个字段都有,旧订单只有旧字段。灰度读(Dark Read / Dual Read):读取时,先读新字段,如果不存在(说明是老数据)fallback 到旧字段。但不对外暴露——这只是用来验证新字段逻辑正确性。
后台迁移(Backfill):用一个后台脚本分批(batch)把所有旧数据的旧字段迁移为新字段。每批处理 1000 条,批次之间休眠 100ms 不抢线上资源。用
bulkWrite+ordered: false保证容错。清理旧格式(Cleanup):所有数据迁移完毕后,下线旧字段的写入代码。再跑一轮迁移确保零遗漏,最后
$unset移除旧字段。
技术映射:这是一个经典的 零停机数据迁移(Zero-Downtime Migration) 模式,适用于任何数据库。核心在于——迁移过程中系统始终兼容新旧两种数据格式。
小白(追问):那如果我要改的不是新增字段,而是字段类型变化——比如 age 从 String 变成 Int,双写没法同时兼容 String 和 Int 吧?
大师:这种情况就要引入 Schema 版本号(schemaVersion)。每个文档加一个版本字段——v1 的 age 是字符串,v2 的 age 是整数。代码里根据schemaVersion做不同的读写逻辑。
技术映射:Schema 版本号是处理不兼容变更的标准模式。v1→v2 时旧数据无需立刻迁移,而是逐步通过后台迁移脚本从 v1 升级为 v2。
小白:那大批量后台迁移怎么做到不影响线上?一跑 updateMany 就把 CPU 占满了吧?
大师:分批(batch)+ 限速(throttle)+ 断点续跑。关键参数:
- 每批次大小:500-2000 条(太大则单批耗时过久影响复制延迟;太小则迁移总时间过长)
- 批次间隔:50-200ms(给线上流量留出 CPU 和 IO 资源)
- 迁移阀值:每次迁移只改"创建时间 < 某个时间点"的数据,避免追着新数据跑(新数据已由双写覆盖)
- 断点续跑:记录迁移到哪个
_id了,重启后从该位置继续
小胖:还有 Change Streams 可以辅助迁移?我知道那是实时推送变更的……
大师:对。迁移期间如果线上有新写入(在你迁移的窗口内),新数据的旧字段和它一起出现。迁移脚本完成旧数据批量处理后,剩余的"增量"可以通过 Change Streams 实时补写——监听 insert/update 事件,发现新数据或更新数据缺少新字段时触发补写。这让迁移的总耗时从"批量全量"变成"批量历史 + 实时增量",对线上影响最小。
大师(总结):Schema 演进四项基本功——双写做兼容、灰度读做验证、后台迁移做补齐、版本号管理做可追溯。额外一个加分项——迁移前先在测试环境做一次完整演练,验证迁移耗时和数据一致性。
3. 项目实战
3.1 环境准备
沿用现有 MongoDB 环境。
3.2 分步实现
步骤一:模拟历史数据——添加 schemaVersion 字段
目标:创建带有老格式 address 的数据,同时添加版本字段。
use local_life db.orders_migration.drop()// 插入 10 万条"历史订单"(旧格式:address 为字符串)for(letbatch=0;batch<10;batch++){constdocs=[]for(leti=0;i<10000;i++){constidx=batch*10000+i docs.push({orderNo:"MIG"+String(idx).padStart(8,'0'),userId:"U"+(idx%500),address:"广东省深圳市南山区科技园路"+(idx%100)+"号",// 旧格式// addressStruct: 不存在(新格式缺失)schemaVersion:1,// 文档版本号totalAmount:NumberDecimal("99.00"),status:"已完成",createdAt:newDate(2025,0,1,0,0,0,idx),updatedAt:newDate()})}db.orders_migration.insertMany(docs,{ordered:false})}print("历史订单:",db.orders_migration.countDocuments(),"条")// 插入一些"新订单"(新格式:addressStruct)for(leti=0;i<100;i++){db.orders_migration.insertOne({orderNo:"MIG_NEW_"+String(i).padStart(5,'0'),userId:"U_NEW",address:"广东省深圳市南山区科技园路1号",// 旧格式(双写)addressStruct:{province:"广东省",city:"深圳市",district:"南山区",detail:"科技园路1号"},// 新格式schemaVersion:2,totalAmount:NumberDecimal("199.00"),status:"已完成",createdAt:newDate(),updatedAt:newDate()})}print("新订单:",db.orders_migration.countDocuments({schemaVersion:2}),"条")步骤二:应用层双写模拟
目标:模拟 Java 代码中的双写逻辑。
// === 应用层代码逻辑模拟 ===// 下单时的双写(同时写新旧格式)functioncreateOrder(orderData){constdoc={orderNo:orderData.orderNo,userId:orderData.userId,// 旧格式:字符串拼接address:`${orderData.province}${orderData.city}${orderData.district}${orderData.detail}`,// 新格式:结构化对象addressStruct:{province:orderData.province,city:orderData.city,district:orderData.district||"",detail:orderData.detail},schemaVersion:2,totalAmount:orderData.totalAmount,status:"待支付",createdAt:newDate(),updatedAt:newDate()}db.orders_migration.insertOne(doc)returndoc.orderNo}// 测试双写createOrder({orderNo:"MIG_DUAL_001",province:"上海市",city:"上海市",district:"浦东新区",detail:"陆家嘴金融中心A座",totalAmount:NumberDecimal("999.00")})// 验证双写结果constdual=db.orders_migration.findOne({orderNo:"MIG_DUAL_001"})print("双写结果:")print(" 旧格式:",dual.address)print(" 新格式:",JSON.stringify(dual.addressStruct))print(" 版本:",dual.schemaVersion)步骤三:灰度读实现
目标:实现兼容新旧格式的读取逻辑。
// === 读取逻辑:优先读新格式,fallback 到旧格式 ===functiongetOrderAddress(orderNo){constorder=db.orders_migration.findOne({orderNo},{addressStruct:1,address:1,schemaVersion:1})if(!order)returnnull// 优先用新格式if(order.addressStruct){return{source:"新格式(v2)",province:order.addressStruct.province,city:order.addressStruct.city,district:order.addressStruct.district,detail:order.addressStruct.detail}}// 降级到旧格式——需要做字符串解析// 生产中建议用正则或假设固定格式constoldAddr=order.addressreturn{source:"旧格式(v1)—降级解析",province:oldAddr.slice(0,3),raw:oldAddr// 实际生产环境可能有专门的地址解析服务}}print("新格式读取:",JSON.stringify(getOrderAddress("MIG_DUAL_001")))print("旧格式读取:",JSON.stringify(getOrderAddress("MIG00000000")))步骤四:分批后台迁移
目标:编写分批迁移脚本,将 v1 数据升级为 v2。
// === 后台分批迁移函数 ===functionbackfillAddress({batchSize=1000,limit=1000,throttleMs=50}){letmigrated=0letlastId=nullwhile(true){constquery={schemaVersion:1,// 只迁移 v1 数据addressStruct:{$exists:false}// 且尚未被迁移}// 基于 _id 做游标式分批(避免 offset 深分页)if(lastId)query._id={$gt:lastId}constbatch=db.orders_migration.find(query).sort({_id:1}).limit(batchSize).toArray()if(batch.length===0)break// 构建 bulkWrite 操作constbulkOps=batch.map(doc=>({updateOne:{filter:{_id:doc._id},update:{$set:{// 这里简化为解析旧格式(实际生产可能需要更复杂的地址解析)addressStruct:{province:doc.address.slice(0,3),city:doc.address.slice(3,6),district:"",detail:doc.address.slice(6)},schemaVersion:2,updatedAt:newDate()}}}}))// 执行批量更新try{constresult=db.orders_migration.bulkWrite(bulkOps,{ordered:false})migrated+=result.modifiedCount}catch(e){print("批量写入部分失败:",e.message)}lastId=batch[batch.length-1]._id migrated+=batch.lengthif(migrated%5000===0){print(`已迁移${migrated}条... (lastId:${lastId})`)}// 批次间隔,避免 CPU 满载sleep(throttleMs)if(migrated>=limit)break}return{migrated}}// 启动迁移conststartTime=Date.now()constresult=backfillAddress({limit:20000,throttleMs:50})constelapsed=(Date.now()-startTime)/1000print(`迁移完成:${result.migrated}条, 耗时${elapsed.toFixed(1)}s`)// 验证迁移结果constremainingV1=db.orders_migration.countDocuments({schemaVersion:1,addressStruct:{$exists:false}})print("剩余 v1 文档:",remainingV1)步骤五:Change Streams 增量补写
目标:用 Change Streams 实时监听新写入并补写新字段。
// === Change Streams 实时补写 ===// 使用 watch() 监听 orders_migration 的插入操作// 对每个新的 insert 补写 addressStructfunctionstartChangeStreamWatcher(){constpipeline=[{$match:{operationType:"insert","fullDocument.schemaVersion":{$exists:false}// 只关注未标记版本的}}]constchangeStream=db.orders_migration.watch(pipeline)print("Change Stream 监听已启动...")// 在 mongosh 中,用 cursor 的 tryNext 非阻塞获取// 这里只演示模式,实际后端代码中用 while(!cursor.isExhausted())// 模拟处理一个事件constnext=changeStream.tryNext()if(next){print("收到变更:",next.operationType,next.fullDocument.orderNo)// 补写 addressStructdb.orders_migration.updateOne({_id:next.fullDocument._id},{$set:{addressStruct:{/* 解析逻辑 */},schemaVersion:2}})}else{print("当前无新变更")}returnchangeStream}// 启动(在 mongosh 中这是阻塞的,仅演示模式)// const stream = startChangeStreamWatcher()// 生成一个新插入事件来测试db.orders_migration.insertOne({orderNo:"MIG_STREAM_TEST",userId:"STREAM_U",address:"北京市朝阳区望京SOHO",// 没有 addressStruct——触发 Change StreamtotalAmount:NumberDecimal("299.00"),status:"待支付",createdAt:newDate()})print("触发事件已插入")步骤六:清理旧字段——最终下线
目标:迁移完成后从所有文档中移除旧格式字段。
// 前置条件:所有文档的 schemaVersion 都已升级到 2constallMigrated=db.orders_migration.countDocuments({$or:[{schemaVersion:{$ne:2}},{addressStruct:{$exists:false}}]})if(allMigrated===0){print("全部文档已迁移,开始清理旧字段...")// 分批 $unset 移除旧字段constcleanupResult=db.orders_migration.updateMany({address:{$exists:true}},[{$set:{schemaVersion:2}},{$unset:"address"}// 或保留 address 作为 fallback 也可])print(`清理旧字段: 修改了${cleanupResult.modifiedCount}个文档`)}else{print(`仍有${allMigrated}个文档未迁移,无法清理`)}// 验证print("残留旧字段的文档:",db.orders_migration.countDocuments({address:{$exists:true}}))3.3 完整代码清单
| 文件 | 用途 |
|---|---|
mongodb-lab/scripts/ch23-create-migration-data.js | 构造迁移测试数据 |
mongodb-lab/scripts/ch23-dual-write.js | 双写逻辑模拟 |
mongodb-lab/scripts/ch23-gray-read.js | 灰度读逻辑 |
mongodb-lab/scripts/ch23-backfill.js | 分批后台迁移脚本 |
mongodb-lab/scripts/ch23-change-stream-watcher.js | Change Streams 增量补写 |
mongodb-lab/scripts/ch23-cleanup.js | 清理旧字段 |
3.4 测试验证
use local_life// 1. 验证双写:新订单同时包含两种格式constdual=db.orders_migration.findOne({orderNo:"MIG_DUAL_001"})print("双写验证:",dual.address&&dual.addressStruct?"PASS":"FAIL")// 2. 验证灰度读:旧格式文档能降级读取constold=db.orders_migration.findOne({orderNo:"MIG00000000"})print("灰度读旧格式:",old.address?"PASS (有旧字段)":"FAIL")// 3. 验证分批迁移:剩余 v1 文档数显著减少constremainingV1=db.orders_migration.countDocuments({schemaVersion:1})print("迁移后 v1 文档:",remainingV1,remainingV1<80000?"PASS (减少了)":"需继续迁移")// 4. 验证断点续跑// 如果迁移中断,重新运行 backfillAddress() 应该从上次 lastId 继续print("\n=== Schema 演进验证完成 ===")4. 项目总结
4.1 Schema 演进策略速查
| 变更类型 | 兼容性 | 策略 | 示例 |
|---|---|---|---|
| 新增字段 | 向后兼容 | 直接加字段,旧文档视为 null | 新增discountRate |
| 删除字段 | 向前兼容 | 代码不再写入,后续批量 $unset | 移除废弃的legacyId |
| 字段改名 | 不兼容 | 双写 + 迁移,最后删旧字段 | addr→address |
| 类型变更 | 不兼容 | 新字段 + schemaVersion + 迁移 | age: String → Int |
| 嵌套变平铺 | 不兼容 | 同上 | {addr:{city:"sz"}}→{city:"sz"} |
| 数组变对象 | 不兼容 | 同上 | tags:["a","b"]→{a:true, b:true} |
4.2 适用场景
本章迁移模式适用:
- 业务需求变更——地址、商品规格、用户画像字段的格式演进。
- 性能优化——将冗余字段(快照)加入到历史订单中避免 $lookup。
- 规范化重构——将嵌入过深的结构平铺到顶层。
- 多版本数据共存——微服务灰度发布中不同版本的服务共享同一个数据库。
不适用场景:
- 全量删除或归档(已不用的数据直接移到归档库)。
- 需要事务全量的跨集合一致性校验——迁移过程中新的写入和旧的迁移交错,不能用事务一次性覆盖。
4.3 注意事项
| 注意事项 | 说明 |
|---|---|
| 迁移脚本必须幂等 | 支持断点续跑,重复执行同一批数据不会产生副作用 |
bulkWrite的ordered: false | 避免单条错误阻塞后续整批的更新 |
| 限速不要完全停止 | throttleMs 太久会拉长迁移总时间;用批次之间的 sleep 而非单条的延迟 |
| 监控复制延迟 | 迁移期间观察 Oplog 窗口和 Secondary 延迟变化 |
| 迁移期间不要删除索引 | 老的索引可能正在被迁移脚本的 find 查询依赖 |
4.4 常见踩坑经验
故障案例一:updateMany 全量迁移导致写入阻塞
某团队用db.orders.updateMany({}, {$set:{newField:"default"}})初始化一个新字段。2000 万条数据跑 updateMany 的场景下,数据库的写入队列被挤爆——正常的业务写入请求排队等待锁释放。根因:updateMany 持有文档的意向锁(Intent Lock)期间,其他写入被阻塞。解决:改为分批 bulkWrite,每批 500-1000 条,批次间 50ms 释放锁窗口给业务。
故障案例二:双写期间新字段和旧字段不一致
某订单中address和addressStruct来自两段不同的代码逻辑——address由下单服务拼接,addressStruct由地址解析服务异步填充。存在时间窗口:用户修改了收货地址后,addressStruct不是最新的。根因:双写逻辑不同步。解决:双写必须在同一个代码路径、同一个事务内完成——要么都写,要么都不写,杜绝异步写入。
故障案例三:迁移脚本忘了加 schemaVersion 检查导致无限循环
某迁移脚本查询{addressStruct: {$exists: false}},没有加schemaVersion: 1。迁移脚本运行过程中,线上新的写入生成 v2 数据但恰好缺失addressStruct——迁移脚本把这些新数据也改成 v1 格式,覆盖了新代码的正确写入。根因:没有区分"需要迁移"和"不需要迁移"的数据。解决:迁移条件加上 schemaVersion 或者 createTime < 迁移截止时间。
4.5 思考题
- 如果迁移过程中需要将一个字段从 Double 改为 Decimal128,但由于 Decimal128 和 Double 无法直接转换(需要以字符串为中间态),迁移脚本应该怎么写?
- 有多个微服务共享一个数据库,其中一个服务先发布,开始写新字段;另一个服务未发布,还在读旧字段——如何避免老服务因新字段类型不兼容而报错?
(答案将在第 24 章末尾揭晓)
上一章思考题答案:
$merge在分片集群中要求on字段必须包含分片键——因为$merge需要确定 upsert 的目标文档在哪个分片上。如果不包含分片键,mongos 无法确定_id对应的数据应该写到哪个 shard——这会导致$merge失败或者 mongos 需要向所有分片广播写入请求(性能差且可能出现分片键冲突)。
allowDiskUse是"临时把溢出的分组数据写磁盘"——这是应急措施,磁盘 I/O 远慢于内存,结果是聚合能完成但慢数倍到数十倍。增大 WiredTiger 缓存是"增加工作内存"——让更多数据在内存中处理,消除磁盘溢写的需求,从根本上加速聚合。两者适用场景:查询量不大但分组过多 →allowDiskUse兜底;聚合是高频业务 → 增大缓存,减少磁盘溢写触发几率。
延伸阅读与资源
python入门:Rquests从菜鸟脚本到企业级SDK的网络实战圣经
Milvus向量数据库实战修炼:从 0 到 1精通向量检索与生产落地
后端工程师的 AI 转型第一课:Ollama 与私有化大模型实战
10倍开发者的 Dify 魔法书:从零构建全栈 AI 应用
后端工程师转型AI第一课-Ollama 与私有化大模型实战
大型语言模型(LLM) vLLM 高性能推理落地实战
Agent开发之LlamaIndex 实战修炼与源码进阶
大语言模型Transformers 实战修炼与源码剖析