写 Python 的人,对线程池一定不陌生。requests 并发下载、数据库批量查询,用 ThreadPoolExecutor 一把梭,确实省心。但真遇到纯 CPU 密集的任务,比如批量图像处理、蒙特卡洛模拟、量化交易策略回测里的参数扫描,线程池的表现通常会让人怀疑人生。这时候 ProcessPoolExecutor 才是更靠谱的答案。它是 concurrent.futures 标准库自带的多进程调度组件,接口和线程池几乎一模一样,但底层会启动真正的独立子进程,让任务分布在不同的 CPU 核上。因为绕过了 CPython 的 GIL,它能把多核机器真正跑满,是 Python 中做 CPU 密集并发时最值得优先考虑的一个组件。
不管是刚开始写 Python 脚本的小白,还是已经在处理数据分析、爬虫解析、模型仿真这类任务的工程师,这篇文章都适合你。我会从原理讲到实操,再把我自己在生产环境里踩过的坑一并整理出来。你不需要先精通操作系统多进程知识,只需要能看懂 Python 函数调用,就能照着落地。
1. 为什么绕不开 ProcessPoolExecutor
1.1 GIL 才是大多数 Python 并行问题的根源
很多新手不理解一个现象:同一个计算函数,用 ThreadPoolExecutor 同时开 16 个线程去跑,CPU 占用率却始终只在一个核上跳动。这就是 GIL 在起作用。CPython 解释器里有一把全局解释器锁,它保证同一时刻只有一个线程能执行 Python 字节码。
你可以把它想象成一家只有一个收银台的超市。收银员在结账,所有顾客排队,线程切换只是在队伍里换来换去,但同一时间能完成结账的只有一个人。对于 I/O 密集操作,比如网络请求、文件读写,线程在等待底层系统返回时会把 GIL 让出来,所以多线程能获得明显的并发提升。可一旦任务里全是纯 Python 的循环计算,GIL 几乎不会被释放,多线程不仅没有加速,还可能因为锁切换而变慢。
ProcessPoolExecutor 的思路完全不同。它直接把任务分给多个独立进程,每个进程都有自己的 Python 解释器和独立内存空间,自然也有各自的 GIL。多个进程可以同时跑在不同 CPU 核上,从根上避开了 GIL 的粒度和竞争问题。
1.2 ThreadPoolExecutor 和 ProcessPoolExecutor 怎么选
很多人一上来就看表象:两者都叫 Executor,都有 submit、map、shutdown,感觉用法一样。但选错了,性能差距能到数倍甚至数十倍。
我整理过一个非常粗略的选型表,你可以直接对照:
| 维度 | ThreadPoolExecutor | ProcessPoolExecutor |
|---|---|---|
| 适用核心场景 | I/O 密集任务:网络请求、文件读取、数据库操作 | CPU 密集任务:复杂计算、循环遍历、模拟仿真 |
| 是否受 GIL 影响 | 受,纯计算时多线程无法并行 | 不受,每个子进程独立解释器 |
| 创建资源成本 | 低,线程开销远小于进程 | 高,需要创建独立进程和解释器 |
| 数据传递 | 共享同一个进程内存,但要注意线程安全 | 任务参数和返回值需要 pickle 序列化 |
| 稳定性影响 | 线程中操作不当可能影响整个进程 | 子进程崩溃未必立刻影响主进程,但池可能废弃 |
| 适合新手程度 | 简单,容易写成资源竞争 | 相对复杂,需要理解对象序列化和程序入口保护 |
一句话概括:任务在等网络、等磁盘、等数据库,用 ThreadPoolExecutor;任务在计算、杀 CPU、跑数学模拟,用 ProcessPoolExecutor。如果你拿不准,可以写一个小函数,用单线程先跑一轮采样,看看 CPU 占用率。如果 CPU 占用接近 100%,说明确实是计算密集,别浪费时间去调线程池了。
1.3 实际场景里它到底能解决什么问题
我处理过不少批处理脚本,最常见的情况不是单次计算不够快,而是数据量大、任务数量多。比如有 1 万份文档需要做特征提取,每份文档里有一段几万轮的循环逻辑,单进程跑要 3 个小时;用 ThreadPoolExecutor 改成 8 线程后,时间只缩短到 2 小时 40 分钟,CPU 还是上不去。后来改成 ProcessPoolExecutor,8 个 worker 直接压满八个核,时间缩短到 22 分钟左右。
类似适合 ProcessPoolExecutor 的典型场景还有:
- 蒙特卡洛模拟、期权定价、风险价值计算这类金融数值模拟。
- 量化交易策略回测时,对一组参数组合做批量扫描。
- 批量图像处理中的像素级滤镜,或者视频抽帧后的逐帧计算。
- 文本批量处理中的大规模分词、特征工程。
- 科学计算里无法直接用 NumPy 向量化替代的循环逻辑。
需要注意,如果某个任务本身非常轻量,比如只是给数字加 1,那开进程池反而不划算。因为每个任务都要序列化数据、创建进程、进入队列等待,当进程创建和通信成本高于计算成本时,整体会比单线程更慢。用之前最好先问问自己:单次任务有没有做到“足够重”,重到值得让一个进程来回折腾一次。
2. 核心细节解析与实操要点
2.1 三种提交任务方式:submit、map、as_completed
ProcessPoolExecutor 最基础的用法非常简单,先看最经典的 submit 场景:
from concurrent.futures import ProcessPoolExecutor def calculate(x): return x * x if __name__ == "__main__": with ProcessPoolExecutor(max_workers=4) as executor: future = executor.submit(calculate, 10) print(future.result())submit 返回一个 Future 对象,可以把它理解成一张“任务小票”。主进程拿到小票后,可以继续干别的事,等需要具体结果时再调用 result() 阻塞等待。如果提交多个任务,通常搭配 as_completed 使用:
from concurrent.futures import ProcessPoolExecutor, as_completed def calculate(x): return x * x if __name__ == "__main__": tasks = range(100) with ProcessPoolExecutor(max_workers=4) as executor: futures = [executor.submit(calculate, i) for i in tasks] for future in as_completed(futures): result = future.result() print(result)as_completed 会按照任务实际完成的时间返回已完成的 Future,而不是按提交顺序。如果任务耗时差异比较大,用这种方式可以第一时间拿到已经结束的结果,主流程不需要等前面的慢任务。
executor.map 也是一个常用接口,它会按输入顺序返回结果:
if __name__ == "__main__": with ProcessPoolExecutor(max_workers=4) as executor: results = executor.map(calculate, range(100)) for res in results: print(res)map 的优点是代码更短,缺点是不能方便地拿到任务状态,只能按顺序等结果。如果某一个任务异常,可能影响整体遍历。所以我个人更推荐批量提交后用 as_completed,尤其是任务多、异常可能多的场景。
2.2 进程模型与“必须写 ifname== 'main'”的原因
ProcessPoolExecutor 并不是每个任务都新建进程,而是启动一组固定数量的 worker 进程,默认情况下进程数由 max_workers 决定。任务提交后会被放进内部的任务队列,空闲的 worker 会从队列里取出下一个任务执行。任务执行过程中,主进程和子进程通过队列传递消息,子进程拿到的是参数对象的序列化副本,运算结果也会被序列化回主进程。
这里最关键的一点是:Python 在不同操作系统上启动子进程的方式不同。在 Linux 上,默认常见的是 fork 模式,子进程直接复制父进程内存镜像,所以不写 ifname== 'main',有些代码也能跑通,容易让人产生侥幸心理。在 Windows 上,或者使用 spawn 模式时,Python 必须重新导入主模块来启动一个干净的子进程。
如果你不保护入口,Windows 下会发生一个很经典的问题:程序导入主模块时又遇到 ProcessPoolExecutor,于是又创建子进程,子进程导入主模块时再创建下一层子进程,最终无限递归,程序还没跑任务就直接崩溃。所以为了兼容任何平台,所有能正常运行的多进程代码都应该把启动逻辑放在 ifname== 'main' 里。这不是“Windows 专属要求”,而是 Python 多进程程序的通用习惯。
另外要特别提醒:如果你在 Jupyter Notebook 里直接运行 ProcessPoolExecutor,也容易出现各种不可思议的重复启动问题。这是因为 Notebook 执行单元的过程更复杂。最好把任务函数写在 .py 文件里,或者用函数封装后通过模块导入方式使用。
2.3 max_workers 到底应该设置成多少
这个参数看起来简单,实际很容易拍脑袋。Python 默认会根据 CPU 数量计算:min(32, os.cpu_count() + 4)。这个默认值偏保守,是为了避免一上来就把机器资源占满。可是如果你用 4 核的笔记本跑 CPU 密集任务,默认会生成 8 个 worker,反而可能因为进程间切换和内存开销让性能下降。
我个人在 CPU 密集场景下的经验是:先设成 CPU 核心数,跑一遍,再往上和往下各试一档。比如:
import os from concurrent.futures import ProcessPoolExecutor workers = os.cpu_count() or 4 with ProcessPoolExecutor(max_workers=workers) as executor: ...如果你不希望对机器其他程序造成太大压力,可以手动减半:
workers = max(2, os.cpu_count() // 2)还有一个容易被忽略的点:如果任务本身占用的内存特别大,比如每个任务都要加载一个几百 MB 的模型文件,那么 worker 数量可不能只看 CPU 核心数。多个进程会把同样的模型加载多份,内存很容易被吃满。我遇到过有人把 8 核机器上所有 worker 都铺满,结果每个任务加载同一个大词典,内存直接爆掉,系统开始疯狂使用交换分区,速度比单进程还慢。这种情况下,workers 要按内存预估来设置,而不是按 CPU 核数。
2.4 pickle 序列化是隐藏的“性能杀手”和“报错源头”
提交给 ProcessPoolExecutor 的每个任务,函数本身和参数都要被序列化,也就是 pickle。子进程算完后,结果又要被序列化传回主进程。这意味着你传的对象体积越大,序列化开销越高;对象不能被 pickle,代码就会直接抛错。
最常见的错误是传 lambda 函数。很多人在写脚本时习惯用 lambda 一写,丢给 executor.map,结果报错说找不到函数。原因是 lambda 无法被正常 pickle。另一个常见错误是传“局部嵌套函数”。比如在函数内部再定义一个 inner 函数,然后提交给进程池,在 spawn 模式大概率会失败。还有不少人不小心把某个类的实例方法作为任务函数传进去,实例对象如果没有实现 pickle 协议,也会失败。
解决方案非常朴素:把任务函数定义到模块顶层,参数尽量用基础类型、路径字符串、简单字典。如果一定要传复杂对象,尽量实现getstate和setstate,但这会显著增加维护成本。更有用的技巧是:不要在任务参数里传大批量数据,把数据写到临时文件,任务传文件路径,让子进程自己读文件。这样既避免了重复多次 pickle,也让内存模型更合理。
3. 实操:用 ProcessPoolExecutor 改造一个 CPU 密集计算任务
3.1 从串行版本开始:蒙特卡洛模拟
我拿一个非常有代表性的计算密集型问题来说明:用蒙特卡洛方法估算圆周率。思路是随机生成平面上的点,统计落在单位圆内的比例,乘 4 得到 π 的近似值。点越多,结果越准,计算量也越大。如果你批量跑 2000 万次随机采样,单线程确实要等一会儿,很适合用来观察进程池效果。
先写一个串行版本:
import math import random import time def run_serial(total_points): rng = random.Random(20240701) inside = 0 for _ in range(total_points): x = rng.random() y = rng.random() if x * x + y * y <= 1.0: inside += 1 return 4 * inside / total_points if __name__ == "__main__": total_points = 20_000_000 start = time.perf_counter() pi_estimate = run_serial(total_points) elapsed = time.perf_counter() - start print(f"串行结果: {pi_estimate:.6f}, 耗时: {elapsed:.2f}s") print(f"误差: {abs(pi_estimate - math.pi):.6f}")这段代码就是典型的 CPU 密集循环。如果换 ThreadPoolExecutor 去并行,因为 Python 的随机数生成和循环计算都在 GIL 内,多线程不仅不能加速,反而可能因为线程切换变得更慢。所以它是最适合改成多进程的一类任务。
3.2 用 ProcessPoolExecutor 拆分成多个并行子任务
我们需要把 2000 万次采样拆成多个小块,每个子进程负责其中一块,最后在主进程汇总。为了任务分配更均匀,我把任务拆成“worker 数量的 4 倍”,也就是每个 worker 会连续处理多个小块,避免某些进程提前结束导致 CPU 空闲。
下面是完整示例:
import math import os import random import time from concurrent.futures import ProcessPoolExecutor, as_completed def count_inside(points, seed): rng = random.Random(seed) inside = 0 for _ in range(points): x = rng.random() y = rng.random() if x * x + y * y <= 1.0: inside += 1 return inside def run_parallel(total_points, workers): # 每个 worker 分 4 个更小任务,让负载更均衡 task_count = workers * 4 chunk_size = total_points // task_count inside_total = 0 with ProcessPoolExecutor(max_workers=workers) as executor: futures = [] for i in range(task_count): futures.append(executor.submit(count_inside, chunk_size, 10000 + i)) # 多少内完成一个就统计一个 for future in as_completed(futures): inside_total += future.result() return 4 * inside_total / total_points if __name__ == "__main__": total_points = 20_000_000 workers = os.cpu_count() or 4 start = time.perf_counter() pi_estimate = run_parallel(total_points, workers) elapsed = time.perf_counter() - start print(f"并行结果: {pi_estimate:.6f}, 耗时: {elapsed:.2f}s") print(f"使用 worker: {workers}") print(f"误差: {abs(pi_estimate - math.pi):.6f}")你可以注意几个细节。count_inside 被定义在模块顶层,参数是整数和小种子,全部能 pickle。每次 submit 时只传递两个数,返回的也是一个整数,重量级数据并没有在主进程和子进程之间反复搬运。用 as_completed 遍历,哪个任务先结束就先累加结果,不需要排着队等最早任务。
3.3 串行与多进程的实测对比
我在一台普通 8 核笔记本上用 Python 3.11 跑上面的代码,2000 万次采样,串行大约需要 9.6 秒左右,改成 8 个 worker 后大约在 2.5 秒左右,差不多能有 3.8 倍加速。如果只用 4 个 worker,耗时大约 3.8 秒。这里没有达到理论上的 8 倍加速,原因在于进程池启动、任务调度、pickle 传参、结果汇聚都有开销,而且一个 Python 进程里随机数生成也并非无限线性可扩展。
具体结论我整理成了下面的表,方便你参考趋势:
| 方案 | 配置 | 参考耗时 | 加速比 |
|---|---|---|---|
| 串行 | 单进程 | 9.6s | 1x |
| ProcessPoolExecutor | 4 workers | 3.8s | 约 2.5x |
| ProcessPoolExecutor | 8 workers | 2.5s | 约 3.8x |
| ProcessPoolExecutor | 16 workers | 2.4s | 基本不再提升 |
你可以看到,worker 数量超过 CPU 核心数后,提升会非常有限。因为机器只有 8 个物理核心,16 个进程同样只能挤在 8 个核上,多出来的进程反而增加系统调度和内存压力。真正的实战中,我建议用 4、8、12、16 几个档位都测一轮,找到你所在机器和任务负载下的“甜点值”。
3.4 改并行后最容易忽略的优化点
代码能跑和跑得高效之间,还有几层细节。
第一,不要在循环里频繁调用 executor.submit。如果你有 10 万个轻量任务,每个任务提交都带着 pickle 和队列通信,主进程很快会成为瓶颈。这时候可以把任务按“批次”合并,例如每 1000 条记录为一组,每次 submit 一个批次任务,让 worker 内部再处理这个批次。批次任务的计算时间变长,单任务并发调度开销就被摊薄了。
第二,使用 with 语句管理 executor 很重要。with 块结束时,会自动调用 executor.shutdown(wait=True),确保所有任务完成且资源释放。如果忘了关闭进程池,进程会一直留在系统里,至少会让脚本退出变慢,严重点还可能造成句柄泄漏。
第三,从性能测试切换到生产代码时,最好保留一个 verbose 参数,方便单进程复现结果。多进程排错困难,如果能用单进程把小样本跑通,再切到多进程,会省去大量排查时间。
4. 常见问题与排查技巧实录
4.1 BrokenProcessPool:子进程到底是怎么“碎”的
这是 ProcessPoolExecutor 用户最容易遇到的异常之一。它的出现场景通常是某个 worker 进程不是正常抛出 Python 异常,而是直接崩溃退出。比如触发了 C 扩展的段错误、进程被操作系统 kill、调用 os._exit 强制退出等。
要注意区分两种情况。如果任务函数内部抛了 ValueError,Future.result() 会把 ValueError 正常抛回主进程,此时 pool 还是好的,其他任务能继续执行。但如果 worker 进程本身死了,executor 检测到后会把整个 pool 标记为不可用,之后就抛 BrokenProcessPool。
我在实际项目中就遇到过:某个第三方加密库在特定输入下会崩溃,导致整个进程池快速“碎裂”,后面几百个任务全部失败。解决办法是在 worker 函数最外层加一层很宽泛的异常保护,把有可能让进程崩掉的操作换成可预期的错误返回:
def safe_worker(item): try: return process_item(item) except Exception as exc: return None, f"failed: {exc}"当然 try except Exception 救不了段错误级别的崩溃,但至少能把多数普通异常挡在 Future 外面。除此之外,还要尽量减少在每个 worker 里加载不稳定的 C 扩展,如果必须加载,可以给每个任务单独创建子进程,而不是复用池。
4.2 任务卡死、超时设置和可怕的 Ctrl+C
有一个很普遍的错误认知:Future.result(timeout=3) 设了超时,如果任务超过 3 秒没完成,Python 就会杀掉子进程。实际上不会。这个 timeout 只控制主进程等待结果的时间。超时后,主进程这边抛 TimeoutError,但那个子进程任务可能还在后台继续运行,占着 CPU 和内存。
如果任务真的会跑很久,而你希望“超时就放弃”,ProcessPoolExecutor 原生并不支持。一个可行设计是把任务拆小,每个小任务都能在合理时间内返回,再由主进程控制整体进度。这样即使某一个小块很慢,你也能通过 as_completed 的超时逻辑提前结束等待。
另一个让人头疼的问题是 Ctrl+C。当你在脚本里运行 ProcessPoolExecutor 时,按 Ctrl+C 往往不会像单线程脚本那样快速退出。因为主进程要等待所有子进程结束,而某些子进程可能还在执行任务。最朴素的处理是捕获 KeyboardInterrupt,在 except 中调用 executor.shutdown(wait=False, cancel_futures=True) 尝试快速结束。但已经提交且开始执行的任务是无法取消的,所以更彻底的办法是让每个任务都足够短,这样中断时的响应时间才能被人类接受。
4.3 别把局部函数和 Lambda 丢给进程池
我见过太多这样的写法:
def outer(): def inner(x): return x * 2 with ProcessPoolExecutor() as executor: result = executor.submit(inner, 10)在很多 Linux 环境里,这段代码用默认 fork 方式可能能跑通,但一放到 Windows 或者改了启动方式就会报 pickle 错误。它背后的原因是,进程池要把 inner 函数序列化到子进程中,而 inner 是一个局部函数,Python 没办法通过模块路径找到它。Lambda 也是一样的道理,它连名字都没有,pickle 自然没法定位。
解决方法是把函数定义到模块顶层,让 Python 可以通过“模块名.函数名”的方式找到它。如果你需要给函数传额外参数,比如常量、配置项,用参数传进去,而不是在 lambda 里闭包捕获。这样做不仅跨平台更安全,也让代码隔离更清晰。
4.4 日志重复和 print 输出混乱
多进程程序里,日志重复是非常常见的现象。本质上每个子进程都有自己的一套标准输出和 logging handler。如果子进程里直接 print,你确实会看到不同进程的输出交替出现。如果主程序入口被多次 import,那么每个子进程都可能再次往同一个日志文件写入,自然很容易出现重复记录。
更好的实践是让所有 worker 只计算、返回结果,不在 worker 里写业务日志;主进程负责统一记录。如果需要实时进度,可以在 worker 里返回状态码,由主进程按 as_completed 汇总后打日志。如果确实需要在子进程里记录大量日志,可以考虑 logging.handlers.QueueHandler + QueueListener,把日志项发送到主进程的日志队列,由主进程统一落盘。这个方案比直接在每个子进程里配 FileHandler 更稳健。
4.5 ProcessPoolExecutor 和 multiprocessing.Pool 怎么选
multiprocessing.Pool 是 Python 更底层的多进程池,ProcessPoolExecutor 是基于它的简化封装。两者各有偏向。
ProcessPoolExecutor 的优势是接口统一,和 ThreadPoolExecutor 几乎一致,迁移成本低,Future 对象在并发处理上更自然。multiprocessing.Pool 的优势则是更底层、更灵活,提供 apply_async、map_async、imap、imap_unordered 等方法,还支持 maxtasksperchild,能让每个 worker 在处理一定数量任务后自动重启,避免某些有状态 worker 慢慢泄漏。
我个人的选择标准很简单:如果只是想让一批 CPU 密集函数并行跑,优先用 ProcessPoolExecutor,因为代码可读性好得多。如果要做非常精细的并发控制,比如需要分块迭代、保留 worker 状态、执行前初始化参数,那 multiprocessing.Pool 更合适。你不需要把两者都研究得很深,但要知道有另外一个方案存在,避免在 ProcessPoolExecutor 的局限里硬磕。
5. 高级玩法与并发架构设计思路
5.1 线程池和进程池混合编排,各干各的
很多真实任务并不只是“纯 CPU”或者“纯 I/O”,而是有 I/O 等待,也有计算逻辑。比如从网上拉一批 JSON 数据,然后对每份 JSON 做复杂特征计算。一个很自然的想法是把整块任务丢给进程池,让每个子进程自己发请求再计算。但这并不明智,因为网络请求在子进程里会占用进程资源,而且启动大量进程去做 I/O 等待,成本远高于线程。
我在做这类数据管道时,通常采用“线程池负责 I/O,进程池负责计算”的两层架构。简单说,主进程开一组线程去下载数据;拿到原始数据后,再提交给进程池做 CPU 密集处理。这样可以同时利用线程的 I/O 并发能力和进程的多核计算能力。代码上可以这样组织:
from concurrent.futures import ProcessPoolExecutor, ThreadPoolExecutor, as_completed def download(url): # 模拟网络 I/O,这里用线程池 return {"url": url, "data": range(1000)} def process(item): # 模拟 CPU 密集处理 return sum(item["data"]) ** 2 def run(urls): with ThreadPoolExecutor(max_workers=8) as thread_pool: with ProcessPoolExecutor(max_workers=4) as process_pool: futures = [] # 先并发下载 for downloaded in thread_pool.map(download, urls): # 每拿到一份数据就丢给进程池 futures.append(process_pool.submit(process, downloaded)) for future in as_completed(futures): print(future.result()) if __name__ == "__main__": sample_urls = [f"http://example.com/{i}" for i in range(20)] run(sample_urls)这种结构看起来很绕,其实核心思想很清晰:把“等”这件事交给线程池,把“算”这件事交给进程池。每个池的并发度可以独立调节,不会让某一边成为瓶颈。当然,这种架构下需要注意“下载结果”和“进程池任务”之间的数据传递。如果下载结果非常大,又要被 pickle 到多个进程,内存会迅速膨胀,所以要做好结果裁剪和批量分批。
5.2 进度展示与结果聚合:用 as_completed 而非 list 硬等
处理大批量任务时,给用户展示进度能让人安心不少。ProcessPoolExecutor 配合 as_completed,是天然适合做进度的方式。你不需要把所有任务都跑完才知道结果,而是每完成一个任务就更新一次计数,这样交互体验会好很多。
一个很实用的写法是维护 future 到任务编号的映射,通过 as_completed 判断哪个任务完成了,马上更新统计:
from concurrent.futures import ProcessPoolExecutor, as_completed item_map = {executor.submit(process_one, item): item for item in items} done_count = 0 for future in as_completed(item_map): item = item_map[future] try: result = future.result() except Exception as exc: log_error(item, exc) continue done_count += 1 if done_count % 50 == 0: print(f"进度: {done_count}/{len(items)}")这里有一个非常好用的小细节:如果 future 已经完成,调用 future.result(timeout=0) 不会阻塞,可以直接获取当前结果。所以你可以把进度展示逻辑放在一个定时循环里。如果任务之间有依赖,需要前面几个任务的结果才能提交后面的任务,那就不要一次性把全部任务提交完,改成分批 submit,或者只通过 done_callback 在回调里提交后续任务。
5.3 单机并发的天花板和之后的扩展方向
ProcessPoolExecutor 的边界很清楚:它只能利用一台机器的 CPU。当任务量大到单机 CPU、内存都不够时,你不会再执着于调大 max_workers,而是要考虑把任务拆分到多台机器上执行。那个阶段常见的方案是任务队列加 worker 服务,比如消息队列分发任务,或者使用 Ray、Dask 这类分布式计算框架。
但在跳到分布式之前,先用好 ProcessPoolExecutor 仍然是性价比最高的选择。很多所谓“慢到不能忍”的脚本,其实只是没有把机器的多核利用起来。你也不需要一次写很复杂的架构,只需要把任务函数定义成顶层可 pickle 的函数,用 with ProcessPoolExecutor 包住 submit 循环,再用 as_completed 汇总结果,你会发现一个很朴素的事实:Python 的 CPU 密集型并发,其实可以很简单。
我在实际项目中踩过几年坑之后,最想说的一点是:不要试图在主进程和子进程之间共享复杂可变对象,也不要把所有并发逻辑堆在一个巨型函数里。多进程并发最舒服的写法,永远是“小任务函数 + 主进程调度”,你提供的数据能 pickle,函数能在模块顶层被找到,然后剩下的交给标准库,它远比自己造轮子要可靠得多。