一、问题背景:一个聚合API拖垮了整个服务
上个月接手一个老服务,核心接口是/api/v1/aggregate,它需要同时调用三个下游服务:用户信息(耗时400ms)、订单列表(耗时600ms)、推荐内容(耗时800ms)。最初的实现是串行调用,用Requests同步阻塞。压测结果惨不忍睹:并发50时,TP99直接破2秒,CPU利用率不到30%,大部分线程都在waiting状态。
# 改造前的核心逻辑(伪代码)
def aggregate(user_id):
user = requests.get(f"http://user-svc/{user_id}").json()
orders = requests.get(f"http://order-svc/{user_id}").json()
reco = requests.get(f"http://reco-svc/{user_id}").json()
return {"user": user, "orders": orders, "reco": reco}
三个下游接口互不依赖,串行浪费了2/3的时间。第一反应是用ThreadPoolExecutor,但线程切换开销大,且服务本身是Gunicorn多进程部署,每个worker再开线程池,线程数会爆炸。于是决定用asyncio重写IO密集部分。
二、环境与版本:Python 3.10 + Flask 2.x 的“兼容性陷阱”
先明确环境:
Python: 3.10.12 (CPython)
Flask: 2.2.5
aiohttp: 3.8.4
gunicorn: 20.1.0 (worker_class=sync, workers=4)
第一个坑:Flask是WSGI框架,本质是同步阻塞模型。直接在视图函数里asyncio.run()会创建新事件循环,每次请求都重新创建,开销极大。正确的做法是在进程启动时创建一个全局事件循环,通过asyncio.run_coroutine_threadsafe()把协程提交到该循环中运行。
我的方案是写一个AsyncLoop管理类:
# async_loop.py
import asyncio
import threading
class AsyncLoopManager:
_loop = None
_thread = None
@classmethod
def start(cls):
"""在后台线程中启动事件循环"""
if cls._loop is not None:
return
cls._loop = asyncio.new_event_loop()
cls._thread = threading.Thread(target=cls._run_loop, daemon=True)
cls._thread.start()
@classmethod
def _run_loop(cls):
asyncio.set_event_loop(cls._loop)
cls._loop.run_forever()
@classmethod
def submit(cls, coro):
"""将协程提交到事件循环,返回concurrent.futures.Future"""
if cls._loop is None:
raise RuntimeError("AsyncLoop not started")
return asyncio.run_coroutine_threadsafe(coro, cls._loop)
@classmethod
def stop(cls):
if cls._loop:
cls._loop.call_soon_threadsafe(cls._loop.stop)
在应用启动时调用AsyncLoopManager.start()(比如在app.py的before_first_request或模块导入时)。
三、方案设计:协程池 + 信号量 + 连接复用
改造分三步骤:
- 将三个下游调用改为
async函数,使用aiohttp.ClientSession替代Requests。注意ClientSession必须复用,不能每次请求都新建,否则TCP连接无法复用,性能反而更差。 - 引入
asyncio.Semaphore控制并发。虽然用了协程,但下游服务同样有承受上限。压测发现,下游对单个实例的并发超过200时,延迟会指数上升。所以设置SEM = asyncio.Semaphore(200)。 - 在Flask视图函数中,不要用
asyncio.run()(会阻塞worker线程),而是用AsyncLoopManager.submit()提交协程,然后等待future结果。
改造后的核心逻辑:
# views.py
import aiohttp
import asyncio
from flask import jsonify, request
from async_loop import AsyncLoopManager
# 全局session,进程内复用
_session = None
SEM = asyncio.Semaphore(200)
async def fetch_json(session, url, timeout=2.0):
async with SEM: # 信号量控制并发
try:
async with session.get(url, timeout=aiohttp.ClientTimeout(total=timeout)) as resp:
return await resp.json()
except Exception as e:
return {"error": str(e)}
async def aggregate_async(user_id):
global _session
if _session is None:
_session = aiohttp.ClientSession(
connector=aiohttp.TCPConnector(limit=300, ttl_dns_cache=300),
timeout=aiohttp.ClientTimeout(total=3.0)
)
# 三个调用并发执行
results = await asyncio.gather(
fetch_json(_session, f"http://user-svc/{user_id}"),
fetch_json(_session, f"http://order-svc/{user_id}"),
fetch_json(_session, f"http://reco-svc/{user_id}")
)
return {"user": results[0], "orders": results[1], "reco": results[2]}
@app.route("/api/v1/aggregate")
def aggregate():
user_id = request.args.get("user_id")
try:
future = AsyncLoopManager.submit(aggregate_async(user_id))
# 阻塞等待结果,但事件循环在另一线程跑,不影响其他请求
data = future.result(timeout=3.5)
return jsonify(data)
except Exception as e:
return jsonify({"error": str(e)}), 500
关键细节:
- aiohttp.ClientSession是线程安全的吗?不是。但我们的session只会在事件循环线程中被使用(所有协程都在该循环内执行),所以安全。
- TCPConnector(limit=300)限制连接池最大300个,避免对下游造成连接风暴。
- future.result(timeout=3.5)必须设置超时,否则如果下游挂掉,请求线程会一直阻塞。
四、踩坑与优化:三个让我抓狂的问题
坑1:asyncio.run()在视图函数中导致“假死”
一开始我用的是asyncio.run(aggregate_async(user_id)),压测时发现,一旦并发超过50,服务就卡死,CPU飙到100%。排查发现:每个请求都会创建新事件循环,导致大量线程切换和垃圾回收。而且asyncio.run()会阻塞当前线程,Gunicorn的sync worker模型下,每个worker同时只能处理一个请求,等于又回到了串行。
解决方案:全局事件循环 + run_coroutine_threadsafe,如上代码所示。
坑2:DNS解析阻塞事件循环
第一次压测时,发现QPS提升不明显。用py-spy dump查看,发现事件循环线程卡在getaddrinfo上。原因是aiohttp默认使用线程池做DNS解析,但线程池默认大小是4,并发一高就排队。
解决方案:显式设置connector = aiohttp.TCPConnector(use_dns_cache=True, ttl_dns_cache=300),并开启aiohttp.resolver.AsyncResolver:
from aiohttp.resolver import AsyncResolver
resolver = AsyncResolver(nameservers=["8.8.8.8", "1.1.1.1"])
connector = aiohttp.TCPConnector(resolver=resolver, use_dns_cache=True)
坑3:Gunicorn worker类型选择
最初用worker_class=sync,每个worker阻塞在future.result()上。Gunicorn的sync worker是单线程,所以理论上4个worker只能同时处理4个请求。但实际压测QPS到了850,为什么?因为请求处理时间从串行的1.8秒降到了0.8秒(并发调用),所以worker占用时间缩短,吞吐量反而上去了。但如果你想进一步压榨性能,可以把worker改成gthread或gevent,但注意它们与asyncio的兼容性。我最终保持sync,因为瓶颈在下游,而非worker线程数。
五、效果数据:从120到850 QPS
压测工具:wrk -t8 -c200 -d30s http://localhost:8080/api/v1/aggregate
| 指标 | 改造前 (Requests串行) | 改造后 (asyncio并发) | 提升 |
|---|---|---|---|
| QPS | 120 | 850 | 7.1x |
| 平均延迟 | 1.8s | 320ms | 5.6x |
| P99延迟 | 3.2s | 180ms | 17.8x |
| 错误率 | 2.1% (超时) | 0.3% | - |
| 下游并发连接数 | 50 (线程池) | 200 (Semaphore) | - |
注意P99从3.2s降到180ms,说明尾部延迟显著改善。为什么P99比平均还低? 因为平均延迟包含了偶尔的下游超时重试,而P99是正常情况下的最高延迟。
资源占用对比:
- CPU:改造前30%(线程等待),改造后45%(协程切换+IO等待)
- 内存:改造前每worker约200MB(线程栈),改造后约150MB(协程栈极小)
六、总结:什么时候该用asyncio?
这次改造让我对asyncio有了新的认知:
- 适用场景:IO密集且任务间无强依赖。如果你的API需要串行调用多个外部服务,且每个服务延迟都在100ms以上,那么asyncio的收益非常明显。
- 不适用场景:CPU密集(计算、加密)、或依赖大量C扩展库(如某些ORM)。asyncio解决不了GIL问题。
- 关键点:
- 永远复用
ClientSession和事件循环 - 用
Semaphore控制并发,保护下游 - 所有IO操作必须设置超时,否则一个慢接口会拖垮整个事件循环
- 不要用
asyncio.run()在同步框架中,使用run_coroutine_threadsafe+ 全局循环
最后说句实话:asyncio的思维模型和同步代码完全不同,调试难度翻倍。如果你的下游接口延迟都在50ms以下,或者并发量低于100,用ThreadPoolExecutor可能更简单。但如果你面对的是我这种“串行调用三个300ms+接口”的场景,asyncio是唯一能让你QPS翻五倍以上的方案。
后续优化方向:如果下游接口支持HTTP/2,可以尝试httpx的AsyncClient,支持多路复用,还能再省20%的延迟。但那是另一个故事了。