线程、进程与 asyncio
从 Promise 和 event loop 出发,理解 Python 任务、线程和进程的适用场景。
学习目标:为工作负载选择并发模型
本节结束时,你能从 JavaScript/TypeScript 的 Promise 和 event loop 迁移到 Python 的 coroutine、Task、线程池和进程池,并根据“等待还是计算”选择工具。你会理解 GIL 的实际边界,知道任务为什么必须有创建、等待、失败、取消和清理的生命周期,还能用 semaphore 或有限队列实现背压,避免 AI 服务、数据库和 GPU 被一口气压垮。
并发不是把代码写得更短,也不是把所有工作都丢进 gather。它改变了完成顺序、错误传播、共享状态和资源使用方式。阅读下面的代码时,始终问三个问题:任务什么时候真正开始;谁收集它的异常;同时在途的任务上限在哪里。
从 JS/TS 迁移的心智模型:Promise.all 不等于无限资源
JS/TS 的 Promise.all 会等待一组 promise,并在一个拒绝时把错误交给调用者;Python 的 asyncio.gather 对 coroutine 做类似的收集,但 coroutine 对象本身只是“可运行的描述”,不是已经启动的后台线程。asyncio.create_task 才会把协程交给事件循环调度。两者都不会自动限制请求数量,也不会让同步 CPU 代码变成并行计算。
const results = await Promise.all(
urls.map(url => fetchJson(url)),
); async def load_all(urls):
tasks = [asyncio.create_task(fetch_json(url)) for url in urls]
return await asyncio.gather(*tasks) 1. 协程、Task 与任务生命周期
直接调用 fetch_json(url) 得到 coroutine;await 会在当前协程中运行并等待它完成;create_task 会把它注册到当前事件循环,当前协程可以继续做别的事。注册任务就产生了责任:正常路径要 await 它,失败路径要读取异常,取消路径要等待清理。只创建任务而不保存引用,会让异常变成难以定位的“Task exception was never retrieved”。
import asyncio
async def run_batch(items: list[str]):
tasks = [asyncio.create_task(fetch_one(item), name=f"fetch:{item}")
for item in items]
try:
return await asyncio.gather(*tasks)
except BaseException:
for task in tasks:
if not task.done():
task.cancel()
await asyncio.gather(*tasks, return_exceptions=True)
raise
gather 默认把第一个异常传出来,但“其他任务如何处理”需要由你的业务决定:批量 embedding 可以保留部分成功并逐项记录错误;事务型写入可能需要取消全部任务。return_exceptions=True 会把异常放进结果列表,它不是吞错开关,调用者仍要逐项检查异常类型并关联输入 ID。
2. GIL 的直觉:线程能等待,进程能分担 Python 计算
CPython 的 GIL 让同一进程中某一时刻通常只有一个线程执行 Python 字节码,因此多个线程并不能让纯 Python 的 CPU 密集循环线性加速。但线程在等待 socket、文件或数据库时可以让其他线程运行;很多数值库、图像库和推理运行时也会在 native 代码中释放 GIL。GIL 不是“线程无用”,而是提醒你先测量工作负载。
进程拥有独立解释器和地址空间,可以真正并行 CPU 密集的 Python 计算;代价是参数和返回值通常要 pickle,模型对象可能被复制或重新加载,Windows 还要注意 spawn 与 if __name__ == "__main__"。如果输入很大,进程间复制本身就可能抵消收益。对 GPU 推理,要遵守运行时规定,不能随意为每条样本创建进程。
from concurrent.futures import ProcessPoolExecutor, ThreadPoolExecutor
def read_url(url: str) -> bytes:
return requests.get(url, timeout=5).content
def cpu_feature(payload: bytes) -> list[float]:
return expensive_python_transform(payload)
with ThreadPoolExecutor(max_workers=8) as threads:
payloads = list(threads.map(read_url, urls))
with ProcessPoolExecutor(max_workers=2) as processes:
features = list(processes.map(cpu_feature, payloads))
线程池把阻塞式 HTTP 包起来,进程池把可序列化的 CPU 工作移出当前解释器;它们都必须有 worker 上限、超时和关闭策略。若同步库能改成 async 客户端,asyncio 往往更节省线程;若无法改库,有限线程池是清晰的适配边界。
3. 背压:生产速度必须服从消费资源
没有背压的代码会把所有输入先转换成任务列表。几百万条事件会同时占用内存;模型服务的连接池、速率配额和 GPU 显存也会先于 CPU 成为瓶颈。有限 asyncio.Queue(maxsize=n) 让生产者在队列满时等待,semaphore 让消费者限制同时访问上游的数量。队列大小、worker 数和超时应是配置,并进入运行日志。
import asyncio
async def worker(queue: asyncio.Queue, results: list[dict]):
while True:
item = await queue.get()
if item is None:
queue.task_done()
return
try:
result = await fetch_json(item)
results.append({"url": item, "ok": True, "value": result})
except Exception as exc:
results.append({"url": item, "ok": False, "error": type(exc).__name__})
finally:
queue.task_done()
async def run_limited(urls: list[str], worker_count: int = 4):
queue = asyncio.Queue(maxsize=worker_count * 2)
results: list[dict] = []
workers = [asyncio.create_task(worker(queue, results))
for _ in range(worker_count)]
for url in urls:
await queue.put(url)
await queue.join()
for _ in workers:
await queue.put(None)
await asyncio.gather(*workers)
return results
这里的 task_done() 必须和每次 get() 配对,否则 queue.join() 会永久等待;哨兵 None 让 worker 有明确的结束信号。真实系统还要在取消时清空或关闭队列,并保证未完成的 HTTP 连接、文件和锁在 finally 中释放。
4. 共享状态、取消与结果顺序
事件循环中的协程虽然通常在 await 点交错运行,但共享 list、dict 的业务不变量仍需要设计;线程更需要 Lock、Queue 或只读数据。不要依赖 GIL 保护“先检查再修改”这样的复合操作。取消不是错误日志噪音,而是客户端断开、部署和超时的正常控制流;捕获后应清理并重新抛出,不能把取消转换成成功结果。
结果顺序也要明确。gather 按输入顺序返回,即使后面的请求先完成;asyncio.as_completed 按完成顺序给出结果,适合尽快发送已完成的 embedding,但必须保留输入 ID。线程池的 map 同样按输入顺序返回,若需要失败不影响其他任务,就用 future 逐项读取并记录上下文。
运行验证:观察并发度、完成顺序和输出
用 20 个 fake URL 运行 run_limited,让每个任务等待 10 毫秒并维护一个在途计数器;当 worker_count=4 时,预期最大在途数不超过 4,最终输出仍包含 20 条并按 URL 关联。再让第 7 个 URL 抛出 TimeoutError,验证其他结果仍能被收集,失败项包含 URL 和异常类型。
运行验证不能只看总耗时。还要输出 max_in_flight、成功数、失败数和取消数;如果把 worker 从 4 调成 40 后吞吐不升反降、p95 延迟和 429 增长,说明上游资源已经饱和。对进程池,再验证函数可 pickle、输入复制成本和 Windows 主入口保护。
常见错误、排错与调试路径
看到所有异步请求一起变慢,先搜索 time.sleep、requests.get 和长 CPU 循环是否出现在事件循环里;它们会阻塞所有协程。看到进程池启动失败,检查函数是否定义在模块顶层、参数是否可序列化以及是否缺少主入口保护。看到任务挂起,检查每个 queue.get 是否有对应 task_done,是否有 worker 永久等待而没有哨兵。
看到内存持续增长,检查是否一次性创建了全部 coroutine、结果列表是否无限增长、生产者是否快过消费者。看到重复或丢失结果,给每个任务绑定输入 ID 和 task name,再记录创建、完成、失败、取消四个事件。调试并发最有效的证据是时间线和资源计数,而不是再加一行无上下文的 print。
练习:限制并发的远程读取
任务是对一组 URL 并发读取,但最多同时运行 4 个任务;单个 URL 失败要能定位,所有任务结束后返回成功 payload 和失败 URL。请分别说明 asyncio 版本和同步线程版本的超时、取消和资源释放策略,并设计一个测试证明 max_in_flight <= 4。
线程、进程与 asyncio 练习
使用 asyncio.Semaphore(4) 包住 fetch_json(url),每个请求设置 timeout;成功结果与异常按 URL 关联,不能让一个失败丢失其他任务的诊断信息。
给我一点提示
为每个 URL 创建带上下文的协程;需要保留全部结果时使用 gather(..., return_exceptions=True) 后逐项处理。
查看参考答案
sem = asyncio.Semaphore(4)
async def limited_fetch(url):
async with sem:
try:
value = await asyncio.wait_for(fetch_json(url), timeout=5)
return {"url": url, "ok": True, "value": value}
except asyncio.CancelledError:
raise
except Exception as exc:
return {"url": url, "ok": False, "error": type(exc).__name__}
results = await asyncio.gather(*(limited_fetch(url) for url in urls)) 本节结论
用 20 个 URL 的 fake client 统计同时在途数量,确认最大值不超过 4;再注入一个超时,确认其他 URL 仍有结果;最后取消外层任务,确认没有把取消记录成普通失败,也没有遗留未等待的 task。
与 AI 数据管线、模型和服务连接
批量读取标注、并发调用 embedding、并行生成特征和在线推理都是并发场景,但资源边界不同:网络 I/O 受连接和 QPS 限制,CPU 特征工程受核数限制,模型推理受 GPU 显存和 batch size 限制,数据库受连接数和锁等待限制。把 worker 数、队列长度、超时和重试预算记录到每次 job,才能解释吞吐变化。
一个稳健的 AI 服务调用链会把任务 ID 与输入样本 ID 贯穿 coroutine、线程 future 或进程结果;失败项进入可重放队列,取消项不伪装为模型拒答,背压指标进入监控。并发的目标是可预测地利用资源,而不是让所有请求尽早出发。
小结
asyncio 适合大量等待,线程适合阻塞库和 I/O,进程适合可分割且可序列化的 CPU 工作;GIL 解释了直觉边界,但测量和资源上限才决定选择。任务必须有生命周期,队列必须有背压,取消必须被清理和传播,输出必须保留输入关联。这样并发才能真正服务于可验证的 AI 数据和推理系统。
延伸阅读
先完成本节练习,再用这些资料查阅完整 API 和真实项目组织方式。
阶段共 8 节课,按顺序完成更容易建立完整的迁移模型。