本文介绍一种内存友好的方式,利用 pymongoarrow 直接将 mongodb 数据流式转换为 arrow 表,再分块加载至 duckdb,规避 jsonl 中间文件与全量内存加载问题。
本文介绍一种内存友好的方式,利用 pymongoarrow 直接将 mongodb 数据流式转换为 arrow 表,再分块加载至 duckdb,规避 jsonl 中间文件与全量内存加载问题。
在 Python 中将大型 MongoDB 集合导入 DuckDB 时,传统做法(如先导出为 JSONL 文件再 read_json_auto)易引发高内存占用和 I/O 瓶颈。更优解是绕过磁盘序列化,采用 零拷贝、列式、流式 的 Arrow 生态链路:MongoDB → PyMongoArrow → DuckDB。
✅ 推荐方案:PyMongoArrow + DuckDB 分块加载
PyMongoArrow 是 MongoDB 官方支持的 Arrow 集成库,可直接从 pymongo cursor 获取 pyarrow.Table 或 pyarrow.RecordBatchReader,天然适配 DuckDB 的高性能 Arrow 接口。
安装依赖
pip install pymongoarrow duckdb pyarrow
核心实现(支持大集合、低内存、自动类型推断)
import duckdb
from pymongoarrow.api import Schema, find_arrow_all
from pymongo import MongoClient
# 1. 连接 MongoDB(确保已启用 Arrow 支持,MongoDB 6.0+ & 合适驱动版本)
client = MongoClient("mongodb://localhost:27017")
collection = client["mydb"]["mycollection"]
# 2. 定义可选 Schema(提升性能与类型稳定性;若省略,PyMongoArrow 自动推断)
# schema = Schema({"_id": pa.string(), "name": pa.string(), "score": pa.float64()})
# 3. 直接读取为 Arrow Table(适合中等规模数据,全量入内存)
# table = find_arrow_all(collection, {}, schema=schema)
# ✅ 更佳实践:使用 RecordBatchReader 流式读取(推荐!内存恒定)
reader = collection.find_arrow_all({}, batch_size=10000) # 按批次拉取(单位:行)
# 4. 创建 DuckDB 连接并建表(首次写入自动建表)
conn = duckdb.connect("analytics.duckdb")
# 初始化表结构(仅执行一次,基于首个 batch 推断)
first_batch = next(reader)
conn.register("temp_batch", first_batch)
conn.execute("""
CREATE OR REPLACE TABLE mongo_table AS
SELECT * FROM temp_batch
""")
# 5. 追加后续批次(避免重复注册/建表)
for batch in reader:
conn.register("temp_batch", batch)
conn.execute("INSERT INTO mongo_table SELECT * FROM temp_batch")
print(f"✅ 成功导入 {conn.execute('SELECT COUNT(*) FROM mongo_table').fetchone()[0]} 条记录")
⚠️ 关键注意事项
- MongoDB 版本要求:服务端需 ≥ 6.0,且部署支持 Arrow 查询(默认开启;如禁用需配置 --enableArrow 启动参数);
- 驱动兼容性:pymongoarrow 依赖 pymongo >= 4.3 和 pyarrow >= 12.0;
- Schema 控制:显式定义 Schema 可避免类型推断偏差(如 ObjectId → string、datetime 保留时区),提升 DuckDB 查询效率;
- 内存安全:find_arrow_all(..., batch_size=N) 返回 RecordBatchReader,每次仅加载一个批次到内存,彻底规避 OOM;
- 性能提示:DuckDB 原生支持 Arrow 表插入,无需序列化/反序列化,吞吐量显著高于 JSONL 方案。
? 总结
放弃“Mongo → JSONL → DuckDB”这一磁盘+内存双重开销路径,转向“Mongo → Arrow(流式)→ DuckDB”架构,不仅降低内存峰值(常驻内存 ≈ 单批次大小),还提升整体 ETL 效率与类型可靠性。对于 TB 级集合,配合 MongoDB 聚合管道预过滤(如 {"$match": {"status": "active"}}),可进一步优化数据传输量。











