发布于5小时前5小时 在 Python Web 服务中,asyncio 接口偶发性变慢,常见原因不是协程数量不足,而是事件循环线程被同步代码阻塞。典型场景包括在 async def 中直接调用 requests、time.sleep、同步数据库驱动或大文件读写。少量请求时问题不明显,并发上升后会表现为所有请求一起变慢、超时,甚至健康检查也无法及时返回。 一、问题定位:先确认是不是事件循环被阻塞 假设服务使用 FastAPI、Starlette 或其他 asyncio 框架,先观察三个指标:事件循环延迟、请求耗时分布和进程线程数。不要只看平均响应时间,阻塞通常会造成 P95、P99 明显升高。 事件循环延迟持续升高:说明 loop 线程没有获得运行机会。 多个无关接口同时变慢:通常比单个接口变慢更能说明是共享事件循环被阻塞。 CPU 使用率不高但请求堆积:常见于同步 I/O 或 sleep 阻塞,而不是纯计算。 开发环境可以开启 asyncio 的慢回调检测。下面的命令会让超过 100 毫秒的事件循环回调被记录: PYTHONASYNCIODEBUG=1 python -X dev app.py 如果使用 uvicorn,可以在启动参数中增加日志级别,并在应用入口配置慢回调阈值: import asyncio loop = asyncio.get_running_loop() loop.set_debug(True) loop.slow_callback_duration = 0.1 线上排查时可以使用 py-spy 获取不暂停进程的调用栈: py-spy top --pid 12345 py-spy dump --pid 12345 --threads 重点检查事件循环线程是否停留在 requests、time.sleep、同步数据库驱动或大块 CPU 计算中。如果堆栈长期停在某个同步函数,基本可以确认阻塞来源。 二、一个容易被忽略的错误示例 下面的代码看起来是异步接口,但 requests.get 和 time.sleep 都是同步调用。async def 只改变了函数的调用方式,并不会自动把函数体中的同步操作变成异步操作。 import time import requests from fastapi import FastAPI app = FastAPI() @app.get('/bad') async def bad(): response = requests.get(' timeout=3) time.sleep(0.5) return {'status': response.status_code} 当一个请求执行 requests.get 时,当前事件循环线程会同步等待网络返回。在等待期间,同一个 worker 上其他协程也无法运行。即使远端接口只有 500 毫秒延迟,100 个并发请求也可能形成明显排队。 三、优先方案:使用真正的异步客户端 如果外部 HTTP 请求是主要操作,应优先使用 httpx.AsyncClient 或 aiohttp,而不是把 requests 简单包进协程。客户端应在应用生命周期内复用,避免每个请求都新建连接池。 from contextlib import asynccontextmanager import httpx from fastapi import FastAPI, HTTPException @asynccontextmanager async def lifespan(app: FastAPI): app.state.http = httpx.AsyncClient( timeout=httpx.Timeout(3.0, connect=1.0), limits=httpx.Limits( max_connections=100, max_keepalive_connections=20 ) ) yield await app.state.http.aclose() app = FastAPI(lifespan=lifespan) @app.get('/good') async def good(): try: response = await app.state.http.get(' response.raise_for_status() except httpx.TimeoutException: raise HTTPException(status_code=504, detail='upstream timeout') except httpx.HTTPError: raise HTTPException(status_code=502, detail='upstream error') return {'status': response.status_code} 这里有几个边界需要明确: timeout 必须配置。没有超时的异步请求可能长期占用连接和任务。 连接池上限需要结合上游限流配置,max_connections 过大不一定更快,可能把本地并发压力放大到上游。 HTTP 连接复用要求客户端在应用关闭时显式 close,否则测试和热重载环境中容易出现未关闭连接警告。 HTTP 502 表示上游返回异常或不可用,HTTP 504 更适合表达超时,便于监控区分故障类型。 四、无法替换同步库时:使用 to_thread,但要限制并发 有些遗留 SDK 只有同步接口,短期内无法更换。这时可以用 asyncio.to_thread 将同步调用放入线程池,避免阻塞事件循环。但线程池不是无限资源,不能把所有同步任务无上限提交进去。 import asyncio from functools import partial import requests from fastapi import FastAPI, HTTPException app = FastAPI() request_slots = asyncio.Semaphore(20) def fetch_sync(url, timeout): return requests.get(url, timeout=timeout) @app.get('/bridge') async def bridge(): async with request_slots: try: call = partial(fetch_sync, ' 3) response = await asyncio.wait_for( asyncio.to_thread(call), timeout=4 ) except asyncio.TimeoutError: raise HTTPException(status_code=504, detail='upstream timeout') except requests.RequestException: raise HTTPException(status_code=502, detail='upstream error') return {'status': response.status_code} Semaphore 控制的是业务并发,不等同于线程池大小。Python 默认线程池大小由运行时决定,实际部署时应根据 CPU、上游容量和请求耗时验证。更严格的做法是为这类同步任务创建专用 ThreadPoolExecutor,避免占满默认线程池,影响其他 to_thread 任务。 from concurrent.futures import ThreadPoolExecutor import asyncio legacy_pool = ThreadPoolExecutor( max_workers=16, thread_name_prefix='legacy-io' ) async def call_legacy(client, payload): loop = asyncio.get_running_loop() return await loop.run_in_executor( legacy_pool, client.send, payload ) 需要注意,取消 await 不一定能终止已经开始执行的同步函数。asyncio 任务被取消后,线程中的 requests 或 SDK 调用仍可能继续运行。因此同步库自身也必须配置连接超时和读取超时,不能只依赖外围的 wait_for。 五、同步数据库驱动和 CPU 计算的处理边界 同步数据库驱动不应直接放进 async def。可选方案是换用对应的异步驱动,例如 asyncpg、aiomysql;如果暂时不能替换,则按照同步 I/O 的方式放入专用线程池,并限制查询并发。 CPU 密集型任务则不适合用 to_thread 作为长期方案。CPython 的 GIL 会限制纯 Python 计算的并行效果,事件循环虽然不会直接执行这段代码,但线程仍可能争用 CPU。对于图像处理、压缩、复杂规则计算等任务,应考虑 ProcessPoolExecutor、独立任务队列或拆分为独立服务。 from concurrent.futures import ProcessPoolExecutor import asyncio process_pool = ProcessPoolExecutor(max_workers=4) def calculate(data): return sum(x * x for x in data) async def calculate_async(data): loop = asyncio.get_running_loop() return await loop.run_in_executor( process_pool, calculate, data ) ProcessPoolExecutor 要求传入函数和参数可以被 pickle,不能直接传递打开的文件、网络连接或包含不可序列化对象的实例。进程池还会产生序列化和进程间通信成本,小任务未必比当前线程直接执行更快。 六、排查时容易误判的几个情况 只把函数声明改成 async def:这不会改变 requests、数据库驱动或文件操作的同步性质。 只增加 worker 数量:可以暂时降低单个进程的影响,但会增加连接数和内存,无法修复阻塞代码本身。 只增加客户端超时时间:这会延长资源占用,不能解决事件循环阻塞,甚至可能扩大请求堆积。 把所有操作都放进线程池:线程池同样有上限,未限制的提交会导致任务排队、上下文切换增加和内存上涨。 忽略响应体读取:HTTP 客户端的连接通常要在响应体消费完或响应关闭后才能回收到连接池。 七、上线前验证方法 改造后至少进行三组验证:无上游延迟、固定上游延迟、上游超时。用一个独立的健康接口观察事件循环是否还能及时响应。例如在测试环境并发请求业务接口的同时,每 100 毫秒访问一次健康接口,记录最大延迟。 curl -s -o /dev/null -w '%{http_code} %{time_total}\n' \ hey -n 500 -c 50 重点比较改造前后的 P95、P99、活动连接数、线程数和上游错误率。若业务接口变快但线程数持续增长,通常说明同步桥接没有设置合适的并发上限或超时。若健康接口仍然出现长尾,则应继续检查中间件、日志处理器和同步缓存客户端等非业务代码。 八、修复建议总结 异步接口中的网络访问优先使用原生异步客户端,并复用连接池。 无法替换的同步 I/O 使用 to_thread 或专用线程池,同时设置并发闸门。 CPU 密集型任务放入进程池或独立任务系统,不要占用事件循环线程。 每一层都设置合理超时,并区分连接超时、读取超时和业务超时。 通过 asyncio 慢回调日志、py-spy 栈信息和 P95/P99 指标确认修复效果,而不是只观察平均耗时。 判断标准很简单:事件循环线程只负责快速调度,阻塞 I/O 和重 CPU 计算必须明确隔离;隔离后还要限制并发、设置超时,并验证取消行为是否真的能释放资源。这样才能避免把一个表面上的异步接口,变成实际上的同步排队点。 网络请求、日志与边界流量分析示意
5小时前5小时 可以先把“疑似阻塞点”从感觉变成可观测数据,再决定哪些接口改线程池。几个实用做法: 1)开启 asyncio 慢回调日志,先抓 event loop 被卡住的位置: import asyncio import logging logging.basicConfig(level=logging.DEBUG) async def main(): loop = asyncio.get_running_loop() loop.set_debug(True) loop.slow_callback_duration = 0.05 # 50ms,按业务调大/调小 # your app start here asyncio.run(main()) 启动时也可以加: PYTHONASYNCIODEBUG=1 python app.py 如果日志里出现某个 callback/task 执行时间异常,就优先看里面是否有 `requests`、同步 ORM、文件 IO、`time.sleep()`、大量 CPU 计算、同步 DNS 等。 2)在线程池改造前,可以加一个轻量 watchdog,定期检测 loop 延迟,方便在压测时定位时间点: import asyncio import time import logging log = logging.getLogger(__name__) async def loop_lag_monitor(interval=0.1, threshold=0.2): loop = asyncio.get_running_loop() expected = loop.time() + interval while True: await asyncio.sleep(interval) now = loop.time() lag = now - expected if lag > threshold: log.warning("event loop lag %.3fs", lag) expected = now + interval 在应用启动时创建: asyncio.create_task(loop_lag_monitor()) 3)确认是同步调用后,再用线程池包起来,不建议直接在协程里调用同步函数: Python 3.9+ 可以用: import asyncio def blocking_call(arg): # requests / 同步 SDK / 同步文件处理等 return do_something(arg) async def handler(arg): result = await asyncio.to_thread(blocking_call, arg) return result 如果需要控制线程池大小: import asyncio from concurrent.futures import ThreadPoolExecutor executor = ThreadPoolExecutor(max_workers=20) def blocking_call(arg): return do_something(arg) async def handler(arg): loop = asyncio.get_running_loop() return await loop.run_in_executor(executor, blocking_call, arg) 4)线程池大小不要盲目调很大。IO 阻塞型可以适当大一点,但同步数据库连接、第三方 API、文件句柄都有上限。建议同时限制并发: import asyncio sem = asyncio.Semaphore(50) async def handler(arg): async with sem: return await asyncio.to_thread(blocking_call, arg) 否则表面上 event loop 不堵了,实际可能把数据库连接池或下游服务打满。 5)常见容易漏掉的阻塞点: `time.sleep()`:改成 `await asyncio.sleep()` `requests`:换 `httpx.AsyncClient` / `aiohttp`,或临时放线程池 同步数据库驱动:换异步驱动,或线程池隔离 CPU 密集任务:线程池不一定有效,考虑 `ProcessPoolExecutor` 或任务队列 日志写入慢:尤其是网络日志、同步 handler,可以改队列日志 同步文件扫描/压缩/图片处理:不要直接放在 async handler 里跑 6)如果是线上问题,推荐在压测环境用 `py-spy` 看卡在哪里,不需要改代码: py-spy top --pid <PID> py-spy dump --pid <PID> py-spy record -o profile.svg --pid <PID> --duration 60 如果看到主线程长期停在某个同步库调用里,基本就能确定是 event loop 被阻塞。 整体思路是:先用 debug/lag/采样确认具体阻塞点,再按类型处理。能换异步库优先换异步库;短期无法替换的同步接口,再用 `to_thread` 或 `run_in_executor` 隔离,并配合并发限制,避免把阻塞问题转移到线程池和下游资源上。
创建帐户或登录后发表意见