在 16.2 和 16.3 节我们已经分别探讨了 threading 和 multiprocessing 模块,它们提供了细粒度的并发控制能力,但同时也要求开发者手动管理线程/进程的创建、回收以及返回值收集。concurrent.futures 模块则在这些底层接口之上封装了一套高层抽象,通过线程池和进程池以统一的方式执行异步任务,让我们写并发代码就像调用普通函数一样简单。
核心组件
- Executor:抽象基类,定义了
submit、map、shutdown等核心方法。实际使用时我们选择它的两个实现类: ThreadPoolExecutor:线程池,适合 IO 密集型任务。ProcessPoolExecutor:进程池,适合 CPU 密集型任务(绕过 GIL)。- Future:代表一个异步任务未来的结果。你可以用它查询任务是否完成、获取返回值或捕获异常。任务提交后立即返回一个
Future对象,主线程可以继续做其他事情。
统一接口的两个关键方法
- *
submit(fn,args, kwargs)
向池中提交一个可调用对象及参数,返回一个 Future。适用于任务数量不多、需要精细控制每个任务结果的场景。
map(fn, iterable, timeout=None, chunksize=1)
类似内置的 map 函数,将 iterable 中的每个元素作为参数调用 fn,返回一个生成器,按提交顺序逐个产出结果。它会阻塞直到当前项完成。适合批量处理且关心返回顺序的场景。chunksize 参数仅对进程池生效:将多个任务打包发送给进程,减少通信开销。
基本用法示例
线程池(IO 密集型:批量下载页面)
import concurrent.futures
import requests
URLS = ['https://example.com', 'https://python.org', 'https://github.com']
def fetch(url):
resp = requests.get(url, timeout=10)
return f"{url}: {resp.status_code}"
with concurrent.futures.ThreadPoolExecutor(max_workers=5) as executor:
future_to_url = {executor.submit(fetch, url): url for url in URLS}
for future in concurrent.futures.as_completed(future_to_url):
url = future_to_url[future]
try:
result = future.result()
print(result)
except Exception as e:
print(f"{url} 出错: {e}")
- 用
with语句管理池的生命周期,结束时自动调用shutdown(wait=True)等待所有任务完成。 as_completed可以按任务完成顺序处理,不必等待前面的任务。
进程池(CPU 密集型:计算斐波那契数)
import concurrent.futures
def fib(n):
a, b = 0, 1
for _ in range(n):
a, b = b, a + b
return a
if __name__ == '__main__':
with concurrent.futures.ProcessPoolExecutor(max_workers=4) as executor:
tasks = [10, 20, 30, 40]
# 使用 map 按顺序获取结果
for n, result in zip(tasks, executor.map(fib, tasks)):
print(f"fib({n}) = {result}")
- 必须把主程序放在
if name == 'main':保护下,否则进程池会在 Windows 上无限递归派生子进程。 map返回顺序与输入顺序一致,即使某个任务先完成,也要等前面的任务结果出来后才产出。
为什么要用线程池/进程池?
- 资源复用:避免频繁创建和销毁线程/进程的巨大开销。池中维持一定数量的工作者,任务结束后仍保留,下次再用。
- 限流保护:通过
max_workers控制并发数上限,避免无节制地创建线程压垮系统,或进程数过多导致 CPU 上下文切换激增。 - 接口统一:无论在线程还是进程中执行任务,代码结构几乎一样,切换只需要改一个类名,大幅降低心智负担。
与原生 threading/multiprocessing 的对比
| 特性 | concurrent.futures | threading / multiprocessing |
|------|-------------------|-----------------------------|
| 手动创建工作者 | 无需 | 需要 Thread/Process 实例 |
| 返回值获取 | Future 直接 result() | 需通过队列 Queue 传递 |
| 异常传递 | 在 result() 时抛出 | 需自己捕捉并通过队列传递 |
| 批量任务 | submit+循环 或 map | 需自己实现任务分发和管理 |
| 执行顺序控制 | as_completed 方便按完成顺序处理 | 需额外编码 |
结论:几乎所有“需要并发执行一批任务”的场景,都可以优先使用 concurrent.futures,除非你需要更底层的控制(如设置线程优先级、守护线程、进程间共享状态等)。
注意事项
- 进程池的序列化要求:提交给进程池的函数及其参数必须是可 pickle 的(普通函数、全局函数、类实例方法等没问题,但
lambda或绑定外部变量的闭包可能 pickle 失败)。 - 线程池的 GIL 限制:CPU 密集型任务在线程池中无法利用多核,应改用进程池。但线程池的创建成本低,更适合 IO 密集任务。
- 善用
with:保证池最终被关闭,避免资源泄漏。若不用with,需手动调用executor.shutdown(wait=True)。 max_workers调优:线程池可适当设为 CPU 核数的几倍(IO 密集时),进程池一般等于os.cpu_count()或略少,避免过度竞争。
concurrent.futures 让并发编程从“手工管理线程/进程”升级为“把任务扔到池子里就好”,是日常开发中最常用的并发利器,也是连接简单并发和后面将要学习的协程/异步 IO 的桥梁。