上个月,一个做数据清洗的朋友找我帮忙看脚本。他有三千多个 CSV 文件要处理,每个文件做一遍清洗、过滤、导出。单进程跑要二十多分钟,他听说 multiprocessing.Pool 能并行,就改成了这样:
from multiprocessing import Pool
def process_file(filepath):
# 读取、清洗、导出
return result
if __name__ == '__main__':
files = [...] # 三千多个路径
pool = Pool(4)
for f in files:
pool.apply(process_file, args=(f,))
pool.close()
pool.join()
他兴冲冲地跑了一遍,结果还是二十多分钟。他以为是进程池没生效,把 Pool(4) 改成 Pool(8),依然没变化。他怀疑是不是文件太小,进程切换开销太大,甚至开始怀疑 Python 的多进程是假的。
我让他把 apply 换成 apply_async,再收集一下结果。改完一跑,时间直接降到六分钟。
他盯着屏幕说:“就改了一个词,差别这么大?”
是的,差别就是这么大。问题不在于进程池有没有生效,而在于 apply 根本就不是用来做批量并发的。
apply 是阻塞的,一次只能等一个
Pool.apply 的行为非常直接:提交一个任务,然后主进程停在这里等,直到这个任务完成,返回结果,才继续往下走。
你可以把它理解成“同步调用”。虽然任务确实是在池子里的某个子进程中执行的,但主进程在等待期间什么也做不了,也不会去提交下一个任务。所以当你写 for f in files: pool.apply(process_file, args=(f,)) 的时候,实际流程是:
- 提交文件 1,主进程等待文件 1 处理完。
- 文件 1 处理完,返回结果,主进程继续循环。
- 提交文件 2,主进程等待文件 2 处理完。
- 以此类推。
池子里有 4 个进程,但同一时间只有一个在干活,另外三个全程空闲。你看到的“并行”其实一点都没发生,整体耗时和单进程循环几乎一样。
apply 的设计目的不是批量处理,而是“我只需要提交一个任务,并且我立刻就要它的结果”。比如你在程序某个分支里需要算一个很重的函数,但只算一次,用 apply 是合适的。一旦进入循环,它就成了性能陷阱。
apply_async 才是并发的入口
apply_async 的名字里带 async,意思很明确:提交任务后不等待,立即返回一个 AsyncResult 对象。任务被丢进池子的队列里,由子进程异步执行。主进程可以马上继续提交下一个任务。
改成这样:
if __name__ == '__main__':
files = [...]
pool = Pool(4)
results = []
for f in files:
res = pool.apply_async(process_file, args=(f,))
results.append(res)
pool.close()
pool.join()
for res in results:
print(res.get())
循环会飞快地跑完,三千个任务几乎瞬间全部提交到池子里。四个子进程开始并行消费队列,真正实现了四路并发。主进程在 pool.join() 处等待所有任务完成,然后通过 res.get() 逐个取出结果。
这里的关键是:提交和执行是分离的。提交动作只是把任务放进队列,不涉及计算,所以极快。真正的计算发生在子进程中,多个子进程同时从队列里取任务,互不阻塞。
AsyncResult 是什么
apply_async 返回的 AsyncResult 是一个句柄,代表“一个将来才会有的结果”。它提供几个核心方法:
get(timeout=None):阻塞等待任务完成,返回函数返回值。如果任务抛了异常,get()会重新抛出这个异常。可以设置timeout,超时抛TimeoutError。ready():非阻塞检查任务是否完成。successful():任务是否成功完成(没有抛异常)。wait(timeout=None):等待任务完成,但不返回结果。
有了这些方法,你可以在提交完所有任务之后,先做点别的事,然后再统一收集结果。也可以边提交边处理已完成的任务,灵活性远高于 apply。
忘了 get 和 join,任务可能悄悄消失
我见过另一个极端:有人用了 apply_async,但提交完就直接让程序结束了。
pool = Pool(4)
for f in files:
pool.apply_async(process_file, args=(f,))
# 没有 close,没有 join,主进程直接退出
主进程退出时,Python 会尝试清理子进程。如果池还没有关闭,正在执行的任务可能被强制终止,队列里还没开始的任务直接丢失。你看到程序正常结束了,但实际只处理了一部分文件,而且没有任何报错。
正确的收尾动作是:
pool.close() # 停止接受新任务
pool.join() # 等待所有已提交任务完成
或者用 with 语句:
with Pool(4) as pool:
results = [pool.apply_async(process_file, (f,)) for f in files]
# 退出 with 块时会自动调用 close 和 join
但要注意,with 块退出时只是等待任务完成,并不会自动调用 get()。如果任务抛了异常,异常会被保存在 AsyncResult 里,不会主动冒出来。你仍然需要在 with 块内部或外部调用 get() 来捕获异常。
异常不会自己跳出来
apply 是同步的,函数一抛异常,主进程立刻就能捕获。但 apply_async 不同,异常发生在子进程里,主进程不知道。只有当你调用 res.get() 的时候,异常才会被重新抛出。
如果你只 join() 不 get(),异常就被静默吞掉了。程序看起来正常结束,但结果文件少了几个,你完全不知道。
所以收集结果时一定要处理异常:
for res in results:
try:
data = res.get()
except Exception as e:
print(f"任务失败: {e}")
如果想在任务失败时立即得到通知,可以用 error_callback:
def on_error(e):
print(f"出错了: {e}")
pool.apply_async(process_file, (f,), error_callback=on_error)
但注意,回调函数是在主进程的一个内部线程中执行的,不是子进程。如果回调里做了耗时操作,会阻塞结果收集。回调里也不应该抛出异常,否则会被忽略。
map 和 apply_async 的关系
很多人会问:那 pool.map 呢?它是不是也是并发的?
pool.map(func, iterable) 是阻塞的。它会自动把 iterable 拆成多个任务,提交到池子里,然后等待所有任务完成,返回一个结果列表。内部确实用了并发,但主进程在 map 调用处一直等到全部结束。所以它适合“我有一批数据,要并行处理,并且我立刻就要全部结果”的场景。
map_async 则是非阻塞版本,返回 AsyncResult,用法和 apply_async 类似。
还有 imap 和 imap_unordered,它们返回迭代器,按顺序(或按完成顺序)产出结果,适合处理大批量数据时边完成边消费,避免一次性把所有结果堆在内存里。
简单对比:
apply:同步,一次一个,返回直接结果。apply_async:异步,一次一个,返回AsyncResult。map:同步,批量,返回结果列表。map_async:异步,批量,返回AsyncResult。imap:异步,批量,返回有序迭代器。imap_unordered:异步,批量,返回无序迭代器。
做批量任务时,apply 基本不该出现。要么用 apply_async 手动管理并发,要么用 map 系列让池子帮你调度。
Windows 上的额外一坑
如果你在 Windows 上跑多进程,所有涉及 Pool 的代码必须放在 if __name__ == '__main__': 下面。因为 Windows 用 spawn 方式启动子进程,子进程会重新导入主模块。如果没有这个保护,子进程导入时又会执行创建 Pool 的代码,导致无限递归创建进程,直接报错或卡死。
Linux 和 macOS 默认用 fork,不会重新导入,所以很多人本地测试没问题,一到 Windows 服务器就炸。这个坑和 apply 无关,但用进程池时几乎都会遇到,值得单独记住。
怎么选:一张表说清楚
| 方法 | 阻塞? | 返回值 | 适用场景 |
apply |
是 | 直接结果 | 单个任务,立即要结果 |
apply_async |
否 | AsyncResult |
批量任务,需要并发,灵活控制 |
map |
是 | 结果列表 | 批量任务,顺序无关,一次性拿全部结果 |
map_async |
否 | AsyncResult |
批量任务,异步,后续统一收集 |
imap |
否 | 迭代器(有序) | 大批量,边完成边消费,保持顺序 |
imap_unordered |
否 | 迭代器(无序) | 大批量,只关心完成,不关心顺序 |
回到我朋友那个场景:三千个文件,每个独立处理,不需要顺序,他只需要最后统计成功和失败的数量。最适合的其实是 imap_unordered:
with Pool(4) as pool:
for result in pool.imap_unordered(process_file, files):
# 每完成一个就处理一个
...
这样内存占用小,结果实时产出,异常也能在迭代时捕获。如果他需要保持顺序,就用 imap。
最后一句
apply 和 apply_async 的区别,本质上是同步等待和异步提交的区别。apply 把并发池用成了串行队列,apply_async 才真正释放了多进程的并行能力。
但异步也带来了新的责任:你必须负责收集结果、处理异常、正确关闭池子。忘了任何一步,任务就可能悄悄失败或丢失。
所以下次写进程池的时候,先问自己一句:我是要“提交一个任务然后等结果”,还是要“提交一堆任务让它们并行跑”?如果是后者,就别再用 apply 了。