发布时间:2026-08-07 14:35:34 分类:营销学堂
multiprocessing 进程池,结果代码写了 80 多行,不仅手动分块、合并结果让人头疼,最后还因为内存溢出(OOM)直接崩溃。multiprocessing 更顺手。multiprocessing 模块本身并不复杂,但数据处理远不止“把计算分布到多个进程”那么简单。当你试图用它来处理大型 DataFrame 时,通常会撞上三座大山:Pool.map 需要通过 pickle 将数据块序列化传给子进程,算完后再序列化传回。对于大 DataFrame,这种序列化/反序列化的开销,往往比实际计算还要耗时。Pool.map 会把所有子进程的结果一次性收集到主进程的内存中。8 个 chunk 的结果同时驻留内存,极易将内存撑爆。multiprocessing 的 bug,而是它的定位。它是一个通用的并行原语,就像一把瑞士军刀,能剪能锯;但你要用它来劈柴,显然不如一把专门的大斧头顺手。import dask.dataframe as dd# 用 Dask 读大 CSV,自动按 128MB 分成多个 partitionddf = dd.read_csv("large_orders.csv", dtype={"user_id": "int32", "amount": "float32"}, blocksize="128MB") # Pandas 风格的语法,Dask 自动处理分块和合并result = ddf.groupby("user_id")["amount"].sum().nlargest(100)# 注意:.compute() 时才真正执行计算top_users = result.compute() multiprocessing 的手动操作,且不易出错。ddf.groupby(...) 时,它并不会立即执行,而是构建了一个任务图(Task Graph),记录下“读文件 → 分组 → 聚合 → 排序”这一串操作。直到你调用 .compute() 时,Dask 才会开始分析这个图,并施展两个魔法:filter 后面接了 groupby,它会把过滤操作下推到 partition 级别,避免不必要的数据传输。multiprocessing 那样序列化传回主进程。对于多阶段计算,这省去了巨大的 I/O 开销。multiprocessing 反而更合适:场景 | 推荐工具 | 核心原因 |
Pandas DataFrame 并行 | Dask | 自动分块、无缝语法、不易 OOM |
纯 Python 函数并行 | multiprocessing | 无额外依赖,调度更轻量 |
超大数据集(>内存) | Dask | 分区惰性加载,完美应对内存瓶颈 |
极轻量级微秒计算 | multiprocessing | 避免 Dask 任务图的调度开销大于计算本身 |
GPU 计算 | multiprocessing | Dask 对 GPU 的调度支持目前尚不够完善 |
multiprocessing,再到 Dask,我们追求的始终是“用更优雅的代码解决更复杂的问题”。