
本文介绍三种主流方法——pandas分块读取、dask并行处理和原生csv模块流式解析,帮助你在内存受限条件下安全、高效地处理千万级csv数据,并附带实用代码示例与关键调优建议。
本文介绍三种主流方法——pandas分块读取、dask并行处理和原生csv模块流式解析,帮助你在内存受限条件下安全、高效地处理千万级csv数据,并附带实用代码示例与关键调优建议。
处理超过1000万行的CSV文件时,直接使用 pd.read_csv() 加载全量数据极易触发 MemoryError。根本解决思路是避免一次性加载全部数据到内存,转而采用流式(streaming)、分块(chunked)或延迟计算(lazy evaluation)策略。以下是三种经过生产验证的高效方案:
✅ 方案一:pandas chunksize —— 简单可控,适合顺序聚合类任务
read_csv() 的 chunksize 参数返回一个可迭代的 TextFileReader 对象,每次仅加载指定行数(如 10,000 行)到内存,处理完即释放。适用于求和、计数、过滤、分组统计等可分解操作。
import pandas as pd
filename = "large_data.csv"
chunksize = 10_000
total_sales = 0
valid_records = 0
for chunk in pd.read_csv(filename, chunksize=chunksize,
usecols=["sales", "status"], # 只读必要列,大幅降内存
dtype={"status": "category"}): # 类别类型节省50%+内存
# 过滤有效记录并累加
filtered = chunk[chunk["status"] == "completed"]
total_sales += filtered["sales"].sum()
valid_records += len(filtered)
print(f"总成交额: {total_sales:.2f}, 有效订单数: {valid_records}")⚠️ 注意事项:
- 始终配合
usecols指定所需列;- 对重复字符串列(如状态、地区)优先设为
category类型;- 避免在循环内拼接 DataFrame(如
pd.concat([df, chunk])),这会引发二次内存膨胀;- 若需最终合并结果,建议先保存各 chunk 处理结果至临时列表/磁盘,最后统一处理。
✅ 方案二:Dask DataFrame —— 类pandas语法,支持并行与分布式
Dask 将大文件逻辑切分为多个分区(partitions),按需计算,不将全部数据载入内存。特别适合需要复杂变换(如多列 join、窗口函数、groupby().apply())但仍希望保持 pandas 风格的场景。
Miller (mlr) 是一个命令行工具,用于查询、整形和重新格式化名称索引数据,如 CSV、TSV、JSON 和 JSON Lines。它将 awk、sed、cut、join 和 sort 的功能整合到一个专为结构化数据处理而构建的单一工具中。
import dask.dataframe as dd
# 自动按64MB块分割(约对应数十万行,取决于字段宽度)
ddf = dd.read_csv("large_data.csv",
blocksize="64MB",
dtype={"user_id": "uint32", "score": "float32"}) # 显式指定低精度类型
# 延迟执行:链式操作不立即计算
result = (ddf[ddf["score"] > 80]
.groupby("user_id")
.agg({"score": "mean", "timestamp": "max"})
.compute()) # .compute() 触发实际计算
print(result.head())✅ 优势:自动并行化、支持
.persist()缓存中间结果、可无缝对接dask.delayed或dask.distributed集群。
⚠️ 注意:首次.compute()仍需足够内存容纳结果;避免对小数据集滥用 Dask(启动开销显著)。
✅ 方案三:标准库 csv 模块 + 生成器 —— 极致轻量,适合自定义解析
当只需提取特定字段、做简单条件判断或写入新文件时,原生 csv.reader 或 csv.DictReader 内存占用最低(每行仅保留字符串对象),且无第三方依赖。
import csv
def process_large_csv(filename: str, target_col: str = "amount") -> float:
total = 0.0
with open(filename, "r", newline="", encoding="utf-8") as f:
reader = csv.DictReader(f) # 或 csv.reader(f)
for i, row in enumerate(reader):
try:
val = float(row[target_col])
if val > 0: # 示例业务逻辑
total += val
except (ValueError, KeyError):
continue # 跳过脏数据
if i % 1_000_000 == 0:
print(f"已处理 {i} 行...")
return total
print("正向金额总和:", process_large_csv("large_data.csv"))? 提示:结合
yield构建生成器函数,可进一步解耦读取与处理逻辑,提升复用性与测试性。
? 总结选型建议
| 场景 | 推荐方案 | 理由 |
|---|---|---|
| 快速统计、过滤、简单聚合 | pandas + chunksize |
语法熟悉、生态完善、调试直观 |
| 复杂ETL、多表关联、需水平扩展 | Dask DataFrame | 延迟计算 + 并行调度 + 兼容pandas API |
| 极致内存敏感、纯文本提取/清洗 |
csv 模块 + 手动解析 |
零依赖、最小内存足迹、完全可控 |
无论选择哪种方式,预览数据结构(head -n 100 file.csv)、估算单行内存(sys.getsizeof(row))、监控进程内存(psutil.Process().memory_info().rss) 都是保障稳定性的关键前置步骤。

















