1. 问题背景:一个慢得像蜗牛的聚合API
上个月接手一个内部数据看板服务,核心接口/api/dashboard需要同时从用户服务、订单服务、库存服务和推荐服务拉取数据,然后组装返回。
原始代码是典型的串行写法:
# before.py - 串行版本
def get_dashboard(user_id: int):
user = requests.get(f"http://user-svc/users/{user_id}", timeout=2).json()
orders = requests.get(f"http://order-svc/orders?uid={user_id}", timeout=2).json()
stock = requests.get(f"http://stock-svc/stock?uid={user_id}", timeout=2).json()
recs = requests.get(f"http://rec-svc/recommend?uid={user_id}", timeout=2).json()
return {"user": user, "orders": orders, "stock": stock, "recs": recs}
每个上游服务平均响应120ms,四个串行就是480ms。高峰期并发上来,Gunicorn worker线程被requests阻塞,CPU闲着呢但请求排队。监控数据显示:P95延迟612ms,QPS只有92,经常触发网关超时告警(我们网关阈值1s)。
2. 环境与版本:Python 3.10 + FastAPI + aiohttp
先说环境,避免版本坑:
- Python 3.10.12(注意:3.10以下没有
asyncio.run()的完整语义,3.8的loop.run_until_complete会有事件循环关闭警告) - FastAPI 0.104.1(用async def定义路由,直接支持asyncio)
- aiohttp 3.9.1(推荐用ClientSession,不要自己搞连接池)
- uvloop 0.19.0(可选项,生产环境建议装,能再压5-8%吞吐)
- 部署:Gunicorn + uvicorn worker,4进程
3. 方案设计:asyncio.gather + Semaphore限流
核心思路很简单:把四个独立请求从串行改成并发。
但有几个设计决策要说明:
第一,用asyncio.gather还是asyncio.wait?
用gather,因为要拿到所有结果统一返回。wait适合只关心完成状态不需要结果的场景。
第二,必须加信号量限流。
如果没有Semaphore,当上游服务抖动时,asyncio会瞬间发出成百上千个并发连接,把上游打挂。我用asyncio.Semaphore(20)限制同一时间最多20个出站请求。
第三,aiohttp连接池复用。
不要每次请求都创建ClientSession,那会重复建立TCP连接和TLS握手。应该在应用启动时创建全局session,设置connector的limit参数。
4. 核心实现:从requests到aiohttp的重构
直接上重构后的代码:
# after.py - asyncio并发版本
import asyncio
import aiohttp
from fastapi import FastAPI
app = FastAPI()
# 全局session,应用启动时创建
session = None
@app.on_event("startup")
async def create_session():
global session
conn = aiohttp.TCPConnector(limit=50, ttl_dns_cache=300)
session = aiohttp.ClientSession(connector=conn, timeout=aiohttp.ClientTimeout(total=2.5))
@app.on_event("shutdown")
async def close_session():
await session.close()
# 信号量:限制并发出站请求数为20
sem = asyncio.Semaphore(20)
async def fetch_json(url: str):
async with sem:
async with session.get(url) as resp:
return await resp.json()
@app.get("/api/dashboard")
async def get_dashboard(user_id: int):
urls = [
f"http://user-svc/users/{user_id}",
f"http://order-svc/orders?uid={user_id}",
f"http://stock-svc/stock?uid={user_id}",
f"http://rec-svc/recommend?uid={user_id}"
]
# gather并发跑四个请求,return_exceptions=True防止单个失败拖垮全部
results = await asyncio.gather(
*[fetch_json(url) for url in urls],
return_exceptions=True
)
# 处理异常情况
user, orders, stock, recs = results
if isinstance(user, Exception):
user = {"error": "user service unavailable"}
# ... 其他类似处理
return {"user": user, "orders": orders, "stock": stock, "recs": recs}
关键点说明:
TCPConnector(limit=50):连接池最大50个连接,超出会等待释放ttl_dns_cache=300:DNS缓存5分钟,避免每次解析ClientTimeout(total=2.5):总超时2.5秒,比requests的2秒略宽松,但比网关1s要短(实际上P95只有189ms,完全够用)asyncio.Semaphore(20):20是压测得出的最佳值,低于10吞吐不足,高于30上游会偶尔超时
5. 踩坑与优化:三个真实教训
坑1:async def函数里用了阻塞调用
第一次重构时,我在fetch_json里用了json.loads(resp.text),这是CPU阻塞操作,会卡住事件循环。正确做法是直接用await resp.json(),它内部是异步解析JSON。
坑2:Gunicorn worker类型没换
我一开始用gunicorn -k sync跑FastAPI,async def路由根本不生效,全部退化到同步执行。必须用-k uvicorn.workers.UvicornWorker或者直接uvicorn app:app --workers 4。最终命令:
gunicorn app:app -k uvicorn.workers.UvicornWorker -w 4 -b 0.0.0.0:8000
坑3:asyncio.run()在FastAPI里的误用
有人会在async def路由里调用asyncio.run(fetch_json(url)),这会导致RuntimeError因为事件循环已存在。正确做法是用await直接调用,asyncio.run只在顶层入口用(比如独立脚本)。
6. 效果数据:从92 QPS到296 QPS
压测工具:wrk -t4 -c100 -d30s
| 指标 | 串行版本 | asyncio版本 | 提升 |
|---|---|---|---|
| 平均延迟 | 482ms | 172ms | 2.8x |
| P95延迟 | 612ms | 189ms | 3.2x |
| QPS | 92 | 296 | 3.2x |
| 上游连接数峰值 | 1(串行) | 20(信号量控制) | - |
补充说明:uvloop开启后,QPS从296涨到318,提升7.4%。另外把日志中的日志序列化从json.dumps换成orjson,延迟又降了5ms——这是细节优化,不展开。
生产环境运行两周,上游服务偶发抖动时,asyncio版本表现稳定:单个上游500ms超时,不影响其他三个请求结果,只是该字段返回错误占位符。而串行版本遇到超时,整个接口直接504。
7. 总结:什么时候该上asyncio
结论很明确:如果你的接口是IO密集型(HTTP调用、数据库查询、文件读写),且并发量要求高,asyncio是必选项,不是可选项。
但要说清楚,asyncio不是万能的:
- CPU密集场景(图像处理、复杂计算)用多进程,别用asyncio
- 简单接口只有一个上游调用,没必要上asyncio,线性代码更清晰
- 团队不熟悉async/await语法,容易写出阻塞代码,需要code review把关
最后留个思考题:如果你的上游服务有依赖关系(比如B需要A的返回值),怎么用asyncio实现?我的方案是用asyncio.create_task + await分阶段处理,但注意不要嵌套过深导致代码难以维护。评论区聊。