Python 并发与异步
Python 并发不是“开很多线程”这么简单。商业项目里,真正要解决的是:接口为什么慢、批处理为什么卡住、模型 API 为什么把服务拖死、CPU 为什么打满但吞吐上不去、async 接口为什么全站变慢、线程池和连接池为什么一起耗尽。
先记住一句话:
Python 并发选型先看瓶颈:等 IO 用线程池或 asyncio,算 CPU 用多进程或底层释放 GIL 的库;生产上必须加超时、限流、背压、异常处理、幂等和观测。
学习目标
学完本页要能做到:
- 区分并发、并行、异步。
- 解释 CPython GIL 限制了什么、没有限制什么。
- 会在线程池、进程池、asyncio 之间做选型。
- 知道事件循环、协程、Task、Future、
await的工作关系。 - 能写带超时、并发限制、异常收集和失败补偿的并发代码。
- 能排查事件循环阻塞、线程池耗尽、下游被打爆、内存上涨和 CPU 高。
并发、并行、异步不是一回事
| 概念 | 重点 | 例子 |
|---|---|---|
| 并发 | 一段时间内处理多个任务 | 一个服务同时处理多个请求 |
| 并行 | 同一时刻多个任务真的同时运行 | 多核 CPU 同时跑多个进程 |
| 异步 | 等待结果时不阻塞当前执行流 | 等 HTTP 返回时先处理别的请求 |
flowchart TD
A["多个任务"] --> B{"主要时间花在哪里"}
B -- "等待网络、数据库、文件" --> C["IO 密集"]
B -- "计算、压缩、解析大数据" --> D["CPU 密集"]
C --> E["线程池或 asyncio"]
D --> F["多进程、C扩展、分布式任务"]如果任务大部分时间在等网络,提升方式是别让 CPU 干等;如果任务大部分时间在算,提升方式是利用多核或更快的底层实现。
三种并发方式怎么选
| 方式 | 适合场景 | 优点 | 风险 |
|---|---|---|---|
threading / 线程池 | 同步 HTTP、同步数据库、文件 IO | 改造成本低,适合阻塞库 | 线程数过多、无超时会耗尽线程 |
multiprocessing / 进程池 | CPU 密集计算、图片处理、压缩、解析 | 绕开单进程 GIL,可用多核 | 启动重、数据序列化成本高 |
asyncio | 大量异步网络 IO、WebSocket、流式接口 | 单线程高并发,资源占用低 | 链路里混阻塞调用会卡住事件循环 |
选择流程:
flowchart TD
A["准备做并发"] --> B{"任务是 CPU 密集吗"}
B -- "是" --> C["优先多进程或底层库"]
B -- "否" --> D{"使用的库支持 async 吗"}
D -- "是" --> E["asyncio"]
D -- "否" --> F["线程池"]
E --> G["设置并发上限、超时、取消、异常处理"]
F --> G
C --> H["控制进程数、分块、序列化成本"]GIL到底限制什么
GIL 是 CPython 的全局解释器锁。它让同一时刻通常只有一个线程执行 Python 字节码。
flowchart TD
A["线程 A"] --> D["竞争 GIL"]
B["线程 B"] --> D
C["线程 C"] --> D
D --> E["同一时刻一个线程执行 Python 字节码"]为什么有 GIL:
- CPython 使用引用计数管理对象生命周期。
- 许多内部对象操作需要保护。
- 一个全局锁让解释器实现更简单,也让单线程性能更稳定。
GIL 的真实影响:
| 任务 | 多线程是否有效 | 原因 |
|---|---|---|
| 网络请求 | 有效 | 等网络时线程不一直执行 Python 字节码 |
| 数据库查询 | 有效 | 等数据库返回期间可切换到其他线程 |
| 文件 IO | 有效 | 等磁盘时可切换 |
| 纯 Python 循环计算 | 通常无效 | 多线程争 GIL,不能真正并行执行字节码 |
| NumPy/Pandas 部分计算 | 可能有效 | 底层 C 扩展可能释放 GIL |
所以面试不能说“Python 多线程没用”。准确说法是:CPython 多线程不适合提升纯 Python CPU 密集计算,但对 IO 密集任务仍然很常用。
线程池全过程
线程池的价值是复用线程并限制并发,而不是每个任务都 new Thread。
flowchart TD
A["提交任务"] --> B["进入任务队列"]
B --> C{"是否有空闲工作线程"}
C -- "有" --> D["工作线程取任务"]
C -- "没有" --> E["任务等待"]
D --> F["执行阻塞 IO"]
F --> G{"是否超时或异常"}
G -- "是" --> H["记录失败"]
G -- "否" --> I["返回结果"]Demo:批量调用资产接口
from concurrent.futures import ThreadPoolExecutor, as_completed
import requests
def fetch_asset(asset_code: str) -> dict:
response = requests.get(
f"http://127.0.0.1:8000/assets/{asset_code}",
timeout=(1, 3),
)
response.raise_for_status()
return response.json()
def fetch_assets(asset_codes: list[str]) -> tuple[list[dict], list[tuple[str, str]]]:
results: list[dict] = []
errors: list[tuple[str, str]] = []
with ThreadPoolExecutor(max_workers=8, thread_name_prefix="asset-fetch") as executor:
future_map = {
executor.submit(fetch_asset, code): code
for code in asset_codes
}
for future in as_completed(future_map):
code = future_map[future]
try:
results.append(future.result())
except Exception as exc:
errors.append((code, str(exc)))
return results, errors为什么要这样写:
| 设计 | 作用 |
|---|---|
max_workers=8 | 控制并发,避免打爆下游 |
timeout=(1, 3) | 连接超时和读取超时,避免线程永久卡住 |
as_completed | 谁先完成先处理,避免慢任务挡住快任务 |
收集 errors | 批处理允许部分失败,便于补偿 |
thread_name_prefix | 线上线程栈更容易识别 |
线程池常见坑
| 坑 | 后果 | 正确做法 |
|---|---|---|
| 没有超时 | 线程被慢接口永久占住 | HTTP/DB/模型调用都设置超时 |
| 线程数太大 | 下游、连接数、文件句柄被打爆 | 结合下游容量设置并发 |
| 任务里继续提交子任务并同步等待 | 线程饥饿死锁 | 拆线程池或改异步编排 |
| 只看成功结果 | 部分失败丢失 | 记录失败明细和补偿任务 |
| 无限重试 | 放大故障 | 指数退避、最大次数、幂等 |
多进程全过程
多进程适合 CPU 密集任务。每个进程有独立解释器和自己的 GIL,可以真正使用多核。
flowchart TD
A["主进程拆分任务"] --> B["序列化参数"]
B --> C["发送给子进程"]
C --> D["子进程独立计算"]
D --> E["序列化结果返回"]
E --> F["主进程合并结果"]Demo:批量计算文件 SHA256
from concurrent.futures import ProcessPoolExecutor
from pathlib import Path
import hashlib
def file_sha256(path: str) -> tuple[str, str]:
digest = hashlib.sha256()
with Path(path).open("rb") as file:
for chunk in iter(lambda: file.read(1024 * 1024), b""):
digest.update(chunk)
return path, digest.hexdigest()
def hash_files(paths: list[str]) -> dict[str, str]:
with ProcessPoolExecutor(max_workers=4) as executor:
return dict(executor.map(file_sha256, paths))
if __name__ == "__main__":
print(hash_files(["a.zip", "b.zip"]))Windows 下必须有 if __name__ == "__main__",否则子进程导入主模块时可能反复创建新进程。
多进程不适合小任务过度拆分。因为进程创建、参数序列化、结果回传都有成本。商业批处理要先分块,例如每个进程处理一批文件或一段数据,而不是一行数据一个进程。
asyncio核心模型
asyncio 的核心是事件循环。协程遇到 await 后主动挂起,把控制权还给事件循环,事件循环继续调度其他可运行任务。
flowchart TD
A["事件循环"] --> B["运行协程 A"]
B --> C{"遇到 await IO"}
C -- "是" --> D["挂起 A,登记等待"]
D --> E["运行协程 B"]
E --> F{"A 的 IO 完成"}
F -- "是" --> G["恢复 A"]关键概念:
| 概念 | 解释 |
|---|---|
| coroutine | 调用 async 函数得到的协程对象,本身不会自动执行 |
| event loop | 调度协程、监听 IO、恢复任务的循环 |
| Task | 把协程包装成可调度任务 |
| Future | 表示未来会完成的结果 |
| await | 等待一个异步结果,同时让出事件循环 |
async函数为什么没执行
async def hello() -> str:
return "hello"
coro = hello()
print(coro)这只创建了协程对象,不会真正运行。要执行:
import asyncio
async def hello() -> str:
return "hello"
async def main() -> None:
result = await hello()
print(result)
asyncio.run(main())asyncio生产Demo:批量生成Embedding
import asyncio
import httpx
async def embed_text(client: httpx.AsyncClient, text: str) -> list[float]:
response = await client.post(
"http://127.0.0.1:8000/embed",
json={"text": text},
)
response.raise_for_status()
return response.json()["vector"]
async def embed_all(texts: list[str]) -> tuple[list[list[float]], list[tuple[str, str]]]:
limits = httpx.Limits(max_connections=20)
timeout = httpx.Timeout(connect=1.0, read=5.0, write=2.0, pool=1.0)
semaphore = asyncio.Semaphore(10)
async with httpx.AsyncClient(limits=limits, timeout=timeout) as client:
async def run_one(text: str) -> tuple[str, list[float] | Exception]:
async with semaphore:
try:
return text, await embed_text(client, text)
except Exception as exc:
return text, exc
raw_results = await asyncio.gather(*(run_one(text) for text in texts))
vectors: list[list[float]] = []
errors: list[tuple[str, str]] = []
for text, result in raw_results:
if isinstance(result, Exception):
errors.append((text, str(result)))
else:
vectors.append(result)
return vectors, errors这个 Demo 适合 RAG 离线入库。重点不是调用 Embedding,而是:
- 连接池限制
max_connections。 - 业务并发限制
Semaphore。 - 超时分成连接、读取、写入、连接池等待。
- 单条失败不会让整批丢失。
- 失败记录可进入补偿任务。
超时、取消和背压
并发系统如果没有超时,就会把等待变成资源泄漏。
单任务超时
import asyncio
async def call_slow_service() -> str:
await asyncio.sleep(10)
return "ok"
async def main() -> None:
try:
result = await asyncio.wait_for(call_slow_service(), timeout=2)
print(result)
except asyncio.TimeoutError:
print("service timeout")批处理限流
import asyncio
async def handle_item(item: int) -> int:
await asyncio.sleep(0.1)
return item * 2
async def handle_batch(items: list[int]) -> list[int]:
semaphore = asyncio.Semaphore(20)
async def guarded(item: int) -> int:
async with semaphore:
return await handle_item(item)
return await asyncio.gather(*(guarded(item) for item in items))背压的意思是:当下游处理不过来时,上游不能无限制继续提交任务。否则内存队列会越来越大,最终 OOM 或雪崩。
flowchart TD
A["上游快速提交任务"] --> B{"并发是否超过上限"}
B -- "否" --> C["执行任务"]
B -- "是" --> D["等待、拒绝或降级"]
C --> E{"下游是否慢"}
E -- "是" --> F["减少并发、延迟重试"]
E -- "否" --> G["正常完成"]async中最危险的坑:阻塞事件循环
错误示例:
import time
async def bad_api() -> dict:
time.sleep(5)
return {"ok": True}time.sleep 会卡住整个事件循环,这 5 秒里其他协程也无法继续运行。
正确写法:
import asyncio
async def good_api() -> dict:
await asyncio.sleep(5)
return {"ok": True}如果必须调用阻塞库:
import asyncio
import requests
def call_sync_api() -> int:
response = requests.get("http://127.0.0.1:8000/health", timeout=3)
return response.status_code
async def proxy() -> int:
return await asyncio.to_thread(call_sync_api)FastAPI 中尤其常见:async def 接口里直接用同步数据库驱动、requests、大文件读取、CPU 计算,都会导致并发请求一起变慢。
FastAPI中def和async def怎么选
| 写法 | 适合 | 原理 |
|---|---|---|
def | 同步数据库驱动、同步 SDK、普通阻塞代码 | FastAPI 通常放到线程池执行,避免阻塞事件循环 |
async def | 异步 HTTP、异步数据库、WebSocket、流式输出 | 在事件循环中 await IO,提高并发 |
不要为了“看起来高级”把所有接口都写成 async def。如果链路里的库都是同步阻塞的,写成 def 反而更稳。
商业场景一:采集任务并发调用医院接口
场景:医疗数据采集平台要从多个医院前置机拉取数据。每个医院接口响应时间不同,有的会超时。
设计要点:
| 问题 | 设计 |
|---|---|
| 下游能力有限 | 每个医院单独限流 |
| 接口可能超时 | 设置连接和读取超时 |
| 部分失败 | 记录失败明细,后续补偿 |
| 重试风险 | 查询类可重试,写入类必须幂等 |
| 结果过大 | 分页拉取,不一次放入内存 |
流程:
flowchart TD
A["读取待采集医院列表"] --> B["按医院分组限流"]
B --> C["并发调用接口"]
C --> D{"成功?"}
D -- "是" --> E["清洗并分批入库"]
D -- "否" --> F["记录失败和错误码"]
F --> G["补偿任务重试"]
E --> H["更新采集批次状态"]商业场景二:AI知识库批量Embedding
RAG 入库时需要把文档 chunk 批量调用 Embedding 接口。这里不是并发越高越好。
必须控制:
- 模型服务 QPS 限制。
- 单次请求 token 限制。
- 失败重试和幂等。
- 向量写入数据库的批大小。
- 任务进度和断点续跑。
错误做法:
# 一次创建几万个任务,没有限流,内存和下游都会炸
await asyncio.gather(*(embed_text(client, text) for text in all_texts))正确思路:分批 + 限流 + 失败记录 + 可重跑。
生产排查流程
flowchart TD
A["Python 并发问题"] --> B{"表现"}
B -- "CPU 高但吞吐低" --> C["是否 CPU 密集却用了线程"]
B -- "程序卡住" --> D["线程池任务是否无超时"]
B -- "async 接口全变慢" --> E["事件循环是否被阻塞"]
B -- "内存上涨" --> F["任务数量、结果列表、队列是否无限增长"]
B -- "下游报错变多" --> G["并发、重试、连接池是否过大"]
C --> H["用 cProfile/py-spy 找热点"]
D --> I["看线程栈、日志耗时、超时配置"]
E --> J["查同步调用、CPU计算、大文件读取"]
F --> K["加批处理、背压、流式处理"]
G --> L["限流、退避、熔断、降级"]排查工具:
| 工具 | 用途 |
|---|---|
| 日志耗时 | 拆请求、数据库、HTTP、AI、文件耗时 |
threading.enumerate() | 看线程数量 |
asyncio.all_tasks() | 看未完成协程 |
py-spy | 不改代码采样查看 Python 热点 |
cProfile | 分析 CPU 函数耗时 |
tracemalloc | 查内存分配来源 |
faulthandler | 卡死时导出线程栈 |
导出线程栈示例:
import faulthandler
import signal
faulthandler.register(signal.SIGUSR1)Linux 下可以给进程发信号:
kill -USR1 <pid>这样能在程序卡住时看到线程停在哪里。
面试标准回答
Python并发怎么选
先判断任务类型。IO 密集任务,例如网络请求、数据库访问、文件 IO,可以用线程池或 asyncio;CPU 密集任务在 CPython 中受 GIL 影响,多线程通常不能提升纯 Python 字节码的并行计算能力,更适合多进程、C 扩展或分布式任务。asyncio 适合大量异步网络 IO,但链路中不能混入阻塞调用。
GIL是什么
GIL 是 CPython 的全局解释器锁,它让同一时刻通常只有一个线程执行 Python 字节码。它简化了解释器对象内存管理,但限制了 CPU 密集 Python 多线程并行执行。不过 IO 密集任务等待网络或磁盘时会让出执行机会,所以多线程对 IO 密集仍然有效。
asyncio为什么能高并发
asyncio 基于事件循环和协程。协程遇到 await 等待 IO 时会挂起,把控制权交还给事件循环,事件循环继续调度其他协程。它不是让单个请求更快,而是让大量 IO 等待时间被复用,从而提高整体吞吐。
async接口全站变慢怎么排查
先怀疑事件循环被阻塞。检查 async def 中是否调用了 time.sleep、同步 HTTP、同步数据库驱动、大文件读取或 CPU 密集计算;再看外部接口是否无超时、连接池是否耗尽、并发是否过大。修复方式是改用异步库、用 asyncio.to_thread 包装阻塞调用、限制并发、设置超时,并把 CPU 密集任务移到进程池或任务系统。
并发任务为什么要做幂等
并发调用中常见超时、重试、取消和部分失败。如果写操作没有幂等键,重试可能重复创建订单、重复发通知、重复写入采集记录。商业系统要给写操作设计业务唯一键或 idempotency key,让重复请求返回同一个结果或被安全忽略。
关联知识点
| 知识点 | 说明 |
|---|---|
| Python Web API 开发 | FastAPI 中同步和异步接口如何选择 |
| Python 数据库访问 | 连接池、事务、慢 SQL 与并发关系 |
| 日志与调试 | 如何通过日志和 request_id 排查并发问题 |
| 测试 | 如何测试超时、异常、Mock 外部依赖 |
| Python与AI开发 | RAG、Embedding、模型接口并发调用 |
| Python面试题 | 标准回答与追问 |
本章小结
Python 并发的关键不是“开更多任务”,而是判断瓶颈、选择模型、控制资源和能排查问题。IO 密集用线程池或 asyncio,CPU 密集用多进程或底层释放 GIL 的库。生产代码必须有超时、限流、背压、异常收集、幂等和观测,否则并发越高,故障越快。
