MongoDB 数据归档实战

发布时间:2026/8/26 16:26:35

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 压缩比高后续数据处理方便· 边导边删本地磁盘不够实时上传释放空间· 流式处理大表分批查询内存占用可控· 断点续传中断后秒级恢复不怕重来· 多进程并行充分利用多核效率翻倍· 异步上传上传不阻塞导出时间利用最大化

相关新闻

今日热榜API:5分钟搭建聚合50+平台热榜的接口服务

今日热榜API:5分钟搭建聚合50+平台热榜的接口服务

2026/8/26 16:16:35

今日热榜API:5分钟搭建聚合50平台热榜的接口服务 【免费下载链接】DailyHotApi 🔥 今日热榜 API,一个聚合热门数据的 API 接口,支持 RSS 模式 及 Vercel 部署 | 前端页面:https://github.com/imsyy/DailyHot 项目地址…

E2E_端到端的ViT_Pytorch实现

E2E_端到端的ViT_Pytorch实现

2026/8/26 16:16:35

✅ 一、PyTorch 代码框架(含对比损失 URDF 正则) 目标:从多视角视频预测 12 维关节角 角速度,并施加 URDF 结构约束。1. 安装依赖 pip install torch torchvision timm einops pytorch3d # pytorch3d 用于可微分 FK&#xff08…

Geoserver2.27.3结合GeoWebCache发布arcgis切片(WMTS服务)

Geoserver2.27.3结合GeoWebCache发布arcgis切片(WMTS服务)

2026/8/26 16:16:35

1.背景说明 原来用的还是Geoserver2.21,漏洞实在太多了,升级到最新的2.27.3(2025年12月18日获取到的最新版)。需要注意两点 (1)java版本从8升级到17.(可以不改服务器java8配置,只让…

多模型对比写同一个函数:Gemini、GPT、DeepSeek 谁更靠谱?开发者横向实测

多模型对比写同一个函数:Gemini、GPT、DeepSeek 谁更靠谱?开发者横向实测

2026/8/26 17:26:38

Q:同一个函数分别交给 Gemini、GPT、DeepSeek 来写,哪个更靠谱? A:如果只看“能不能快速给出一版可运行代码”,三者都能完成基础任务;但如果放到真实开发里,差别就出来了。我最近拿同一个函数做…

ArcGIS JS 基础教程(29):CSVLayer 表格点位图层

ArcGIS JS 基础教程(29):CSVLayer 表格点位图层

2026/8/26 17:26:38

ArcGIS JS 基础教程(29):CSVLayer 表格点位图层零、写在前面一、功能介绍二、功能实现2.1 创建并配置(3D)2.2 让点位「浮」起来(高程模式)2.3 自定义字段类型(numericFields 等&…

Java 第k个最小元素(K’th Smallest Element)

Java 第k个最小元素(K’th Smallest Element)

2026/8/26 17:26:38

目录 【朴素方法】使用排序——时间复杂度为 O(n log(n)),空间复杂度为 O(1) 【预期方法】使用最大堆 - 时间复杂度为 O(n * log(k)),空间复杂度为 O(k) 【替代方案 1】使用快速选择 【替代方案 2】使用计数排序 如果您喜欢此文章,请收藏…

哈森股份(603958)深度研究报告

哈森股份(603958)深度研究报告

2026/8/26 17:26:38

一、投资要点1.1 核心结论哈森股份(603958.SH)是国内中高端女鞋领域的代表性企业之一,旗下拥有哈森、卡迪娜、卡文等自有品牌,并代理多个国际品牌。公司于2016年登陆上交所主板,是A股市场少数以中高端女鞋为主营的上市…

Python 第k个最小元素(K’th Smallest Element)

Python 第k个最小元素(K’th Smallest Element)

2026/8/26 17:26:38

目录 【朴素方法】使用排序——时间复杂度为 O(n log(n)),空间复杂度为 O(1) 【预期方法】使用最大堆 - 时间复杂度为 O(n * log(k)),空间复杂度为 O(k) 【替代方案 1】使用快速选择 【替代方案 2】使用计数排序 如果您喜欢此文章,请收藏…

UVa 10240 The n-Dimensional Cities

UVa 10240 The n-Dimensional Cities

2026/8/26 17:16:37

题目描述 在 nnn 维空间中,最多可以有 (n1)(n1)(n1) 个点两两等距。Talisman\texttt{Talisman}Talisman 国的 nnn 维生物建造了 (n1)(n1)(n1) 座城市,这些城市两两等距。Talisman\texttt{Talisman}Talisman 的道路按照如下规则修建: 111. 每一…

[光学原理与应用-521]:对光的错误理解与纠偏

[光学原理与应用-521]:对光的错误理解与纠偏

2026/8/26 1:50:39

首先光是一种能量的载体和形态,宏观上观察到的光是由无数个微观的光量子组成的,每个光子在产生的瞬间,其在真空的空间中以确定不变的速度沿着一个初始的方向一直向前,在微观层面,每个光量子的运动轨迹是以波函数所展现…

SIP通话转接原理与REFER方法实战解析

SIP通话转接原理与REFER方法实战解析

2026/8/26 1:49:16

1. 通话转接不是“挂断再拨号”,而是SIP会话的动态重定向你有没有遇到过这样的场景:客服坐席A正在和客户通电话,突然需要把这通对话无缝转给专家坐席B,客户完全感知不到中间的断连——既没听到忙音,也没被要求重新拨号…

Kolla-ansible单节点OpenStack部署实战:从环境准备到排坑指南

Kolla-ansible单节点OpenStack部署实战:从环境准备到排坑指南

2026/8/24 21:16:09

1. 为什么选择Kolla-ansible来部署单节点OpenStack?如果你正在寻找一种能把OpenStack从“概念”快速变成“可用的实验环境”的方法,那么Kolla-ansible几乎是当前最主流、最省心的选择。我见过太多人卡在手动编译依赖、配置服务、处理版本冲突的泥潭里&am…

Python random 模块常用函数详解:从入门到实战

Python random 模块常用函数详解:从入门到实战

2026/8/26 0:05:45

目录 1. 引言2. 准备工作3. 基础随机函数4. 序列相关函数5. 随机种子与复现6. 实战案例7. 注意事项8. 常见问题与排查9. 总结 1. 引言 摘要: 本文系统介绍 Python 标准库 random 模块中最常用的随机数生成函数。内容涵盖基础随机函数(random()、unifor…

Hermes接入团队协作后,我推翻了三个效率假设

Hermes接入团队协作后,我推翻了三个效率假设

2026/8/26 0:05:45

聊《Hermes真能提效吗?先看流程里最慢的那一步》之前,先说一句实在的:别急着背概念,先看它在真实项目里到底解决什么问题。摘要团队把 Hermes 接进项目三个月后,交付速度没有提升反而慢了。复盘后发现,最先…

免费AI大模型调教指南:打造专属网文写作助手

免费AI大模型调教指南:打造专属网文写作助手

2026/8/26 0:05:45

1. 先搞清楚“AI小说扩展模式”到底能帮你做什么如果你是一个刚开始写网文、或者卡在L3级别以下的作者,最头疼的可能是情节推进不下去、人物对话干瘪,或者世界观设定不够丰满。自己对着空白文档硬憋,效率很低。这时候,一个能理解你…

摆脱论文困扰!盘点2026年全网爆红的的AI论文写作工具

摆脱论文困扰!盘点2026年全网爆红的的AI论文写作工具

2026/8/22 2:02:26

一天写完毕业论文在2026年已不再是天方夜谭。2026年最炸裂、实测能大幅提速的AI论文写作工具,覆盖选题构思、文献整理、内容生成、格式排版等核心场景,真正帮你高效搞定论文难题。 一、全流程王者:一站式搞定论文全链路(一天定稿首…

导师推荐!2026最新AI论文工具测评与实用推荐

导师推荐!2026最新AI论文工具测评与实用推荐

2026/8/22 4:13:47

2026年真正好用的AI论文工具,核心看生成的论文质量、低AI味、格式正确、学术适配四大指标。综合实测,千笔AI、ThouPen、豆包、DeepSeek、Grammarly 是当前最值得推荐的梯队,覆盖从免费到付费、从中文到英文、从文科到理工的全场景需求。 一、…

告别游戏崩溃:XCOM 2模组管理器的智能革命

告别游戏崩溃:XCOM 2模组管理器的智能革命

2026/8/22 1:32:34

告别游戏崩溃:XCOM 2模组管理器的智能革命 【免费下载链接】xcom2-launcher The Alternative Mod Launcher (AML) is a replacement for the default game launchers from XCOM 2 and XCOM Chimera Squad. 项目地址: https://gitcode.com/gh_mirrors/xc/xcom2-lau…