一、问题背景:一个慢查询API拖垮了整个微服务
去年底我们接手了一个内部数据中台项目,其中一个名为/api/v1/batch-query的接口,接收一组用户ID,并行请求三个下游服务(用户画像、行为标签、信用评分),聚合后返回。
最初版本用requests库发起HTTP请求,在并发数达到50时,接口响应时间急剧恶化。生产环境中,该接口的TP99长期维持在2.3s以上,高峰期甚至导致上游服务雪崩。监控显示,单次请求中,网络IO等待时间占比超过80%。这意味着,99%的CPU时间都在等网络响应。
显然,这是一个典型的IO密集型场景,非常适合用asyncio进行优化。
二、环境与版本
本次优化的运行环境及关键依赖版本如下:
| 组件 | 版本 |
|---|---|
| Python | 3.10.12 |
| FastAPI | 0.104.1 |
| Uvicorn | 0.24.0 |
| aiohttp | 3.9.1 |
| httpx | 0.25.2 |
| uvloop | 0.19.0 |
| 操作系统 | Ubuntu 22.04 LTS |
| 压测工具 | wrk 4.2.0 |
服务部署在4核8G的云服务器上,下游三个服务均为同一机房内的HTTP接口,平均延迟约80ms。
三、方案设计:用asyncio替换同步阻塞
核心思路:将每个下游请求封装为异步协程,通过asyncio.gather并发调度,同时复用TCP连接以减少握手开销。
具体设计如下:
1. 使用aiohttp.ClientSession作为连接池,设置连接池大小=200,每主机最大连接数=100,启用keepalive。
2. 引入uvloop替换默认的asyncio事件循环,利用libuv提升事件调度性能。
3. 设置合理的超时:连接超时5s,读取超时10s,避免单个慢请求阻塞整个协程组。
4. 控制并发量:使用asyncio.Semaphore限制同时发起的请求数,防止下游被打爆。
四、核心实现:Before vs After代码
Before(同步版本)
# sync_batch_query.py (优化前)
import requests
from fastapi import FastAPI, HTTPException
from typing import List, Dict
app = FastAPI()
TIMEOUT = (5, 10) # connect, read timeout
@app.post("/api/v1/batch-query")
async def batch_query(user_ids: List[str]) -> Dict:
try:
# 串行请求三个下游服务
profile = requests.post(
"http://profile-service/get",
json={"ids": user_ids},
timeout=TIMEOUT
)
behavior = requests.post(
"http://behavior-service/get",
json={"ids": user_ids},
timeout=TIMEOUT
)
credit = requests.post(
"http://credit-service/get",
json={"ids": user_ids},
timeout=TIMEOUT
)
return {
"profiles": profile.json(),
"behaviors": behavior.json(),
"credits": credit.json()
}
except requests.Timeout:
raise HTTPException(status_code=504, detail="下游服务超时")
问题分析:虽然FastAPI的handler函数写成了async def,但内部使用的是同步requests库。当三个requests.post被执行时,当前线程会阻塞等待IO,实际上还是串行执行。FastAPI的异步能力在此处完全被浪费。
After(异步版本)
# async_batch_query.py (优化后)
import asyncio
import aiohttp
from fastapi import FastAPI, HTTPException
from typing import List, Dict
app = FastAPI()
# 全局连接池,避免每次请求都新建
CONNECTOR = aiohttp.TCPConnector(
limit=200, # 连接池总大小
limit_per_host=100, # 每主机最大连接数
ttl_dns_cache=300, # DNS缓存5分钟
force_close=False, # 保持keepalive
)
# 信号量控制并发数,防止下游雪崩
SEMAPHORE = asyncio.Semaphore(50)
async def fetch(session: aiohttp.ClientSession, url: str, payload: Dict) -> Dict:
async with SEMAPHORE:
try:
async with session.post(url, json=payload, timeout=aiohttp.ClientTimeout(
total=15, connect=5, sock_read=10
)) as resp:
resp.raise_for_status()
return await resp.json()
except asyncio.TimeoutError:
raise HTTPException(status_code=504, detail=f"服务超时: {url}")
except aiohttp.ClientError as e:
raise HTTPException(status_code=502, detail=str(e))
@app.post("/api/v1/batch-query")
async def batch_query(user_ids: List[str]) -> Dict:
async with aiohttp.ClientSession(connector=CONNECTOR) as session:
tasks = [
fetch(session, "http://profile-service/get", {"ids": user_ids}),
fetch(session, "http://behavior-service/get", {"ids": user_ids}),
fetch(session, "http://credit-service/get", {"ids": user_ids}),
]
# 并发执行三个协程
results = await asyncio.gather(*tasks, return_exceptions=True)
# 检查是否有异常
for r in results:
if isinstance(r, HTTPException):
raise r
return {
"profiles": results[0],
"behaviors": results[1],
"credits": results[2]
}
关键改动:
1. aiohttp.ClientSession复用全局连接池,避免每次请求都三次握手。
2. asyncio.gather并发调度三个协程,IO等待时自动切换。
3. asyncio.Semaphore(50)控制同时进行的对外请求数,防止下游服务因突发流量被打挂。
4. 使用raise_for_status()和自定义异常处理,保持错误信息清晰。
五、踩坑与优化
踩坑1:忘了替换事件循环
第一次测试时发现性能只提升了30%,排查后才发现,uvloop需要手动设置:
import uvloop
asyncio.set_event_loop_policy(uvloop.EventLoopPolicy())
在旧版Python中,默认事件循环是SelectorEventLoop,性能远不如uvloop。使用uvloop后,事件循环的epoll调度效率提升约40%。
踩坑2:连接池耗尽导致连锁失败
初期没有使用Semaphore,结果在压测时,200并发请求瞬间创建了600个到下游的TCP连接(每个请求3个下游),下游服务扛不住,出现大量Connection Refused。加上Semaphore后,将并发数控制在50,下游服务稳定运行。
踩坑3:aiohttp的timeout参数位置
aiohttp的timeout必须通过aiohttp.ClientTimeout对象传入,而不是简单的元组。开始时误写成session.post(url, timeout=(5,10)),结果timeout完全无效,因为aiohttp不识别元组格式。
优化:使用httpx替代aiohttp
如果项目要求更高,可以考虑httpx(0.25.2版本),它同时支持同步和异步API,且API设计更贴近requests。但aiohttp在连接池控制上更精细,适合高并发场景。
六、效果数据:从380到1600 QPS
使用wrk进行压测,命令如下:
wrk -t4 -c200 -d30s --latency http://target-host/api/v1/batch-query
压测结果对比:
| 指标 | 同步版本 | 异步版本 | 提升倍数 |
|---|---|---|---|
| QPS | 380 | 1600 | 4.2x |
| TP50 | 480ms | 120ms | 4.0x |
| TP99 | 2.3s | 450ms | 5.1x |
| 错误率(5xx) | 3.2% | 0.1% | - |
CPU利用率:同步版本在50并发时CPU利用率就达到85%(大量上下文切换),异步版本在200并发时CPU利用率仅60%(协程切换开销远小于线程切换)。
内存占用:同步版本每个请求占用约3MB(线程栈),异步版本每个请求仅占用约50KB(协程栈),内存消耗降低60倍。
七、总结
-
asyncio不是银弹:它只对IO密集型任务有效。如果你的API涉及大量CPU计算(如加密、图像处理),应该考虑多进程或线程池。
-
连接池和并发控制是核心:没有连接池的异步编程,就像没有缓冲区的管道,流量一冲就垮。asyncio.Semaphore是保护下游服务的必备工具。
-
版本兼容性很重要:Python 3.10+的asyncio已经比较成熟,但如果使用3.7或3.8,需要注意一些API差异(如asyncio.run底层实现不同)。
-
压测数据要说谎:务必在生产环境流量模式下进行压测,单纯用wrk压本地接口,忽略了网络抖动、下游限流等因素,容易得出过于乐观的结论。
最终,这个优化版本上线后,该接口的TP99稳定在500ms以下,再也没有出现过因该接口导致的上游雪崩。异步编程的核心价值,不是让代码更快,而是让等待不再阻塞资源。
如果你们也有类似的阻塞IO噩梦,不妨试试asyncio。记住:别用同步库写异步函数,那只是心理安慰。