pandas 大文件的分块处理:chunksize 迭代器与内存峰值压制实战

在数据分析的日常工作中,我们经常会从生产数仓、日志归档存储或外部客户那里接收到一个体积庞大的单体 CSV 或 TXT 数据文件(例如一个 15 GB 的全量用户行为日志)。
很多刚入行的数据分析师在拿到这个文件后,第一反应是在 Jupyter Notebook 里敲下一行:df = pd.read_csv("huge_log_15gb.csv")。
紧接着,你的电脑屏幕会发生熟悉的静止:
- 鼠标指针变成彩虹等待圈;
- 系统的风扇发出刺耳的轰鸣;
- 操作系统内存监控中,Python 进程的内存瞬间吃满 16GB / 32GB,触发交换分区(Swap)剧烈抖动;
- 最终终端直接弹出一行冷酷的提示:MemoryError: Unable to allocate 15.0 GiB for an array…,内核直接被操作系统 OOM Killer 强制杀死。
面对超过物理内存限制的超大文件,难道只能去搭建沉重的分布式 Spark 集群吗?
完全不需要!利用 Pandas 原生提供的 chunksize 分块流式迭代器(Chunking Iterator) 与 轻量级规约聚合技术,单台普通笔记本电脑就能在极低内存占用(< 300 MB)下,行云流水般完成数十 GB 级别大文件的数据清洗与统计分析。
流式分块处理的物理机制:流水线视窗(Sliding Pipeline)
分块读取的本质,是将一次性吞下大象的“暴食模式”,重构为小口吞咽的“流水线模式”:
[ 磁盘上的 15 GB 超大 CSV 文件 (包含 5000 万行) ]
│
▼ (每次仅读取 chunksize = 100,000 行到内存)
┌─────────────────────────────────────────────────────────────┐
│ 内存中的微型 DataFrame (仅占用约 80 MB 内存) │
│ 1. 过滤脏数据 (去除无效状态与空值) │
│ 2. 原地局部聚合: df_chunk.groupby('store_id')['amount'].sum()│
└─────────────────────────────────────────────────────────────┘
│
▼ (仅将微小的局部聚合结果累加到全局汇总器)
[ 全局精简指标字典 / 最终结果落地文件 ] (内存常驻 < 20 MB)
无论磁盘上的原始文件是 10 GB 还是 100 GB,内存峰值永远被死死压制在单批次 chunk 所占用的几十兆字节内!
实战演练一:超大文件多维 GroupBy 聚合统计
假设我们需要处理一个 15 GB 的交易明细日志,统计:每个省份(province)在各渠道(channel)的累计有效订单总金额与去重用户数。
import pandas as pd
import numpy as np
def process_huge_csv_by_chunks(file_path: str, chunk_size: int = 200_000):
"""
分块流式聚合超大 CSV 文件
"""
# 1. 显式指定需要的列与紧凑数据类型 (进一步降低每块内存)
use_cols = ['province', 'channel', 'user_id', 'pay_amount', 'status']
dtypes = {
'province': 'category',
'channel': 'category',
'user_id': 'int64',
'pay_amount': 'float32',
'status': 'category'
}
# 2. 开启 chunksize 模式,此时 reader 是一个 TextFileReader 迭代器
reader = pd.read_csv(
file_path,
usecols=use_cols,
dtype=dtypes,
chunksize=chunk_size
)
# 全局指标中间累加器
intermediate_results = []
total_processed_rows = 0
for i, chunk in enumerate(reader):
total_processed_rows += len(chunk)
# 步骤 A:块内快速过滤 (仅保留有效支付数据)
valid_chunk = chunk[chunk['status'] == 'PAID']
if valid_chunk.empty:
continue
# 步骤 B:块内局部预聚合 (减少中间行数)
# 注意:此处将 user_id 收集为集合 set,用于后续全局精准去重
local_agg = (
valid_chunk.groupby(['province', 'channel'], observed=True)
.agg(
local_amount=('pay_amount', 'sum'),
local_users=('user_id', lambda x: set(x.unique()))
)
.reset_index()
)
intermediate_results.append(local_agg)
if (i + 1) % 10 == 0:
print(f"已流式处理 {total_processed_rows:,} 行数据,当前内存稳定处于安全区…")
# 3. 最终全局规约合并 (Reduce Phase)
combined_df = pd.concat(intermediate_results, ignore_index=True)
final_summary = (
combined_df.groupby(['province', 'channel'], observed=True)
.agg(
total_pay_amount=('local_amount', 'sum'),
# 将各块的 user 集合求并集,计算全局精准 UV 去重!
unique_buyer_count=('local_users', lambda sets: len(set().union(*sets)))
)
.reset_index()
)
return final_summary
实战演练二:分块清洗与流式写入新文件(ETL Pipeline)
如果任务是“清洗一个 20GB 的脏文件,并输出为一个精简干净的新 Parquet 文件”,可以使用 pyarrow 实现分块流式持久化:
import pandas as pd
import pyarrow as pa
import pyarrow.parquet as pq
def stream_clean_and_export_parquet(input_csv: str, output_parquet: str):
reader = pd.read_csv(input_csv, chunksize=100_000)
parquet_writer = None
for chunk in reader:
# 清洗流水线
clean_chunk = (
chunk.dropna(subset=['order_id'])
.assign(net_amount=lambda x: x['amount'] – x['fee'])
)
# 转换为 PyArrow Table 并流式写入 Parquet 格式
table = pa.Table.from_pandas(clean_chunk)
if parquet_writer is None:
# 首次初始化写入器并固化 Schema
parquet_writer = pq.ParquetWriter(output_parquet, table.schema, compression='snappy')
parquet_writer.write_table(table)
if parquet_writer:
parquet_writer.close()
print("超大文件已成功流式重构为高性能 Parquet 格式!")
生产落地的三条核心性能护栏
- 设置太小(如 chunksize=1000):Python 循环开销占比过大,整体变慢;
- 设置太大(如 chunksize=5,000,000):单块内存可能接近物理极限;
- 黄金折中值:通常在 100,000 ~ 500,000 行 之间,单块内存约 50MB ~ 150MB,CPU 解析效率与内存安全达到最佳平衡。

