当数据规模膨胀到千万行以上时,常见的Python脚本往往在读取阶段就消耗掉全部内存,随后进入漫长的垃圾回收与换页等待。真正可行的路线是把计算拆成可控的块,在有限资源内流水式推进,而不是试图一次性把所有内容塞进内存。本节我们先建立分块处理的基本认知,再逐步展开工程层面的优化手法。

一、为什么全量加载会失效
Python的列表与pandas的DataFrame本质都基于连续或分散的对象引用。以pandas为例,一个包含字符串的列在底层会包装为Python对象数组,单条记录的内存开销远超原始字节。当数据量达到千万行,即使数值型字段也会因为DataFrame的块管理结构产生额外元数据成本,轻松突破单机数十GB内存上限。
除了内存,全量加载还会引发隐性性能陷阱。比如使用pd.read_csv不限制体积时,解析器需要在内存中动态扩容,触发多次重新分配;后续若再做df.merge或df.groupby,中间结果同样以全量形态存在,导致CPU缓存命中率骤降。理解这一点,才能明白分块不是权宜之计,而是大规模数据下的必然范式。
二、基于chunksize的分块读取
pandas自带的read_csv提供了chunksize参数,它返回一个可迭代的TextFileReader对象,每次只产出指定行数的DataFrame。这样主进程内存中始终只有一块数据,处理完即可释放。下面示例展示如何对超大文件做分块求和并合并统计。
import pandas as pd
total_sum = 0
count = 0
# 每次读取十万行,避免内存峰值
chunk_iter = pd.read_csv('big_data.csv', chunksize=100000)
for chunk in chunk_iter:
# 向量化计算当前块的数值列之和
total_sum += chunk['value'].sum()
count += len(chunk)
print('平均值:', total_sum / count)
上述代码的核心在于chunk['value'].sum()走的是numpy底层循环,比纯Python的for累加快一个数量级。分块后每个chunk独立参与向量化,整体时间复杂度接近全量计算,但内存占用仅为单块大小。
需要注意的是,如果最终要拼接所有块的结果,应避免在循环里频繁pd.concat,因为拼接本身会产生复制。正确做法是像示例那样只保留标量聚合,或把每块结果追加到列表,最后做一次合并。很多初学者在循环内concat导致耗时翻倍,这便是典型的误操作。
三、多进程加速分块计算
当单块内计算较重(如复杂特征工程),可以利用multiprocessing把不同块派发到多核。由于每块数据独立,进程间几乎无需通信,扩展性很好。以下代码演示进程池处理分块文件。
import pandas as pd
from multiprocessing import Pool
def process_chunk(file_path_and_skip):
path, skip = file_path_and_skip
# 读取指定偏移的小块
chunk = pd.read_csv(path, skiprows=skip, nrows=100000)
# 模拟耗时变换
return chunk['value'].mean()
if __name__ == '__main__':
tasks = [('big_data.csv', i) for i in range(0, 1000000, 100000)]
with Pool(4) as p:
results = p.map(process_chunk, tasks)
print('各块均值:', results)
这里用skiprows与nrows手动切分,配合四进程并行,理论加速比接近线性。但需留意磁盘IO可能成为瓶颈:多进程同时读同一块机械盘反而会因寻道竞争变慢,此时可先按块拆成多个物理文件,或改用SSD与内存映射。
另外,Windows下多进程会重新导入模块,务必把任务函数放在if __name__ == '__main__'保护之后,否则容易递归派生进程。Linux的fork模式虽无此忧,但也要注意每个子进程持有独立内存副本,块不能过大。
四、numpy内存布局与类型压缩
分块之外,字段类型选择直接决定内存水位。pandas默认用int64、float64存储数字,若业务允许,可降级为int32、float32,甚至用category类型表达低基数字符串。在块内处理前显式转换,能削减一半以上占用。
import pandas as pd
chunk_iter = pd.read_csv('big_data.csv', chunksize=100000)
for chunk in chunk_iter:
# 压缩类型后再运算
chunk['value'] = chunk['value'].astype('float32')
chunk['type'] = chunk['type'].astype('category')
# 后续分析逻辑
pass
numpy的连续数组在CPU预取机制下表现优异,而category列内部使用整数编码,比较与分组操作远快于对象列。将这两点融入分块流程,既省内存又提速,是生产环境常用的组合拳。
还需关注内存对齐与副本问题:某些pandas操作会静默产生副本,如chunk[['a','b']]切列视图在修改时可能触发copy-on-write。明确使用.loc或.copy()控制生命周期,可以减少无谓分配,让分块流水线更顺畅。
五、磁盘与GC层面的辅助调优
即便算法分块合理,频繁的小文件读写和Python的自动垃圾回收也会引入抖动。建议将临时块以feather或parquet格式落盘,这类列式存储读写均快于CSV,且保留类型信息,避免重复解析。
import pandas as pd
chunk = pd.read_csv('part.csv')
# 转为parquet便于后续快速载入
chunk.to_parquet('part.parquet', index=False)
GC方面,可在分块循环里手动调用gc.collect()及时释放循环引用,尤其当块内创建了大量临时DataFrame时。同时适当调大csv读取的engine为C版,并关闭不必要的类型推断,都能让千万级任务稳定跑完而不必堆砌服务器。
综合来看,分块计算并非单点技巧,而是从读取、类型、并行到落盘的一整套权衡。按本文顺序梳理自己的管线,多数Python大数据慢局都能在普通开发机上化解。