1. 问题背景:IO密集接口为何成了性能瓶颈
上个月接手一个给运营看板用的聚合接口 /api/summary,逻辑很简单——拿到用户请求参数后,依次调用用户服务、订单服务、风控服务三个RPC接口,最后把结果拼装返回。
三个下游接口的耗时分别是:用户服务300ms,订单服务450ms,风控服务400ms。串行调用总耗时大概1.2秒。压测时用4个线程跑,QPS死活上不去,CPU占用不到10%,但响应时间就是下不来。
我第一反应是加缓存,但产品要求数据准实时,缓存最多30秒就过期,而且下游数据变更频繁,缓存命中率低。后来看了火焰图,发现线程几乎全在等待Socket接收数据——典型的IO密集型瓶颈。
当时线上Python版本是3.8.10,用Flask 2.0.2 + requests发HTTP调用。requests的同步阻塞模型决定了每个请求要占一个线程,而GIL又让线程切换开销不小。这种场景,就是asyncio最对口的问题。
2. 环境与版本:改造前的技术栈
| 组件 | 版本 |
|---|---|
| Python | 3.8.10(后续升级到3.10.11) |
| Web框架 | Flask 2.0.2 |
| HTTP客户端 | requests 2.27.1 |
| 部署方式 | Gunicorn 20.1.0 + 4 workers |
| 压测工具 | wrk 4.2.0 |
| 下游服务 | 三个内部RPC(HTTP JSON协议) |
这里要说明一点,Flask本身是同步框架,直接在视图函数里跑asyncio代码会有事件循环冲突。这个细节后面会细说,也是很多新手会踩的坑。
3. 方案设计:协程化改造的核心思路
我的改造方案分三步:
- 用httpx.AsyncClient替换requests,创建3个并发任务分别调用下游服务
- 用asyncio.gather + return_exceptions 做并发控制,一个服务挂了不拖垮整体
- 用信号量(Semaphore)限制最大并发数,防止下游服务被冲垮
整体代码结构如下:
# before: 同步串行版
def get_summary_sync(user_id):
user_info = requests.get(f"http://user-service/api/user/{user_id}", timeout=3).json()
order_info = requests.get(f"http://order-service/api/orders?user_id={user_id}", timeout=3).json()
risk_info = requests.get(f"http://risk-service/api/risk/{user_id}", timeout=3).json()
return {
"user": user_info,
"orders": order_info,
"risk_score": risk_info["score"]
}
同步版代码逻辑清晰,但三个请求串行执行,总耗时等于三次网络往返之和。如果其中一个下游超时3秒,整个接口就跟着卡3秒。
4. 核心实现:asyncio并发改造的具体代码
改造后的异步版本核心逻辑:
# after: asyncio并发版
import asyncio
import httpx
from flask import Flask, jsonify
app = Flask(__name__)
# 全局复用AsyncClient,避免每次请求都创建新连接池
_client = httpx.AsyncClient(
timeout=httpx.Timeout(3.0, connect=1.0),
limits=httpx.Limits(max_connections=200, max_keepalive_connections=50)
)
# 信号量:限制同时最多50个并发请求发往下游
_semaphore = asyncio.Semaphore(50)
async def fetch_json(client, url):
async with _semaphore:
resp = await client.get(url)
resp.raise_for_status()
return resp.json()
async def get_summary_async(user_id):
urls = [
f"http://user-service/api/user/{user_id}",
f"http://order-service/api/orders?user_id={user_id}",
f"http://risk-service/api/risk/{user_id}"
]
# 并发发起3个请求,单个失败不影响其他结果
results = await asyncio.gather(
*(fetch_json(_client, url) for url in urls),
return_exceptions=True
)
user_info, order_info, risk_info = results
# 处理异常情况
user_data = user_info if isinstance(user_info, dict) else {"error": "user service failed"}
order_data = order_info if isinstance(order_info, dict) else {"error": "order service failed"}
risk_score = risk_info.get("score", 0) if isinstance(risk_info, dict) else 0
return {
"user": user_data,
"orders": order_data,
"risk_score": risk_score
}
@app.route("/api/summary")
async def summary():
user_id = request.args.get("user_id")
if not user_id:
return jsonify({"error": "missing user_id"}), 400
data = await get_summary_async(user_id)
return jsonify(data)
注意两点:一是 httpx.AsyncClient 要全局复用,不能每次请求都新建,否则TCP连接没法复用,性能反而更差;二是 Semaphore(50) 这个值是我压测调出来的,下游服务只能扛住约80QPS,设成200直接把它打挂了,后面会详细说。
5. 踩坑与优化:Flask异步视图的兼容问题
第一个坑:Flask + asyncio不兼容。上面代码直接用是跑不起来的,因为Flask的视图函数是同步执行的,它不知道如何处理协程对象。我一开始直接运行,报错 TypeError: 'coroutine' object is not callable。
解决方法是安装 flask[async] 扩展,或者用 asgiref 的 sync_to_async 包装。我选的是第二个方案,因为不想动Flask版本:
# 用asgiref包装异步视图
from asgiref.sync import async_to_sync
@app.route("/api/summary")
def summary_sync():
user_id = request.args.get("user_id")
if not user_id:
return jsonify({"error": "missing user_id"}), 400
data = async_to_sync(get_summary_async)(user_id)
return jsonify(data)
第二个坑:超时控制不当导致雪崩。最开始我没设 connect=1.0,只设了 timeout=3.0。结果下游用户服务在某次发布时连接不释放,所有请求都卡在connect阶段,直接把我的worker线程池占满了,接口全部504。
后来加上连接超时 + 总超时双重控制,并且对每个子任务单独设置超时,避免一个慢服务拖垮整体。
第三个坑:Gunicorn worker类型必须改。Flask同步worker里跑协程没意义,要把worker类型改成 gthread 或者直接用 uvicorn。我最终用的是 gunicorn -k gthread --threads 8,让每个worker内部用线程池调度协程。
6. 效果数据:改造前后的性能对比
压测环境:本机4核8G,wrk发压5分钟,HTTP keep-alive开启。数据如下:
| 指标 | 同步版(before) | 异步版(after) | 提升幅度 |
|---|---|---|---|
| 平均响应时间 | 1.2s | 260ms | 78%↓ |
| P95响应时间 | 2.1s | 380ms | 82%↓ |
| QPS(4并发) | 120 | 850 | 608%↑ |
| 超时率(>3s) | 5.2% | 0.01% | - |
压测时同步版4个线程跑到第2分钟就开始出现大量超时,因为requests的线程被阻塞,无法及时释放连接。异步版用同样4个并发连接,由于事件循环调度,可以同时处理上百个在途请求。
另外我特意测试了单下游故障场景:模拟风控服务宕机,同步版整个接口按3秒超时等待,异步版因为有 return_exceptions=True,其他两个服务结果照常返回,风控部分返回默认值,整体响应时间保持在300ms以内。
7. 总结与建议
这次改造让我对asyncio有了实感:它解决的不是计算密集型问题,而是让IO等待时间重叠起来。如果你的接口瓶颈在外部调用次数多、单次耗时稳定、并发量上不去,用asyncio换掉阻塞IO是最直接的收益。
最后几点建议:
- 不要盲目追求全异步,如果接口里有大量CPU计算,协程反而帮倒忙
- 一定要用信号量限流,否则下游服务会被你的并发打爆(我亲测打挂了风控服务一次)
- 监控要跟上,我加了每个下游子调用耗时的埋点日志,后续排查问题全靠它
- Python 3.8的asyncio够用,但3.10+的 asyncio.timeout 上下文管理器更好用,建议升级
如果有人问为什么不用FastAPI直接上异步?我只能说,老项目迁移成本摆在那,用Flask + 协程包装已经是最小改动方案了。