
本文介绍三种主流方案:pandas分块读取、dask延迟计算和原生csv模块流式处理,涵盖代码示例、性能权衡与关键调优技巧,助你稳定处理1000万+行csv数据。
本文介绍三种主流方案:pandas分块读取、dask延迟计算和原生csv模块流式处理,涵盖代码示例、性能权衡与关键调优技巧,助你稳定处理1000万+行csv数据。
处理超大CSV文件(如1000万行以上)时,直接使用 pd.read_csv() 加载全部数据极易触发MemoryError。根本解决思路是避免一次性加载全量数据到内存,转而采用流式读取、延迟计算或按需解析策略。以下是经过实践验证的三种高效方案:
✅ 方案一:pandas chunksize —— 简单可控,适合聚合/过滤类任务
chunksize 参数将CSV按行分批加载为DataFrame迭代器,每批仅占用对应内存,适合统计汇总、条件筛选、批量写入等场景。
import pandas as pd
filename = "large_data.csv"
chunksize = 50_000 # 推荐 10k–100k 行/块,过大仍可能OOM;过小则I/O开销上升
# 示例:计算某数值列的全局总和
total_sales = 0
for chunk in pd.read_csv(filename, chunksize=chunksize, usecols=["sales"]): # 显式指定列,大幅减内存
total_sales += chunk["sales"].sum()
print(f"总销售额: {total_sales:,}")
关键优化点:
- 使用
usecols仅读取必需列(可减少30%–70%内存); - 设置
dtype(如{"id": "uint32", "price": "float32"})避免默认object或float64浪费; - 对字符串列启用
dtype={"category"}或low_memory=False防止类型推断错误; - 处理完每块后显式
del chunk并调用gc.collect()(Python 3.12+ 可省略)。
✅ 方案二:Dask DataFrame —— 类pandas语法,支持并行与分布式
Dask将大数据集切分为多个分区(partitions),所有操作均延迟执行(lazy evaluation),.compute() 时才真正计算,天然支持多核加速。
Miller (mlr) 是一个命令行工具,用于查询、整形和重新格式化名称索引数据,如 CSV、TSV、JSON 和 JSON Lines。它将 awk、sed、cut、join 和 sort 的功能整合到一个专为结构化数据处理而构建的单一工具中。
import dask.dataframe as dd
# 自动按64MB块分割(推荐值),也可设 nrows=1e6 指定行数
ddf = dd.read_csv(
"large_data.csv",
blocksize="64MB",
dtype={"user_id": "uint32", "amount": "float32"},
assume_missing=True # 处理空值更鲁棒
)
# 链式操作不立即执行
result = ddf[ddf["amount"] > 100]["amount"].mean().compute() # 触发实际计算
print(f"高消费用户平均金额: {result:.2f}")
适用场景:
✔️ 需要复杂链式操作(filter → groupby → agg → join);
✔️ 单机多核CPU充足,追求比pandas chunk更快的吞吐;
⚠️ 注意:首次 .compute() 仍会占用峰值内存,但远低于全量加载。
✅ 方案三:原生 csv 模块 + 生成器 —— 极致内存控制,适合逐行逻辑
当只需提取特定字段、做简单转换或写入新文件时,csv.reader 或 csv.DictReader 内存占用最低(每行仅数百字节)。
import csv
def process_large_csv(filename):
with open(filename, "r", newline="", encoding="utf-8") as f:
reader = csv.DictReader(f) # 或 csv.reader(f)
for i, row in enumerate(reader, 1):
# 示例:提取ID和状态,跳过无效行
if not row.get("id") or row.get("status") != "active":
continue
yield int(row["id"]), row["name"]
# 流式处理,内存恒定 ~1KB
for user_id, name in process_large_csv("large_data.csv"):
if user_id % 100000 == 0:
print(f"已处理 {user_id} 条活跃用户")
优势:
- 内存占用与文件大小无关,仅取决于单行长度;
- 可无缝集成
concurrent.futures实现多线程解析(注意GIL限制); - 配合
yield实现管道化(pipeline)处理,如:read → validate → transform → write。
? 终极建议与避坑指南
-
优先尝试
pandas + chunksize:学习成本低、生态完善,90%场景已足够; -
避免
pandas.read_csv(..., iterator=True):已被弃用,统一用chunksize; -
警惕
dask的“假并行”:单机下若I/O瓶颈明显,多线程反而拖慢;建议先用dd.read_csv(..., sample_nrows=10000)探查schema; -
预处理优于运行时处理:对超大文件,考虑用
awk/sed/csvkit提前清洗(如删空行、抽样建索引); -
终极手段:数据库导入:
sqlite3(轻量)、duckdb(分析快)或postgresql COPY,适合反复查询场景。
选择方案的核心依据是:你的计算模式(单次聚合?多次交互?实时流?)和硬件约束(内存 vs CPU vs 磁盘IO)。没有银弹,但理解每种机制的内存模型,就能精准控压、稳定交付。










