直接to_parquet()会oom或极慢,因pyarrow默认全量加载内存且不启用分块写入,加上类型推断不准和压缩低效;应改用分批write_table+类型降级+zstd压缩。

为什么直接to_parquet()会OOM或极慢?
千万级DataFrame(比如 2000 万行 × 50 列)直接调用 df.to_parquet("out.parq") 很容易触发内存暴涨甚至崩溃——Pandas 默认用 pyarrow 引擎,但不启用分块写入时,会把整个 DataFrame 加载进内存再序列化,中间还可能复制数据(尤其含 object 类型列)。另外,未指定压缩或 schema 时,PyArrow 可能默认用 snappy(快但压缩率低)或不做类型推断,导致文件体积翻倍。
- 确认当前引擎:
pd.options.io.parquet.engine,优先设为"pyarrow"("fastparquet"对复杂类型支持弱且并发写不稳) - 强制分块写入:不用一次性全量导出,改用
pd.concat(...).to_parquet()是陷阱;正确做法是用PartitionedDataSet思路或手动分批 - 提前清理无用列、转换低精度类型(如
int64→int32,object→category),能减少 30–60% 内存占用
如何用分块 + 类型优化安全导出?
核心不是“怎么调用”,而是“怎么切、怎么转、怎么压”。下面这段实操代码覆盖真实瓶颈点:
import pandas as pd
import pyarrow as pa
<h1>假设 df 是你的原始 DataFrame</h1><p>df = df.copy() # 避免 SettingWithCopyWarning
df["date"] = pd.to_datetime(df["date"]) # 显式转 datetime,避免 PyArrow 推断为 string
df["category_col"] = df["category_col"].astype("category") # object → category 省内存
df["id"] = df["id"].astype("int32") # 检查值域后降精度</p><h1>分块写入:每 50 万行一批,用 pyarrow 的 write_table 手动控制</h1><p>batch_size = 500_000
for i in range(0, len(df), batch_size):
batch = df.iloc[i:i+batch_size]
table = pa.Table.from_pandas(batch)
pa.parquet.write<em>table(
table,
f"output/part</em>{i//batch_size:04d}.parquet",
compression="zstd", # 比 snappy 压缩率高 2–3×,CPU 开销可控
use_dictionary=True, # 对 category / 重复字符串启用字典编码
version="2.6", # 兼容 Spark 3.0+ 和 DuckDB
)</p>
注意:pa.parquet.write_table() 比 df.to_parquet() 更底层、更可控;compression="zstd" 需要 pip install zstandard;version="2.6" 是目前最稳妥的跨系统兼容版本。
分区(partition)写入是否必要?
如果你后续主要按某列(如 date、region)过滤查询,分区能极大加速读取——但写入时会生成大量小文件,反而拖慢导出和元数据管理。千万级数据下,只在满足以下条件时才开分区:
- 分区字段取值离散度低(如
date.dt.date只有 365 个值,而非timestamp每行都不同) - 你确定下游工具(DuckDB / Polars / Spark)会利用分区剪枝
- 单个分区数据量 ≥ 100MB(否则小文件过多,HDFS/S3 列表操作变瓶颈)
开启方式不是传 partition_cols 给 to_parquet(),而是用 pyarrow.dataset.write_dataset():
from pyarrow import dataset as ds
table = pa.Table.from_pandas(df)
ds.write_dataset(
table,
"output_partitioned",
format="parquet",
partitioning=ds.partitioning(pa.schema([("date", pa.date32())]), flavor="hive"),
use_threads=True,
)
注意:flavor="hive" 生成 date=2024-01-01/ 这类目录,Spark/DuckDB 能自动识别;但 Pandas 自带的 read_parquet() 不支持直接读 Hive 分区目录,得用 ds.dataset().to_table() 或 DuckDB。
导出后校验与读取性能关键点
文件写完不代表可用。常见失效场景:schema 不一致(比如某批里 float64 列出现 NaN 后被推成 float32)、压缩损坏、分区路径错位。必须做三件事:
- 用
pyarrow.parquet.read_metadata("part_0000.parquet")检查各分片 schema 是否完全一致 - 用
duckdb.query("SELECT count(*) FROM 'output/*.parquet'").df()快速验证总行数(比 Pandasread_parquet快 5–10×) - 读取时禁用全局
use_threads=False(默认 True 反而因锁竞争变慢),并指定列子集:pd.read_parquet("output/", columns=["id", "date"])
最后提醒:Parquet 不是万能“加速器”。如果原始 DataFrame 已经是内存瓶颈,先解决数据建模问题(比如是否真需要导出全部列?能否预聚合?)比调参更重要。文件大小和读取速度之间永远存在权衡,zstd 压得越狠,解压 CPU 越高——在 I/O 密集型场景(如 S3)值得,但在本地 SSD 上可能 snappy 更合适。
Python免费学习笔记(深入):立即使用
在学习笔记中,你将探索 Python 的核心概念和高级技巧!











