8.3 并行与并发处理:多灶开火


8.3 并行与并发处理:多灶开火

本节摘要:向量化榨干单灶效率之后,还有多灶开火——多进程绕开解释器的全局锁、进程池摊派独立任务。本节讲清多进程与多线程的分工场面、进程池的标准写法,以及"并行不划算"的几种情况。承接 8.2 的内存账,通往 8.4 的排错手册。

深夜的批处理窗口总是不够用

深夜的批处理窗口定在零点到六点,任务却越排越满:三十个月的月报要逐月重算、数百个门店的数据要逐店清洗。单灶干到天亮也干不完——可这些任务彼此独立,天然适合多灶开火:每月一个灶、每店一个灶,同时开工,窗口自然够用。本节讲开灶的规矩:哪类活能摊、用进程还是线程、池子怎么开。往前接 8.1 的提速与 8.2 的内存约束(灶多内存也多),往后 8.4 的排错里,并行任务出错往往更难查,所以规矩要立在开灶之前。

参数拆解:进程池的旋钮

multiprocessing.Pool:Pool(processes=None) 不给数就按机器核数开灶;pool.map(函数, 任务列表, chunksize=...) 把任务摊给各灶,chunksize 控制每次派几件,任务碎时给大一点省通信;用完记得 close 加 join,或干脆用 with 管理生命周期。**concurrent.futures.ProcessPoolExecutor**:进程池的现代化封装,submit 返回未来对象、as_completed 边完成边收——哪个灶先出菜就先上哪盘,适合耗时不均的任务。**线程版 ThreadPoolExecutor**:同一套接口,适用场面不同(见下)。**joblib.Parallel(n_jobs=-1)**:机器学习圈常用的第三方并行封装,一行摊派、结果自动收集,与 sklearn 生态最合拍。

from multiprocessing import Pool import pandas as pd def monthly_report(月份: str) -> dict: """独立任务:单月清洗加聚合,返回可合并的结果。""" df = pd.read_csv(f"orders_{月份}.csv") df = df.dropna(subset=["金额"]) return {"月份": 月份, "总额": df["金额"].sum()} if __name__ == "__main__": months = [f"2024-{m:02d}" for m in range(1, 13)] with Pool(processes=4) as pool: # 开四个灶 results = pool.map(monthly_report, months, chunksize=1) print(pd.DataFrame(results).sort_values("月份").head(3))

实操示例:一批独立文件的并行清洗

场景:三百个门店文件,逐个串行清洗要很久;按"每店一灶"摊给进程池,收齐结果统一汇总。演示 futures 写法——它对"各店耗时悬殊"最友好。

from concurrent.futures import ProcessPoolExecutor, as_completed import glob def clean_one(path: str) -> dict: df = pd.read_csv(path) df = df.drop_duplicates().dropna(subset=["金额"]) return {"门店": path, "行数": len(df), "总额": df["金额"].sum()} files = glob.glob("stores/*.csv") summary = [] with ProcessPoolExecutor(max_workers=8) as ex: futures = {ex.submit(clean_one, f): f for f in files} for fut in as_completed(futures): summary.append(fut.result()) # 谁先完成先收谁 report = pd.DataFrame(summary) print(report["总额"].sum()) # 各灶结果合并为总账

开灶前的收益估算

并行不是免费午餐,动手前估算一轮:加速上限受灶数限制,但实际收益被三处摊薄——任务的拆分与合并成本、子进程的启动开销、内存的复制占用。经验判断式:单任务秒级以上、任务数明显多于灶数、任务间零依赖,三条同时满足,并行的收益大概率值回票价;任何一条不满足,先回 8.1 看向量化还有没有空间。估算花不了几分钟,却能省掉一个下午的并行调试——多灶的锅,坏起来比单灶难查得多。

坑点与翻车

**翻车一:给 CPU 密集任务开线程。**解释器的全局锁让同一时刻只有一个线程在执行字节码——线程并行对纯计算无效,开八个线程等于排队。分工口诀:计算密集用进程,等网络等磁盘的输入输出密集用线程。**翻车二:任务太碎还要并行。**每个任务毫秒级完成,摊派与收结果的通信开销反而吃掉全部收益——并行前估一下单任务耗时,秒级以下先考虑合并任务或 chunksize。**翻车三:大表当参数传进池子。**每提交一次,表就被完整复制一份发往子进程,八个灶就是八份内存——把"读数据"也放进子任务函数里(像示例那样传路径不传表),或者用 7.4 的 Dask 处理跨进程共享。**翻车四:各灶随机种子相同。**每个子进程继承同样的随机状态,抽样或模拟出的结果八灶一模一样——子任务里显式给 random_state 或用 os.urandom 播种。另外两处纪律:Windows 上进程池必须包在 if 判断主模块的语句里,否则无限递归开灶;子进程里的报错不会自动冒泡,fut.result() 时才抛——收结果处要接 8.4 的异常处理。

替代方案

开灶之前按序想三问:能不能向量化(8.1,永远的首选)?能不能分块串行(6.3,开销最小)?数据量是不是大到该用 Dask(7.4,灶台自管调度)?多数"想并行"的场景,答案在前三问里。第三方封装里,joblib 适合函数式的小摊派,swifter 与 pandarallel 能给 apply 自动选并行档——便利有余,可控性不足,生产脚本建议还是用标准库的池子,行为可预期。真正要跨机器并行,进程池就不顶用了,那是 7.4 大灶台的领域。

收档清单

  • 分工口诀:计算密集开进程,输入输出密集开线程;
  • 池子两式:Pool 配 map 管均匀任务,futures 配 as_completed 管悬殊任务;
  • 传路径不传表:大对象进子进程就是复制,内存按灶数翻倍;
  • 播种要显式:子进程随机态同源,抽样模拟各灶要各自播种;
  • 先想三问:向量化、分块串行、Dask,都不合适才真开灶。

多灶开火之后,出错的面也大了——8.4 把报错家族与调试手册备齐,让翻车可查可救。


作者与出处
原作者: 灏天文库
来源:灏天文库
整理: 灏天文库整理
由灏天文库平台收录,内容或由平台用户上传,仅供学习交流
发布者: 作者: 灏天文库 转发
评论区 (0)
U