Learn
Python/14-concurrency

并发编程

Python 有三套并发生态,它们的适用场景完全不同:

工具适合关键限制
threadingI/O 密集(网络、文件、DB)GIL 让 CPU 密集任务跑不快
multiprocessingCPU 密集(计算、压缩、图像)进程间通信要序列化,开销大
asyncio高并发 I/O(成千上万的连接)需要代码是 async 协程

本章的"网络"示例都使用模拟数据,因为本环境的 Docker 沙箱禁用了网络。这样反而能专注于并发模式本身。

1. threading —— 线程

线程适合 I/O 等待多的任务。多个线程在 I/O 阻塞时轮流使用 CPU。

基础线程
import threading, time
 
def worker(name, secs):
    print(f"  [{name}] 开始")
    time.sleep(secs)            # 模拟 I/O
    print(f"  [{name}] 结束")
 
start = time.perf_counter()
threads = [threading.Thread(target=worker, args=(f"T{i}", 0.2)) for i in range(3)]
for t in threads: t.start()
for t in threads: t.join()      # 等待所有线程结束
print(f"总耗时: {time.perf_counter() - start:.2f}s")
ℹ️GIL

全局解释器锁(GIL)保证同一时刻只有一个线程执行 Python 字节码。所以CPU 密集任务用多线程没意义——但 I/O 阻塞时 GIL 会被释放,其它线程可以跑。

2. Lock —— 线程安全

多线程改共享数据时需要加锁:

Lock 保护共享计数
import threading
 
counter = 0
lock = threading.Lock()
 
def inc(n):
    global counter
    for _ in range(n):
        with lock:        # 等价于 acquire()/release()
            counter += 1
 
threads = [threading.Thread(target=inc, args=(10_000,)) for _ in range(5)]
for t in threads: t.start()
for t in threads: t.join()
print("counter =", counter, "(期望 50000)")

3. ThreadPoolExecutor —— 线程池

手写线程管理很烦,concurrent.futures 帮你做:

ThreadPoolExecutor
import time
from concurrent.futures import ThreadPoolExecutor, as_completed
 
def fetch(name, secs):
    time.sleep(secs)
    return f"{name}={secs}s"
 
start = time.perf_counter()
with ThreadPoolExecutor(max_workers=3) as pool:
    futures = [pool.submit(fetch, f"req{i}", 0.2) for i in range(3)]
    for fut in as_completed(futures):
        print("  完成:", fut.result())
print(f"总耗时: {time.perf_counter() - start:.2f}s")
💡as_completed vs map
  • as_completed(fs):哪个先完成就 yield 哪个,适合"逐个处理结果"
  • pool.map(fn, items):按入参顺序返回结果,适合"批量转换"

4. multiprocessing —— 多进程

CPU 密集任务用多进程绕过 GIL:

多进程计算
import multiprocessing as mp
import time
 
def cpu_heavy(n):
    s = 0
    for i in range(n):
        s += i * i
    return s
 
N = 2_000_000
 
start = time.perf_counter()
results = [cpu_heavy(N) for _ in range(4)]
print(f"串行: {time.perf_counter() - start:.2f}s, 结果示例: {results[0]}")
 
start = time.perf_counter()
with mp.Pool(processes=4) as pool:
    results = pool.map(cpu_heavy, [N] * 4)
print(f"并行: {time.perf_counter() - start:.2f}s, 结果示例: {results[0]}")
⚠️沙箱说明

本沙箱用 multiprocessing 的默认 fork 方式。在 Jupyter / Windows / 某些容器里可能要换成 spawn。如果并行跑不出来,多半是 if __name__ == '__main__': 没加。

5. asyncio —— 高并发 I/O(第 11 章的延伸)

当并发量非常大(几千上万个连接)时,线程的开销就吃不消了。asyncio 用单线程 + 事件循环撑住海量连接。

asyncio 高并发
import asyncio, time
 
async def fetch(i):
    await asyncio.sleep(0.1)         # 模拟一次 I/O
    return f"result-{i}"
 
async def main():
    start = time.perf_counter()
    # 同时发起 50 个"请求"
    results = await asyncio.gather(*(fetch(i) for i in range(50)))
    print(f"完成 {len(results)} 个请求")
    print(f"总耗时: {time.perf_counter() - start:.2f}s(串行要 5s)")
    print("前 3 个:", results[:3])
 
asyncio.run(main())

6. ProcessPoolExecutor

concurrent.futures 也有进程池版本,API 与 ThreadPoolExecutor 完全一致:

ProcessPoolExecutor
from concurrent.futures import ProcessPoolExecutor
 
def cube(x):
    return x ** 3
 
with ProcessPoolExecutor(max_workers=4) as pool:
    results = list(pool.map(cube, range(10)))
print(results)

7. 怎么选?

任务类型推荐方案
100 个 HTTP 请求asyncio 或 ThreadPoolExecutor
上万个长连接(WebSocket)asyncio
CPU 重活(加密、图像处理)ProcessPoolExecutor / multiprocessing
简单脚本偶尔并发几个 I/OThreadPoolExecutor 最省事
既要 I/O 又要 CPU各自处理:进程里跑事件循环,或线程里跑进程池
💡一个反直觉的事实

I/O 等待多时,asyncio 比 threading 更快,是因为它没有线程切换的开销。一个 asyncio 程序可以轻松撑 1 万个并发连接,而开 1 万个线程内存就爆了。

🎯 练习

实现一个"并行抓取"工具:接收 10 个 URL,模拟每个 0.1s 的网络延迟,用 asyncio.gather 并发跑,统计总耗时。

并行抓取
import asyncio, time
 
async def fetch(url):
    # 模拟一次网络请求
    pass
 
async def main():
    urls = [f"https://api.example.com/items/{i}" for i in range(10)]
    start = time.perf_counter()
    # 在这里并发抓取,期望 ~0.1s 而非 1.0s
    pass
 
asyncio.run(main())
🎯提示

results = await asyncio.gather(*(fetch(u) for u in urls))。

小结

  • ✅ threading 适合 I/O 密集,受 GIL 限制不适合 CPU 密集
  • ✅ Lock 保护共享数据;with lock: 比手写 acquire/release 安全
  • ✅ ThreadPoolExecutor / ProcessPoolExecutor 统一了 API
  • ✅ asyncio 适合超高并发 I/O,开销远小于线程
  • ✅ CPU 密集 → multiprocessing / ProcessPoolExecutor
  • ✅ 三套工具各有适用场景,没有银弹

下一章 单元测试:用 unittest / pytest 写出可信赖的代码。