文档数据库适合承载结构变化较快、嵌套层级较深的数据;关系型数据库则更适合跨实体关联、约束校验、事务处理与标准化报表。当系统从单一业务服务发展为多个服务协作,团队常常需要将部分核心数据迁入关系型数据库。
最容易失败的做法,是把每个文档字段都直接摊平成一张宽表,然后一次性导入。这样做通常会遇到四类问题:
- 文档的字段并不完全一致,历史数据可能缺字段、字段类型变化,甚至同名字段含义不同。
- BSON 中的
ObjectId、日期、二进制值、Decimal128等类型不能简单等同于普通 JSON 字符串。 - 迁移扫描进行时,源库仍可能有新增和更新;只执行一次全量扫描无法自然得到一致切面。
- 一旦任务中断,若没有稳定主键、幂等写入和失败记录,重跑可能产生重复数据或覆盖较新的数据。
因此,更稳妥的目标应当是建立一条可重复执行的迁移流水线:源文档可追溯、目标写入可幂等、异常数据可定位、迁移结果可对账,并能为最终切换留出增量同步路径。
设计原则:原文档保真与关系字段并存
对于首次迁移,不建议立刻把所有嵌套结构拆成大量关联表。一个较低风险的中间形态是:
- 为每份源文档保存一个唯一的来源键;
- 将原始文档以 JSONB 保存,作为审计、回溯和后续补数依据;
- 仅抽取已确认稳定、会参与查询或约束的字段为关系列;
- 将复杂数组、可变对象和暂未统一语义的字段保留在 JSONB 中;
- 后续根据真实查询需求,再从 JSONB 拆分出子表。
这种模型并不是否定关系建模,而是把“数据安全落地”和“业务语义建模”拆成两个可验证阶段。
以orders集合为例,假设其中多数文档具有_id、status、updatedAt,其余商品明细、地址和扩展字段结构并不稳定。目标表可以先这样设计:
CREATETABLEcustomer_document(source_keyTEXTPRIMARYKEY,statusTEXT,source_updated_at TIMESTAMPTZNOTNULL,payload JSONBNOTNULL,imported_at TIMESTAMPTZNOTNULLDEFAULTCURRENT_TIMESTAMP);CREATEINDEXidx_customer_document_statusONcustomer_document(status);CREATEINDEXidx_customer_document_source_updated_atONcustomer_document(source_updated_at);source_key不应仅仅使用_id的显示字符串。MongoDB 的_id允许使用不同 BSON 类型;将_id编码为规范扩展 JSON 后再保存,可以避免字符串_id与ObjectId恰好拥有相同可见文本时发生碰撞。
payload保存扩展 JSON 形式的完整文档。这样即使当前没有为某个字段设计关系列,也不会在迁移时丢失 BSON 类型信息。需要注意,JSONB 中保存的是 JSON 表示,不是 MongoDB 的原生 BSON;其价值在于可追溯和可解析,而不是承诺与源存储的物理字节完全一致。
迁移前的准备工作
1. 明确文档资格与时间字段
下面的示例要求每个待迁移文档都拥有updatedAt,且其值是 BSON 日期。该约束使增量同步和冲突处理有明确依据。
在执行前,应先检查不符合条件的数据:
db.orders.countDocuments({updatedAt:{$not:{$type:"date"}}})若结果不为零,不要静默把这些记录当成正常数据导入。可选择补齐时间字段、单独隔离,或为其设计另一套明确的排序与冲突规则。本文脚本会把此类文档写入拒绝文件,而不是写入目标表。
2. 创建最小环境与依赖
创建虚拟环境并安装依赖:
python-mvenv .venv..venv/bin/activate pipinstall"pymongo""psycopg[binary]"通过环境变量提供连接信息,避免将凭据写入代码或提交到仓库:
exportMONGO_URI='mongodb://migration_user:password@mongo.example.internal:27017/?authSource=admin'exportMONGO_DATABASE='business'exportMONGO_COLLECTION='orders'exportPOSTGRES_DSN='postgresql://migration_user:password@pg.example.internal:5432/warehouse'exportBATCH_SIZE='500'exportMAX_DOCS='0'MAX_DOCS=0表示不设上限。首次运行时建议设置一个较小的正整数,只验证连接、字段映射和拒绝记录是否符合预期;确认后再执行完整扫描。
可重跑的批量导入脚本
以下脚本以_id升序读取集合,以批为单位提交 PostgreSQL 事务。目标端采用ON CONFLICT写入,因此同一来源键重复处理不会新增重复行。对于已有记录,只有源文档的updatedAt不早于目标记录时才更新,避免旧扫描结果覆盖较新的内容。
importjsonimportosimportsysfromdatetimeimportdatetime,timezonefrombsonimportjson_utilfrompymongoimportMongoClient,ASCENDINGfrompsycopgimportconnectfrompsycopg.types.jsonimportJsonb REQUIRED=["MONGO_URI","MONGO_DATABASE","MONGO_COLLECTION","POSTGRES_DSN"]missing=[namefornameinREQUIREDifnotos.getenv(name)]ifmissing:raiseRuntimeError("缺少环境变量: "+", ".join(missing))batch_size=int(os.getenv("BATCH_SIZE","500"))max_docs=int(os.getenv("MAX_DOCS","0"))reject_path=os.getenv("REJECT_FILE","rejected_documents.jsonl")upsert_sql=""" INSERT INTO customer_document (source_key, status, source_updated_at, payload) VALUES (%s, %s, %s, %s) ON CONFLICT (source_key) DO UPDATE SET status = EXCLUDED.status, source_updated_at = EXCLUDED.source_updated_at, payload = EXCLUDED.payload, imported_at = CURRENT_TIMESTAMP WHERE EXCLUDED.source_updated_at >= customer_document.source_updated_at """defcanonical_key(value):returnjson_util.dumps(value,json_options=json_util.CANONICAL_JSON_OPTIONS)defextended_payload(document):text=json_util.dumps(document,json_options=json_util.CANONICAL_JSON_OPTIONS)returnjson.loads(text)defflush(pg_conn,rows):ifnotrows:returnwithpg_conn.transaction():withpg_conn.cursor()ascur:cur.executemany(upsert_sql,rows)defmain():mongo=MongoClient(os.environ["MONGO_URI"])collection=mongo[os.environ["MONGO_DATABASE"]][os.environ["MONGO_COLLECTION"]]processed=0rejected=0rows=[]withconnect(os.environ["POSTGRES_DSN"])aspg_conn,\open(reject_path,"a",encoding="utf-8")asreject_file:cursor=collection.find({}).sort("_id",ASCENDING).batch_size(batch_size)fordocincursor:updated_at=doc.get("updatedAt")ifnotisinstance(updated_at,datetime):reject_file.write(json_util.dumps(doc)+"\n")rejected+=1continueifupdated_at.tzinfoisNone:updated_at=updated_at.replace(tzinfo=timezone.utc)status=doc.get("status")ifnotisinstance(status,str):status=Nonerows.append((canonical_key(doc["_id"]),status,updated_at,Jsonb(extended_payload(doc)),))iflen(rows)>=batch_size:flush(pg_conn,rows)processed+=len(rows)rows.clear()print(f"已提交:{processed}",file=sys.stderr)ifmax_docs>0andprocessed+len(rows)>=max_docs:breakflush(pg_conn,rows)processed+=len(rows)print(json.dumps({"processed":processed,"rejected":rejected}))if__name__=="__main__":main()运行方式如下:
python migrate_documents.py脚本中的单批事务范围是一个工程取舍:批越大,事务提交次数越少,但锁持有时间、失败回滚范围和内存占用也会增加。应结合目标库负载、文档大小和维护窗口调整,不能假定某个固定批量适用于所有环境。
对账:从行数一致走向内容可信
迁移完成后,首先比较“符合迁移资格的源文档数”和“目标记录数”:
db.orders.countDocuments({updatedAt:{$type:"date"}})SELECTCOUNT(*)FROMcustomer_document;这只能发现明显漏数,不能证明内容正确。还应至少执行以下检查:
- 按
status分组比较数量,检查常见筛选维度是否发生映射错误。 - 按日期范围比较
updatedAt的最小值、最大值和分布。 - 从源端按
_id抽取固定样本,在目标端按source_key查询,比较关键业务字段与 JSONB 中的对应值。 - 统计拒绝文件行数,并逐条决定是修复源数据、增加映射规则,还是纳入独立迁移通道。
例如在目标端检查状态分布:
SELECTstatus,COUNT(*)FROMcustomer_documentGROUPBYstatusORDERBYstatusNULLSFIRST;对账结论应被记录为可复查的迁移证据,包括执行时间、源端筛选条件、拒绝记录位置、目标端统计 SQL 与处理决定。不要只依赖控制台中的“任务成功”日志。
在线写入场景:全量扫描不是最终切换方案
上面的脚本适合离线迁移,或作为在线迁移的历史数据装载阶段。但只要源集合仍在写入,就不能把一次扫描完成视为切换完成。
常见的收敛方式有两类:
- 停写切换:在维护窗口冻结源端写入,完成最后一次导入和对账后再切换读写流量。实现简单,但需要业务能够接受短暂停机。
- 增量同步切换:应用写入时同时记录可靠的变更事件,或从具备恢复位点的变更记录中持续消费;历史装载完成后,继续追赶增量,待延迟归零并对账通过后切换。
若使用updatedAt做增量条件,应采用(updatedAt, _id)这样的复合游标,而不是单独使用时间条件。原因是多份文档可能有相同更新时间;仅记录最后一个时间值会造成漏读或重复读。重复读可以由目标端幂等写入吸收,漏读则会直接造成数据不一致。
此外,删除操作不能仅靠扫描当前集合发现。若业务需要将删除同步到目标端,必须定义墓碑记录、软删除字段,或单独消费删除事件。具体选择取决于源端的保留策略与合规要求。
常见问题
为什么不直接将整份文档全部拆表?
完全拆表要求先明确每个嵌套对象的基数、主键、更新语义和查询需求。对历史文档结构差异较大的集合,过早拆表会放大迁移风险。先保留 JSONB 并抽取稳定字段,可以让后续建模建立在已验证的数据之上。
为什么要求updatedAt不能为空?
幂等写入需要判断哪份数据更新。没有可靠版本号、更新时间或变更序列时,迁移程序无法安全判定旧数据是否会覆盖新数据。若业务没有updatedAt,可改用单调递增版本号或独立变更日志,但必须明确其排序语义。
目标端 JSONB 能否替代所有关系表?
不能。JSONB 适合保留动态结构和审计载荷,但高频关联、唯一性约束、复杂聚合以及需要严格类型治理的字段,仍更适合进入明确的关系列或关联表。本文方案是迁移落地层,不是永久的数据模型结论。
任务中断后是否可以直接重跑?
在来源键稳定、目标表唯一约束存在且冲突更新规则正确的前提下,可以重跑。仍应确认源端的updatedAt语义可信,并检查拒绝文件;“可重跑”不等于可以忽略异常数据。
总结
可靠的数据迁移应当把导入程序视为一条可审计的数据产品链路,而不是临时脚本。保留规范化来源键和原始载荷,抽取少量稳定关系字段,以幂等写入保证重跑安全,再用行数、分组、样本和异常清单完成对账,可以显著降低文档模型迁往关系模型时的不可逆风险。
对于在线系统,历史全量导入只是第一阶段;只有补齐增量、处理删除语义、完成一致性对账并执行受控切换,迁移才算真正结束。