进程池与线程池
约 1461 字大约 5 分钟
2026-05-10
手动创建线程或进程适合学习底层机制,但实际项目中,更常用的是池。池会提前维护一组工作线程或工作进程,把任务提交进去执行。
- ThreadPoolExecutor
- submit 与 Future
- ascompleted
- 异常处理
- 1用 ThreadPoolExecutor(max_workers=10) 并发请求一批 URL,几乎是手写 Thread 的精简版。
- 2对比 executor.map(func, items) 和 executor.submit(func, x) + as_completed:前者顺序结果、后者完成顺序。
- 3future.result(timeout=3) 给单个任务设超时;future.cancel() 取消未开始的任务。
- 4异常会在 future.result() 时重新抛出 —— 一定要包 try/except,否则后台静默失败。
- 5ProcessPoolExecutor 适合 CPU 密集任务;submit 传的函数和参数必须能 pickle。
手动创建线程或进程适合学习底层机制,但实际项目中,更常用的是池。池会提前维护一组工作线程或工作进程,把任务提交进去执行。
Python 标准库 concurrent.futures 提供了统一接口:
ThreadPoolExecutorProcessPoolExecutor
ThreadPoolExecutor
线程池适合 IO 密集型任务。
from concurrent.futures import ThreadPoolExecutor
import time
def download(url):
time.sleep(1)
return f'下载完成:{url}'
urls = ['a.com', 'b.com', 'c.com']
with ThreadPoolExecutor(max_workers=3) as executor:
results = executor.map(download, urls)
for result in results:
print(result)max_workers 控制最大并发线程数。
submit 与 Future
submit() 会提交任务并返回 Future 对象。
from concurrent.futures import ThreadPoolExecutor
def square(n):
return n * n
with ThreadPoolExecutor(max_workers=2) as executor:
future = executor.submit(square, 5)
print(future.result())Future 表示一个未来才会完成的结果。调用 result() 会等待任务完成并返回结果。
as_completed
as_completed() 可以按任务完成顺序处理结果。
from concurrent.futures import ThreadPoolExecutor, as_completed
import time
def task(n):
time.sleep(n)
return n
with ThreadPoolExecutor(max_workers=3) as executor:
futures = [executor.submit(task, n) for n in [3, 1, 2]]
for future in as_completed(futures):
print(future.result())输出顺序大概率是 1, 2, 3,因为它按完成顺序返回。
异常处理
线程或进程中的异常会在调用 future.result() 时重新抛出。
from concurrent.futures import ThreadPoolExecutor, as_completed
def task(n):
if n == 2:
raise ValueError('出错了')
return n
with ThreadPoolExecutor(max_workers=3) as executor:
futures = [executor.submit(task, n) for n in [1, 2, 3]]
for future in as_completed(futures):
try:
print(future.result())
except ValueError as exc:
print(f'任务失败:{exc}')ProcessPoolExecutor
进程池适合 CPU 密集型任务。
from concurrent.futures import ProcessPoolExecutor
def calculate(n):
return sum(i * i for i in range(n))
if __name__ == '__main__':
numbers = [10_000_000, 10_000_000, 10_000_000]
with ProcessPoolExecutor(max_workers=3) as executor:
for result in executor.map(calculate, numbers):
print(result)使用进程池时,同样建议把入口放在 if __name__ == '__main__': 下面。
map 与 submit 的区别
| 方法 | 特点 |
|---|---|
map() | 简洁,按输入顺序返回结果 |
submit() | 灵活,返回 Future,方便单独处理异常、取消、超时 |
简单批量任务用 map();需要精细控制用 submit()。
超时控制
future.result(timeout=3)如果任务 3 秒内没有完成,会抛出 TimeoutError。
选择线程池还是进程池
- 网络请求、文件 IO、数据库查询:线程池
- 大量计算、图片处理、压缩加密:进程池
- 不确定任务是否 CPU 密集:先测量,不要猜
注意事项
- 池的大小不是越大越好。
- 线程池适合 IO 密集型;进程池适合 CPU 密集型。
- 进程池中的任务函数和参数需要能被序列化。
- 使用
with可以自动关闭池。 - 处理结果时要捕获
future.result()抛出的异常。
为什么推荐先学 Executor
直接管理 Thread 或 Process 容易写出很多样板代码:创建、启动、join、收集结果、处理异常。concurrent.futures 把这些操作包装成统一接口。
你只需要先理解两个概念:
Executor:负责管理一组线程或进程;Future:代表一个还没完成或已经完成的任务结果。
ThreadPoolExecutor:请求多个网页
from concurrent.futures import ThreadPoolExecutor, as_completed
import requests
urls = [
'https://httpbin.org/delay/1',
'https://httpbin.org/get',
'https://httpbin.org/uuid',
]
def fetch(url):
response = requests.get(url, timeout=5)
response.raise_for_status()
return url, response.status_code
with ThreadPoolExecutor(max_workers=3) as executor:
futures = [executor.submit(fetch, url) for url in urls]
for future in as_completed(futures):
url, status = future.result()
print(url, status)这个例子适合 IO 密集型任务:每个线程大部分时间在等网络返回。
ProcessPoolExecutor:计算多个任务
from concurrent.futures import ProcessPoolExecutor
def count(n):
total = 0
for i in range(n):
total += i * i
return total
if __name__ == '__main__':
nums = [10_000_000, 12_000_000, 14_000_000]
with ProcessPoolExecutor(max_workers=3) as executor:
results = list(executor.map(count, nums))
print(results)这个例子适合 CPU 密集型任务。注意多进程依然建议放在 if __name__ == '__main__': 下面。
submit 和 map 怎么选
map 更适合一组同类任务,并且你希望按输入顺序拿结果:
results = executor.map(func, items)submit 更灵活,适合:
- 每个任务参数不完全一样;
- 想用
as_completed谁先完成先处理; - 想分别处理每个任务的异常;
- 想给 Future 建立额外映射关系。
常见写法:
future_to_item = {
executor.submit(handle, item): item
for item in items
}
for future in as_completed(future_to_item):
item = future_to_item[future]
try:
result = future.result()
except Exception as exc:
print('任务失败', item, exc)
else:
print('任务成功', item, result)max_workers 怎么定
没有一个永远正确的数字。可以先用保守值:
- 网络 IO:5~20 起步,看接口限制和机器情况;
- 文件 IO:不要太高,避免磁盘抖动;
- CPU 计算:接近 CPU 核心数;
- 调用外部服务:以对方限流规则为准。
如果你不确定,先小一点。并发过高时,程序可能不是更快,而是更容易失败。
总结
concurrent.futures 是 Python 并发编程中非常实用的高级接口。它屏蔽了线程和进程管理细节,让我们把注意力放在任务提交、结果获取和异常处理上。
- Executor 让线程池和进程池有统一用法,Future 用来拿结果和异常。
- IO 密集型任务优先 ThreadPoolExecutor,CPU 密集型任务考虑 ProcessPoolExecutor。
- `map` 简洁但不灵活,`submit + as_completed` 更适合逐个处理结果和异常。
版权所有
版权归属:Shuo Liu
