1. concurrent.futures 是什么? #

执行器 适用场景
ThreadPoolExecutor 网络请求、文件读写、数据库查询
ProcessPoolExecutor 数值计算、图像处理、数据分析

2. 前置知识:I/O 与 CPU #

顺序执行:  任务1 → 任务2 → 任务3 → ...  (慢)
并发执行:  任务1 ┐
            任务2 ├─ 同时进行,谁先完成先处理
            任务3 ┘

3. ThreadPoolExecutor #

# 导入线程池执行器
from concurrent.futures import ThreadPoolExecutor
# 导入 time 模拟网络延迟
import time

# 模拟网络请求函数
def fetch(url):
    # 模拟等待 1 秒
    time.sleep(1)
    # 返回完成信息
    return f"{url} 完成"

# 待请求的 URL 列表
urls = ["http://a.com", "http://b.com", "http://c.com"]

# 创建最多 3 个并发线程的线程池
with ThreadPoolExecutor(max_workers=3) as executor:
    # map 批量提交,按输入顺序返回结果
    for result in executor.map(fetch, urls):
        print(result)
# 使用 submit 逐个提交任务
with ThreadPoolExecutor(max_workers=3) as executor:
    # 为每个 URL 提交任务,得到 Future 列表
    futures = [executor.submit(fetch, url) for url in urls]
    # 遍历 Future,阻塞等待并获取结果
    for future in futures:
        print(future.result())

4. ProcessPoolExecutor #

# 导入进程池执行器
from concurrent.futures import ProcessPoolExecutor
# 导入 math 用于开方计算
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]

# Windows 上进程池必须放在 main 保护块内
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
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:
    # 提交所有任务,用字典关联 Future 与原始参数
    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. 如何选择执行器 #

场景 执行器 max_workers 建议
批量 HTTP 请求 ThreadPoolExecutor CPU 核心数 × 2~4
批量读文件 ThreadPoolExecutor CPU 核心数 × 2~4
批量数值计算 ProcessPoolExecutor CPU 核心数
图像批量处理 ProcessPoolExecutor CPU 核心数

7. 批量下载 #

# 导入线程池和 as_completed
from concurrent.futures import ThreadPoolExecutor, as_completed
# 导入 time 计时
import time
# 导入 random 模拟随机耗时
import random

# 模拟下载函数
def download(url, file_id):
    # 随机等待 0.5~1.5 秒
    time.sleep(random.uniform(0.5, 1.5))
    # 返回文件信息字典
    return {"file_id": file_id, "url": url, "size": random.randint(100, 1000)}

# 5 个待下载文件 (url, id)
files = [(f"http://example.com/f{i}.zip", i) for i in range(1, 6)]
# 记录开始时间
start = time.time()

# 最多 3 个并发下载
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:.2f} 秒")

8. 常见错误与注意事项 #

# 导入 os 获取 CPU 核心数
import os

# 获取 CPU 核心数,获取失败时默认 4
cpu_count = os.cpu_count() or 4
# I/O 密集型建议的 worker 数
io_workers = min(32, cpu_count * 4)
# CPU 密集型建议的 worker 数
cpu_workers = cpu_count

9. 总结 #

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 最佳实践 #