1. concurrent.futures 是什么? #
concurrent.futures 是 Python 3.2+ 标准库,提供线程池 和进程池 的高级并发接口。
比手动 threading.Thread 更简洁:自动管理线程/进程创建、复用和销毁,用 with 语句即可。
两个核心执行器:ThreadPoolExecutor(I/O 密集)和 ProcessPoolExecutor(CPU 密集)。
提交任务后返回 Future 对象,用于获取结果、处理异常和设置超时。
执行器
适用场景
ThreadPoolExecutor
网络请求、文件读写、数据库查询
ProcessPoolExecutor
数值计算、图像处理、数据分析
2. 前置知识:I/O 与 CPU #
I/O 密集型 :大部分时间在等待(网络、磁盘),适合线程池,等待时释放 GIL。
CPU 密集型 :大部分时间在计算,适合进程池,每个进程有独立 GIL,可真正并行。
GIL 详情见 13.threading.md;简单记:等 I/O 用线程池,算数据用进程池 。
100 个各需 1 秒的 I/O 任务,10 并发线程理论上约 10 秒完成,顺序执行需 100 秒。
顺序执行: 任务1 → 任务2 → 任务3 → ... (慢)
并发执行: 任务1 ┐
任务2 ├─ 同时进行,谁先完成先处理
任务3 ┘3. ThreadPoolExecutor #
用 with ThreadPoolExecutor(max_workers=N) as executor: 创建线程池,with 结束自动关闭。
max_workers 控制最大并发线程数,I/O 任务可设为 CPU 核心数的 2–4 倍。
executor.map(func, items) 批量提交,结果保持输入顺序 。
executor.submit(func, arg) 提交单个任务,返回 Future,灵活度更高。
from concurrent.futures import ThreadPoolExecutor
import time
def fetch (url ):
time.sleep(1 )
return f"{url} 完成"
urls = ["http://a.com" , "http://b.com" , "http://c.com" ]
with ThreadPoolExecutor(max_workers=3 ) as executor:
for result in executor.map (fetch, urls):
print (result)
with ThreadPoolExecutor(max_workers=3 ) as executor:
futures = [executor.submit(fetch, url) for url in urls]
for future in futures:
print (future.result())4. ProcessPoolExecutor #
用法与 ThreadPoolExecutor 几乎相同,换成 ProcessPoolExecutor 即可。
默认 max_workers 等于 CPU 核心数,适合充分利用多核做计算。
被执行的函数和参数必须可被 pickle 序列化 ,不能传文件句柄、数据库连接等。
每个进程有独立内存,修改参数不影响原数据;进程创建开销比线程大。
from concurrent.futures import ProcessPoolExecutor
import math
def is_prime (n ):
if n < 2 :
return False
for i in range (2 , int (math.sqrt(n)) + 1 ):
if n % i == 0 :
return False
return True
numbers = [112272535095293 , 112582705942171 ]
if __name__ == "__main__" :
with ProcessPoolExecutor() as executor:
results = list (executor.map (is_prime, numbers))
print (list (zip (numbers, results)))注意: Windows 上进程池代码必须放在 if __name__ == "__main__": 保护块内。
5. as_completed 与错误处理 #
as_completed(futures) 按完成顺序 (而非提交顺序)返回 Future,适合"谁先好谁先处理"。
任务中的异常不会立即抛出,调用 future.result() 时才会抛出,必须用 try/except 捕获。
future.result(timeout=秒) 可设置超时,超时抛出 TimeoutError。
批量任务中一个失败不应影响其他任务,逐个 result() 并分别处理异常。
from concurrent.futures import ThreadPoolExecutor, as_completed
def safe_divide (x, y ):
return x / y
tasks = [(10 , 2 ), (20 , 4 ), (30 , 0 ), (40 , 5 )]
with ThreadPoolExecutor() as executor:
futures = {executor.submit(safe_divide, x, y): (x, y) for x, y in tasks}
for future in as_completed(futures):
x, y = futures[future]
try :
print (f"{x} /{y} = {future.result()} " )
except ZeroDivisionError:
print (f"{x} /{y} = 错误:除数为零" )6. 如何选择执行器 #
网络请求、API 调用、文件读写、数据库查询 → ThreadPoolExecutor。
数学运算、图像处理、机器学习推理、大数据计算 → ProcessPoolExecutor。
不确定时问自己:任务大部分时间在等 还是在算 ?
max_workers 参考:I/O 任务 min(32, cpu_count * 4),CPU 任务 cpu_count。
场景
执行器
max_workers 建议
批量 HTTP 请求
ThreadPoolExecutor
CPU 核心数 × 2~4
批量读文件
ThreadPoolExecutor
CPU 核心数 × 2~4
批量数值计算
ProcessPoolExecutor
CPU 核心数
图像批量处理
ProcessPoolExecutor
CPU 核心数
7. 批量下载 #
下面综合 ThreadPoolExecutor + submit + as_completed + 异常处理,是项目中的典型模式。
用字典 {future: 元数据} 关联 Future 与业务信息,完成时通过 future 反查来源。
记录开始/结束时间可直观看到并发带来的提速效果。
实际项目中将 time.sleep 替换为 requests.get(url) 即可用于真实下载。
from concurrent.futures import ThreadPoolExecutor, as_completed
import time
import random
def download (url, file_id ):
time.sleep(random.uniform(0.5 , 1.5 ))
return {"file_id" : file_id, "url" : url, "size" : random.randint(100 , 1000 )}
files = [(f"http://example.com/f{i} .zip" , i) for i in range (1 , 6 )]
start = time.time()
with ThreadPoolExecutor(max_workers=3 ) as executor:
futures = {executor.submit(download, url, fid): fid for url, fid in files}
for future in as_completed(futures):
fid = futures[future]
try :
r = future.result()
print (f"文件 {r['file_id' ]} 完成,{r['size' ]} KB" )
except Exception as e:
print (f"文件 {fid} 失败: {e} " )
print (f"总耗时: {time.time() - start:.2 f} 秒" )8. 常见错误与注意事项 #
CPU 密集任务用线程池不会加速,GIL 限制了并行计算,应改用进程池。
ProcessPoolExecutor 的函数必须在模块顶层定义,且参数可 pickle;Windows 需 if __name__ == "__main__":。
忘记对 future.result() 做异常处理,会导致一个任务失败中断整个流程。
始终用 with ThreadPoolExecutor(...) as executor: 确保资源正确释放。
import os
cpu_count = os.cpu_count() or 4
io_workers = min (32 , cpu_count * 4 )
cpu_workers = cpu_count9. 总结 #
concurrent.futures 是项目并发首选:比手写 threading 简洁,比 asyncio 上手快。
核心 API:map()(批量有序)、submit()(单个灵活)、as_completed()(完成顺序)、future.result()(取结果)。
I/O 用 ThreadPoolExecutor,CPU 用 ProcessPoolExecutor,配合 with 和异常处理。
底层线程概念见 13.threading.md,网络 I/O 实战见 15.urllib.md。
9.1 速查 #
from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor, as_completed
with ThreadPoolExecutor(max_workers=4 ) as ex:
results = list (ex.map (func, items))
with ThreadPoolExecutor() as ex:
future = ex.submit(func, arg)
result = future.result(timeout=10 )
with ThreadPoolExecutor() as ex:
futures = [ex.submit(func, x) for x in items]
for f in as_completed(futures):
print (f.result())9.2 最佳实践 #
I/O 并发优先 ThreadPoolExecutor,不要为 CPU 计算开很多线程
批量任务用 as_completed + try/except,单个失败不影响整体
进程池函数放模块顶层,Windows 加 if __name__ == "__main__":
合理设置 max_workers,过多线程反而因切换开销变慢