
用pandas读取CSV时,最自然的写法是pd.read_csv()。它会返回一张完整的DataFrame,后续筛选和分组都很方便。但当文件逐渐接近机器可用内存时,问题不只在文件本身有多大:解析字符串、创建索引、生成中间列和执行聚合都可能继续占用内存。
如果任务只是按地区计数、求和或计算退款数,没有必要让全部订单同时留在内存。我们可以把CSV切成若干行块,每次读取一块、产生一份很小的部分结果,最后再把这些部分结果相加。这种“读取一块—计算一块—合并结果”的流程,就是本文要实现的分块聚合。
案例会处理24万行模拟订单,并把分块结果与一次性全量读取逐项比较。这里用24万行是为了让代码在普通电脑上快速复现,不代表chunksize只能用于这个规模;实际应用中,块大小应根据列数、数据类型、可用内存和单块处理逻辑共同确定。
输入是一份包含订单编号、日期、地区、渠道、金额和退款标记的CSV。目标是按地区统计订单数、退款数和净销售额,同时避免把全部明细一次性载入DataFrame。
流程需要满足三个判断标准:所有行都被处理且没有重复;分块聚合与全量基准逐项一致;最大单块DataFrame的估算占用明显低于全量DataFrame。这里的全量读取只用于教程中的正确性和内存对照,真实内存受限任务不需要保留这一步。
给pd.read_csv()传入chunksize后,返回值不再是一张完整DataFrame,而是一个可以迭代的TextFileReader。每次循环只取得指定行数的一块数据。只要当前块、部分结果和少量运行状态能够放进内存,逻辑数据集就可以大于单次可用内存。
分块聚合可以理解为两个阶段。第一阶段在每块内部按地区计算订单数、退款数和净销售额;第二阶段把八份部分结果按地区再次求和。计数和求和具有可合并性,所以这种两阶段结果与对全部明细直接聚合相同。
并非所有统计量都能这样处理。把各块中位数再取平均,通常不等于总体中位数;标准差也不能只靠各块标准差直接相加。遇到分位数、全局排序、复杂窗口或需要跨块协调的任务,应推导正确的合并公式、使用近似算法,或改用支持外存计算的数据库和分析引擎。
本文使用固定随机种子生成24万条模拟电商订单,数据不代表真实业务分布。金额用整数“分”保存,退款行在计算净额时取负数。这样做既贴合金额存储的常见实践,也能避免浮点数分块求和产生微小舍入差。
读取阶段只选择region、amount_cents和is_refund三列,因为其余字段不参与本次聚合。每块为3万行,理论上得到8块。代码同时记录每块DataFrame的深度内存估算,以及处理完每一块后的地区累计净销售额,分别用于内存对照和处理过程可视化。
完整代码分四段连续运行。示例会重建/tmp/chunked_csv_demo专用目录;如果在自己的项目中修改WORK_DIR,应先确认它没有指向真实数据目录。
import osfrom pathlib import Path# 把Matplotlib缓存放入可写目录,确保后台绘图不触发权限警告os.environ[”MPLCONFIGDIR”] = ”/tmp/wechat_python_chunk_mpl_20260811”import shutilimport matplotlibmatplotlib.use(”Agg”)import matplotlib.pyplot as pltimport numpy as npimport pandas as pdfrom pandas.testing import assert_frame_equal# 固定随机种子,保证模拟订单与所有统计结果可以复现PYTHON_SEED = 20260811ROW_COUNT = 240_000CHUNK_SIZE = 30_000rng = np.random.default_rng(PYTHON_SEED)# 重建专用演示目录,避免旧文件影响本次运行WORK_DIR = Path(”/tmp/chunked_csv_demo”)if WORK_DIR.exists():shutil.rmtree(WORK_DIR)WORK_DIR.mkdir(parents=True)CSV_PATH = WORK_DIR / ”online_orders.csv”# 生成固定种子的模拟电商订单,不代表真实业务分布regions = [”East”, ”North”, ”South”, ”West”]channels = [”App”, ”MiniProgram”, ”Web”]orders = pd.DataFrame({”order_id”: np.arange(1, ROW_COUNT + 1),”order_date”: pd.Timestamp(”2026-01-01”)+ pd.to_timedelta(rng.integers(0, 180, ROW_COUNT), unit=”D”),”region”: rng.choice(regions, ROW_COUNT, p=[0.28, 0.22, 0.27, 0.23]),”channel”: rng.choice(channels, ROW_COUNT, p=[0.52, 0.31, 0.17]),”amount_cents”: rng.gamma(3.5, 4200, ROW_COUNT).round().astype(”int64”),”is_refund”: rng.random(ROW_COUNT) < 0.055,})orders.to_csv(CSV_PATH, index=False)del ordersprint(f”数据来源:固定随机种子的模拟订单(seed={PYTHON_SEED})”)print(f”CSV行数:{ROW_COUNT:,}”)print(f”CSV大小:{CSV_PATH.stat().st_size / 1024**2:.2f} MB”)print(f”分块设置:每块 {CHUNK_SIZE:,} 行”)
输出:
数据来源:固定随机种子的模拟订单(seed=20260811)CSV行数:240,000CSV大小:9.41 MB分块设置:每块 30,000 行
磁盘上的CSV为9.41 MB,但这不能直接代表读取后的内存占用。文本解析成字符串、整数、布尔值和索引后,DataFrame有自己的内存结构;后面的对照会直接计算DataFrame各列的占用。
# 只读取聚合所需列,避免无关字段进入每一个内存块use_columns = [”region”, ”amount_cents”, ”is_refund”]dtype_map = {”region”: ”string”, ”amount_cents”: ”int64”, ”is_refund”: ”bool”}partial_results = []chunk_metrics = []running_totals = pd.Series(0, index=regions, dtype=”int64”)running_records = []# read_csv返回可迭代的TextFileReader,每次只产生一个DataFramereader = pd.read_csv(CSV_PATH,usecols=use_columns,dtype=dtype_map,chunksize=CHUNK_SIZE,)for chunk_id, chunk in enumerate(reader, start=1):# 金额使用整数分,避免分块相加时引入浮点舍入误差chunk[”net_cents”] = np.where(chunk[”is_refund”],-chunk[”amount_cents”],chunk[”amount_cents”],)partial = (chunk.groupby(”region”, as_index=False).agg(order_count=(”amount_cents”, ”size”),refund_count=(”is_refund”, ”sum”),net_cents=(”net_cents”, ”sum”),))partial_results.append(partial)# deep=True把字符串内容也纳入DataFrame内存估算chunk_memory_mb = chunk.memory_usage(deep=True).sum() / 1024**2chunk_metrics.append({”chunk”: chunk_id, ”rows”: len(chunk), ”memory_mb”: chunk_memory_mb})# 保存每块之后的累计结果,用于观察流式聚合如何收敛current = partial.set_index(”region”)[”net_cents”].reindex(regions, fill_value=0)running_totals = running_totals.add(current, fill_value=0).astype(”int64”)for region in regions:running_records.append({”chunk”: chunk_id,”region”: region,”net_sales_million”: running_totals[region] / 100 / 1_000_000,})# 第二次groupby只合并四行一块的部分结果,而不是原始明细stream_result = (pd.concat(partial_results, ignore_index=True).groupby(”region”, as_index=False)[[”order_count”, ”refund_count”, ”net_cents”]].sum().sort_values(”region”).reset_index(drop=True))chunk_metrics = pd.DataFrame(chunk_metrics)running_results = pd.DataFrame(running_records)print(f”实际分块数:{len(chunk_metrics)}”)print(f”每块行数范围:{chunk_metrics['rows'].min():,}—{chunk_metrics['rows'].max():,}”)print(f”累计订单数:{stream_result['order_count'].sum():,}”)print(f”累计退款数:{stream_result['refund_count'].sum():,}”)print(f”净销售额:{stream_result['net_cents'].sum() / 100:,.2f} 元”)
输出:
实际分块数:8每块行数范围:30,000—30,000累计订单数:240,000累计退款数:13,104净销售额:31,445,148.05 元
八块的行数之和正好是24万,说明没有漏读。每块内部只留下四个地区的部分统计,所以第二次groupby处理的是最多32行小表,而不是再次处理24万行原始明细。
# 全量读取只用于本教程校验,实际内存受限任务不需要这一步full_data = pd.read_csv(CSV_PATH, usecols=use_columns, dtype=dtype_map)full_data[”net_cents”] = np.where(full_data[”is_refund”],-full_data[”amount_cents”],full_data[”amount_cents”],)full_result = (full_data.groupby(”region”, as_index=False).agg(order_count=(”amount_cents”, ”size”),refund_count=(”is_refund”, ”sum”),net_cents=(”net_cents”, ”sum”),).sort_values(”region”).reset_index(drop=True))# 整数金额使两条路径能够逐项精确比较,而不只比较总和assert_frame_equal(stream_result, full_result, check_dtype=True)# 这里比较DataFrame自身占用,不把解释器和绘图库内存算进去full_memory_mb = full_data.memory_usage(deep=True).sum() / 1024**2peak_chunk_mb = chunk_metrics[”memory_mb”].max()memory_ratio = full_memory_mb / peak_chunk_mbmax_difference = int((stream_result[”net_cents”] - full_result[”net_cents”]).abs().max())print(”校验结果:分块聚合与全量聚合逐项一致”)print(f”最大金额差:{max_difference} 分”)print(f”全量DataFrame:{full_memory_mb:.2f} MB”)print(f”最大单块DataFrame:{peak_chunk_mb:.2f} MB”)print(f”本例DataFrame占用比:{memory_ratio:.1f} 倍”)
输出:
校验结果:分块聚合与全量聚合逐项一致最大金额差:0 分全量DataFrame:16.13 MB最大单块DataFrame:2.02 MB本例DataFrame占用比:8.0 倍
assert_frame_equal()比较的不只是总销售额,而是四个地区的订单数、退款数和净销售额全部字段。结果逐项一致,最大金额差为0分。内存数字来自DataFrame.memory_usage(deep=True),它适合比较本例两张DataFrame,但不是整个Python进程的峰值内存。
# 图一展示各地区累计净销售额随块推进的变化fig, ax = plt.subplots(figsize=(9.2, 5.2))for region in regions:region_data = running_results[running_results[”region”] == region]ax.plot(region_data[”chunk”],region_data[”net_sales_million”],marker=”o”,linewidth=2.2,label=region,)ax.set_title(”Cumulative Net Sales During Chunk Processing”, pad=14, weight=”bold”)ax.set_xlabel(”Chunk completed”)ax.set_ylabel(”Net sales (million CNY)”)ax.set_xticks(chunk_metrics[”chunk”])ax.grid(axis=”y”, alpha=0.25)ax.legend(frameon=False, ncol=4, loc=”upper left”)fig.tight_layout()progress_plot = Path(”/tmp/python_chunk_progress.png”)fig.savefig(progress_plot, dpi=180, bbox_inches=”tight”)plt.close(fig)# 图二对比全量DataFrame与最大单块DataFrame的深度内存估算fig, ax = plt.subplots(figsize=(7.6, 5.2))labels = [”Full read”, f”Chunked\n({CHUNK_SIZE:,} rows)”]values = [full_memory_mb, peak_chunk_mb]bars = ax.bar(labels, values, color=[”#E76F51”, ”#2A9D8F”], width=0.58)ax.set_title(”Estimated DataFrame Memory Footprint”, pad=14, weight=”bold”)ax.set_xlabel(”Reading strategy”)ax.set_ylabel(”Memory (MB, deep=True)”)ax.set_ylim(0, full_memory_mb * 1.16)ax.bar_label(bars, labels=[f”{value:.2f} MB” for value in values], padding=4, weight=”bold”)ax.grid(axis=”y”, alpha=0.22)fig.tight_layout()memory_plot = Path(”/tmp/python_chunk_memory.png”)fig.savefig(memory_plot, dpi=180, bbox_inches=”tight”)plt.close(fig)# 输出图像大小和最终差异,便于核对绘图与聚合结果print(f”Progress chart:{progress_plot.name} ({progress_plot.stat().st_size:,} bytes)”)print(f”Memory chart:{memory_plot.name} ({memory_plot.stat().st_size:,} bytes)”)print(f”Regions compared:{len(stream_result)}”)print(f”Rows processed:{chunk_metrics['rows'].sum():,}”)print(f”Validation difference:{max_difference} 分”)
输出:
Progress chart:python_chunk_progress.png (117,613 bytes)Memory chart:python_chunk_memory.png (44,929 bytes)Regions compared:4Rows processed:240,000Validation difference:0 分

累计轨迹展示了“部分结果不断合并”的过程。每完成一块,四个地区各增加一个累计点;到第八块时,曲线终点就是最终聚合值。曲线近似平稳上升来自本次随机模拟,不能据此推断真实地区销售趋势,因为横轴是文件块序号,不是时间。

在相同三列和相同数据类型下,全量DataFrame估算为16.13 MB,最大单块为2.02 MB,比例正好约8倍,与24万行拆成8个等长块一致。这个比例不应机械外推到所有任务:部分结果大小、中间对象、字符串长度和处理函数都会影响真实峰值。
本次实跑说明,分块并不意味着牺牲聚合准确性。因为订单数、退款数和整数金额总和都可以在块间相加,两阶段结果与全量结果逐项一致。chunksize控制的是单次返回的行数;它不会自动替你设计正确的部分统计和合并逻辑。
内存比较也要准确表述。2.02 MB和16.13 MB是所选三列加净额列的DataFrame深度占用,不包含Python解释器、CSV解析器、Matplotlib、部分结果列表和其他对象。因此可以说“本例最大单块DataFrame约为全量的八分之一”,不能说“整个程序内存一定降低八倍”。
块越小,单块占用通常越低,但循环、解析和合并次数会增加;块越大,调度开销较少,却更接近内存上限。实践中可以先用一个保守值运行,记录单块行数、处理耗时和进程峰值,再逐步调整,而不是寻找适用于所有数据的固定数字。
pd.read_csv()的chunksize参数真正有价值的地方,是把一个必须整体进入内存的问题改写成一串可独立处理的小问题。对于计数、求和、最小值、最大值等可合并统计,先在块内聚合、再合并部分结果,通常是一条清楚而可靠的路径。
使用时需要同时守住三个边界:只读取需要的列,为关键字段明确数据类型,并确认目标统计量有正确的块间合并规则。若任务需要大量跨块连接、全局排序、精确分位数或复杂窗口,手写分块流程会迅速变难,应考虑数据库、DuckDB、Polars流式执行或Dask等更合适的工具。
