MongoDB 数据归档实战 一、背景需要对 MongoDB 数据库进行全量备份迁移数据规模如下· 总存储近 6TB· 表数量230 张· 最大单表11 亿条记录目标将数据导出为 Parquet 格式上传至 S3 进行长期保存。二、为什么选择 Parquet 格式在选型时考虑过 JSON、CSV、Parquet 三种格式最终选择 Parquet 的原因1. 存储效率高· 列式存储 压缩算法Zstd压缩比非常高· 相比 JSON 格式存储空间大幅节省· 实际压缩效果6TB MongoDB 数据压缩后约 500GB压缩比约 12:12. 便于后续数据处理· Hive / Spark 原生支持 Parquet可直接建外部表查询· 列式存储支持列裁剪只读取需要的字段查询效率高· Athena / Presto 可直接查询 S3 上的 Parquet 文件3. Schema 信息自包含· Parquet 文件自带 Schema 元数据不需要额外维护表结构· 支持复杂数据类型嵌套结构、数组等三、核心挑战与解决方案挑战 1本地磁盘空间不足问题6TB 数据无法全部存储在本地磁盘解决方案边导出边上传导出 - 上传 - 删除本地文件· 单个文件限制为 2GB超过后自动切割· 文件导出完成后立即上传 S3· 上传成功后删除本地文件释放磁盘空间· 实际磁盘占用不超过 50GB多进程并发时# 核心技术文件分卷 实时上传MAX_PARQUET_SIZE 2 * 1024 * 1024 * 1024 # 2GBif current_size table.nbytes MAX_PARQUET_SIZE:pqwriter.close()upload_file_to_s3(current_file) # 上传到 S3os.remove(current_file) # 删除本地文件# 创建新文件继续写入file_index 1new_file f{collection}_{file_index:03d}.parquet挑战 2MongoDB Schema 不一致问题MongoDB 是无 Schema 数据库同一张表不同文档的字段可能完全不同解决思路· 1. 先推断出所有可能的字段通过采样· 2. 统一 Schema所有字段都定义为 string 类型· 3. 导出时缺失的字段填充为 None推断字段的方法· 优先使用 MongoDB 聚合查询 $objectToArray覆盖率最高· 兜底方案多轮采样最小/最大 _id 随机采样 分段采样# 核心技术MongoDB 聚合查询推断字段pipeline [{$limit: batch_size * 10},{$project: {fields: {$objectToArray: $$ROOT}}},{$unwind: $fields},{$group: {_id: None, fieldNames: {$addToSet: $fields.k}}}]result list(coll.aggregate(pipeline, allowDiskUseTrue))fields result[0][fieldNames] # 获取所有字段名为什么要强制 string 类型踩坑经验如果让 PyArrow 自动推断类型会遇到类型不匹配错误。· 场景第一批数据某列全是 NonePyArrow 推断为 null 类型· 问题后续批次该列有实际值类型变为 string导致写入失败· 解决强制所有字段为 string 类型避免类型冲突# 核心技术PyArrow 强制 Schema# 先定义 Schema所有字段都是 stringarrow_schema pa.schema([pa.field(field_name, pa.string())for field_name in inferred_fields])# 用预定义的 Schema 创建 Arrow Tabletable pa.Table.from_pandas(df, schemaarrow_schema, preserve_indexFalse)# 这样即使某列全是 None也会被创建为 string 类型而非 null 类型挑战 3查询效率大表扫描慢问题11 亿条数据全表扫描非常慢解决方案按 _id 排序 流式游标 禁用超时· 利用 _id 索引避免全表扫描· 设置 no_cursor_timeoutTrue禁用游标超时默认 10 分钟· 流式遍历每批处理 20,000 条数据内存占用可控· 兜底机制如果游标异常重建游标从 last_id 继续# 核心技术流式游标 禁用超时BATCH_SIZE 20000# 设置游标选项query {_id: {$gt: last_id}} if last_id else {}cursor coll.find(query,sort[(_id, 1)],no_cursor_timeoutTrue # ⭐ 禁用 10 分钟超时限制).batch_size(BATCH_SIZE)# 流式遍历try:for doc in cursor:process_document(doc)last_id doc[_id]except Exception as e:# 兜底如果游标异常网络断开等重建游标继续if cursor in str(e).lower():cursor coll.find({_id: {$gt: last_id}}, sort[(_id, 1)],no_cursor_timeoutTrue)挑战 4导出过程中断怎么办问题导出 6TB 数据需要数小时期间可能因为网络、机器重启等原因中断解决方案断点续传机制· 每批数据处理完后保存检查点记录 last_id、已处理行数、文件索引· 程序重启后读取检查点从 last_id 继续导出· 避免重复导出节省时间# 核心技术检查点机制文件锁 原子重命名def save_checkpoint(collection, data):保存检查点支持多进程并发path f{CHECKPOINT_ROOT}/{collection}.jsontemp_path path .tmpwith open(temp_path, w) as f:fcntl.flock(f.fileno(), fcntl.LOCK_EX) # 文件锁json.dump(data, f)os.fsync(f.fileno()) # 强制刷盘fcntl.flock(f.fileno(), fcntl.LOCK_UN)os.replace(temp_path, path) # 原子重命名# 检查点内容checkpoint {last_id: str(last_id),file_index: 3,processed: 5000000}挑战 5导出效率太低问题单进程串行导出 230 张表预计耗时数十小时解决方案多进程并行导出· 使用 multiprocessing.Pool每个进程导出一张表· 进程数CPU 核心数 - 1避免 CPU 跑满影响系统· 性能提升5-10 倍8 核机器# 核心技术multiprocessing.Poolfrom multiprocessing import Pool, cpu_countworkers cpu_count() - 1 # 7 个进程假设 8 核tables [table1, table2, ..., table230]with Pool(processesworkers) as pool:results pool.map(export_collection, tables)# 每个进程独立导出一张表互不干扰挑战 6上传 S3 阻塞导出问题上传 2GB 文件需要数分钟阻塞导出进程解决方案异步上传· 使用 ThreadPoolExecutor 创建上传线程池2 个线程· 文件导出完成后提交到线程池异步上传· 导出进程继续处理下一批数据不等待上传完成# 核心技术ThreadPoolExecutor 异步上传from concurrent.futures import ThreadPoolExecutorupload_executor ThreadPoolExecutor(max_workers2)def upload_async(file_path, collection):异步上传到 S3s3_key f{S3_PREFIX}{collection}/{os.path.basename(file_path)}s3_client.upload_file(file_path, S3_BUCKET, s3_key)os.remove(file_path) # 上传成功后删除# 提交异步上传任务upload_executor.submit(upload_async, current_file, collection)# 不等待上传完成继续导出下一批数据四、技术方案架构整体流程· 1. 字段推断聚合查询 / 多轮采样· 2. 流式导出按 _id 分批查询20,000 条/批· 3. 类型统一强制 string 类型避免冲突· 4. 写入 Parquet每批写入2GB 自动切割· 5. 异步上传上传线程池不阻塞导出· 6. 断点续传检查点机制支持中断恢复· 7. 多进程并行230 张表并发导出关键技术栈技术用途PyArrowArrow Table 创建、Parquet 写入PandasDataFrame 数据处理PyMongoMongoDB 查询、聚合multiprocessing.Pool多进程并行导出ThreadPoolExecutor异步上传 S3fcntl / msvcrt跨平台文件锁Unix/Windowsboto3S3 上传客户端五、踩过的坑坑 1PyArrow 类型推断导致 Schema 不匹配现象导出到一半报错 Table schema does not match原因第一批数据某列全是 None推断为 null 类型后续批次有值类型变为 string解决先定义 Schema全 string再创建 Table坑 2MongoDB 游标超时已解决问题导出大表时可能遇到游标超时默认 10 分钟解决设置 no_cursor_timeoutTrue 禁用超时兜底即使禁用超时仍然捕获游标异常支持重建游标继续应对网络断开等情况坑 3多进程并发写入检查点文件冲突现象检查点文件内容损坏导致无法恢复原因多个进程同时写入同一个检查点文件解决文件锁 临时文件 原子重命名坑 4字段推断不完整现象导出到一半发现新字段Schema 不匹配原因单次采样覆盖不全解决多策略采样聚合查询 最小/最大 _id 随机 分段坑 5多进程负载不均衡现象3 个进程并发导出2 个进程很快完成1 个进程还在导出大表原因使用 pool.map() 静态分配任务某个进程可能分到多个大表解决改用 pool.imap_unordered() 动态任务队列进程完成一个任务后立即取下一个# 修改前静态分配负载不均衡results pool.map(export_collection, tables)# 进程 1: table1, table4, table7, ... (可能都是大表)# 进程 2: table2, table5, table8, ... (可能都是小表)# 进程 3: table3, table6, table9, ... (可能都是小表)# 修改后动态任务队列负载均衡results list(pool.imap_unordered(export_collection, tables))# 进程完成任务后立即从队列取下一个充分利用所有进程坑 6多进程 tqdm 进度条混乱现象多个进程的进度条相互覆盖输出跳来跳去原因tqdm 默认不支持多进程所有进程都在同一行显示进度解决使用 position 参数为每个进程分配不同的行位置# 使用 position 参数解决多进程进度条冲突position os.getpid() % 10 # 根据进程 PID 分配行位置pbar tqdm(totaltotal, desccollection, positionposition, leaveTrue)# 效果每个进程占一行互不干扰# trend_snapshot: 27%|██▋ | 320M/1191M [7:34:3225:04:31, 9643row/s]# competitor: 15%|█▌ | 30M/200M [0:15:001:25:00, 33333row/s]坑 7MongoDB fork 警告现象多进程启动时出现 UserWarning: MongoClient opened before fork原因主进程在 fork 前创建了 MongoDB 连接子进程继承了这个连接影响仅仅是警告不影响功能子进程会重新创建连接解决获取表列表后立即关闭主进程的连接再 fork 子进程# 获取表列表from common.plugins import trend_mongodball_collections db.list_collection_names()# ⭐ 关闭主进程连接避免 fork 继承trend_mongodb.client.close()# fork 子进程此时主进程没有活动连接with Pool(workers) as pool:results list(pool.imap_unordered(export_collection, tables))六、实践效果指标结果数据规模6TB / 230 张表 / 单表最大 11 亿条导出格式ParquetZstd 压缩压缩效果6TB 压缩后约 500GB压缩比约 12:1导出速度多进程并行速度提升 5-10 倍磁盘占用峰值不超过 50GB边导出边上传边删除数据完整性断点续传 强制 Schema零数据丢失七、总结这次数据归档任务的几个关键点· 选对格式Parquet 压缩比高后续数据处理方便· 边导边删本地磁盘不够实时上传释放空间· 流式处理大表分批查询内存占用可控· 断点续传中断后秒级恢复不怕重来· 多进程并行充分利用多核效率翻倍· 异步上传上传不阻塞导出时间利用最大化