
本文介绍三种主流方案:pandas分块读取、dask分布式处理和原生csv模块流式解析,涵盖代码示例、性能对比与关键注意事项,助你稳定处理1000万+行csv数据。
本文介绍三种主流方案:pandas分块读取、dask分布式处理和原生csv模块流式解析,涵盖代码示例、性能对比与关键注意事项,助你稳定处理1000万+行csv数据。
处理超大CSV文件(如1000万行以上)时,直接使用 pd.read_csv(filename) 极易触发 MemoryError。根本原因在于pandas默认将整个文件加载进内存并构建完整DataFrame。解决思路是避免一次性全量加载,转而采用流式、分块或延迟计算策略。以下是三种经过生产验证的高效方案:
✅ 方案一:pandas chunksize —— 简单可控,适合聚合/过滤类任务
chunksize 参数让 read_csv() 返回一个可迭代的 TextFileReader 对象,每次仅加载指定行数(如10,000行)到内存,处理完即释放。适用于统计汇总、条件筛选、逐块写入数据库等场景。
import pandas as pd
filename = "large_data.csv"
chunksize = 10_000
total_sales = 0.0
valid_records = []
for chunk in pd.read_csv(filename, chunksize=chunksize,
usecols=["order_id", "amount", "status"], # 只读必要列,大幅降内存
dtype={"order_id": "category", "status": "category"}): # 类型预设节省内存
# 示例1:累加数值列
total_sales += chunk["amount"].sum()
# 示例2:筛选有效订单并暂存(谨慎!避免累积过多)
valid_chunk = chunk[chunk["status"] == "completed"]
valid_records.append(valid_chunk)
print(f"总销售额: {total_sales:.2f}")
# 若需合并所有有效记录,用 pd.concat(valid_records, ignore_index=True),但注意内存峰值
⚠️ 关键提示:
- 始终配合
usecols(指定列)和dtype(显式类型)减少内存占用; - 避免在循环中无节制追加DataFrame(如
results.append(chunk)),易导致内存持续增长; - 处理逻辑尽量向量化(如
.sum(),.groupby().agg()),避免对每行调用Python函数。
✅ 方案二:Dask DataFrame —— 类pandas语法,支持并行与磁盘计算
Dask将大数据集划分为多个分区(partitions),按需加载、惰性求值,并支持多线程/多进程加速。特别适合需要复杂变换(如多列join、窗口函数)且无法全量驻留内存的场景。
Miller (mlr) 是一个命令行工具,用于查询、整形和重新格式化名称索引数据,如 CSV、TSV、JSON 和 JSON Lines。它将 awk、sed、cut、join 和 sort 的功能整合到一个专为结构化数据处理而构建的单一工具中。
import dask.dataframe as dd
# 自动按64MB分块(推荐),也可指定 nrows 或 blocksize="128MB"
ddf = dd.read_csv("large_data.csv",
blocksize="64MB",
dtype={"user_id": "string", "score": "float32"}) # 支持部分类型推断优化
# 惰性构建计算图(不立即执行)
result = ddf.groupby("category")["score"].mean().compute() # .compute() 触发实际计算
print(result.head())
✅ 优势:语法与pandas高度兼容;自动并行化;支持 .persist() 缓存中间结果;可无缝对接Xarray、Scikit-learn等生态。
⚠️ 注意:首次.compute()有启动开销;小文件(blocksize(通常32–128MB)。
✅ 方案三:原生 csv 模块 + 生成器 —— 内存极致精简,适合逐行解析
当只需提取特定字段、做简单校验或写入新格式(如JSONL、数据库INSERT)时,绕过pandas开销,用标准库流式处理最轻量:
import csv
def process_large_csv(filename: str, target_column: str = "email"):
with open(filename, "r", encoding="utf-8") as f:
reader = csv.DictReader(f) # 或 csv.reader(f) + 手动映射列名
for i, row in enumerate(reader):
if i % 100_000 == 0:
print(f"Processed {i} rows...")
# 示例:提取邮箱并简单清洗
email = row.get(target_column, "").strip().lower()
if "@" in email:
yield email # 使用生成器,内存恒定O(1)
# 流式消费(不缓存全部结果)
valid_emails = list(process_large_csv("users.csv")) # 仅当需全部结果时才list()
# 或直接写入文件:
# with open("clean_emails.txt", "w") as out:
# for email in process_large_csv("users.csv"):
# out.write(email + "\n")
✅ 优势:内存占用最低(仅单行);启动最快;无第三方依赖。
⚠️ 局限:无内置缺失值处理、类型转换、向量化运算;需手动实现业务逻辑,开发成本略高。
? 总结建议
| 场景 | 推荐方案 | 关键操作 |
|---|---|---|
| 快速统计、分组聚合、简单过滤 | pandas.read_csv(chunksize=...) |
务必用 usecols + dtype
|
| 复杂ETL、多表关联、需复用中间结果 | Dask DataFrame | 设置合理 blocksize,善用 .persist()
|
| 提取关键字段、数据清洗、导入数据库 | 原生 csv 模块 + 生成器 |
用 yield 流式产出,避免 list() 全量加载 |
最后提醒:无论选用哪种方式,预先用 head -n 1000 large_data.csv > sample.csv 抽样检查结构,确认分隔符、编码、空值标记等,能避免90%的解析失败。










