并发编程
Python 有三套并发生态,它们的适用场景完全不同:
| 工具 | 适合 | 关键限制 |
|---|---|---|
threading | I/O 密集(网络、文件、DB) | GIL 让 CPU 密集任务跑不快 |
multiprocessing | CPU 密集(计算、压缩、图像) | 进程间通信要序列化,开销大 |
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 —— 线程安全
多线程改共享数据时需要加锁:
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 帮你做:
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 用单线程 + 事件循环撑住海量连接。
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 完全一致:
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/O | ThreadPoolExecutor 最省事 |
| 既要 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 写出可信赖的代码。