一、问题背景:一个看似“没问题”的聚合API
我们有个订单服务,前端需要一次性展示订单详情、用户信息、物流轨迹、优惠券状态。最初实现是同步串行调用四个下游服务:
def get_order_detail(order_id):
user = requests.get(f"http://user-service/user/{order_id}").json()
logistics = requests.get(f"http://logistics-service/track/{order_id}").json()
coupon = requests.get(f"http://coupon-service/coupon/{order_id}").json()
...
线上运行半年,日均调用量不大(约200万),一直没出问题。直到双十一前压测,50并发下QPS 37,P99 4.8s——这数据根本没法上线。四个下游服务响应时间都在200-500ms,串行累加后单请求耗时轻松超过1.5s。
二、环境与版本:先交代清楚
- Python 3.10.11(3.10+的asyncio才比较稳定,3.8的坑太多)
- Flask 2.1.2(注意:Flask本身是WSGI同步框架,asyncio需要自己桥接)
- aiohttp 3.8.4(替代requests的关键)
- wrk 4.2.0(压测工具)
- 服务器:8核8G CentOS 7.9,下游服务在同一内网,延迟约50ms
三、方案设计:不是简单替换,而是分层改造
直接上asyncio不能解决所有问题。我拆了三步:
- 把四个串行IO改为并发协程:用
asyncio.gather同时发起四个请求 - 限制并发度:防止突发流量把下游打爆,用
Semaphore(100)做信号量 - 超时熔断:每个请求用
asyncio.wait_for设置500ms超时,超时直接返回降级数据
核心设计原则:异步化不改变业务逻辑,只改变IO等待方式。下游接口不兼容协程没关系,aiohttp直接替代requests,返回体结构相同。
四、核心实现:before/after代码对比
Before:同步串行版(50行)
# app.py - 同步版本
import requests
from flask import Flask, jsonify
app = Flask(__name__)
def fetch_user(order_id):
return requests.get(f"http://user-service/user/{order_id}", timeout=1).json()
def fetch_logistics(order_id):
return requests.get(f"http://logistics-service/track/{order_id}", timeout=1).json()
def fetch_coupon(order_id):
return requests.get(f"http://coupon-service/coupon/{order_id}", timeout=1).json()
@app.route("/order/")
def get_order(order_id):
user_data = fetch_user(order_id)
logistics_data = fetch_logistics(order_id)
coupon_data = fetch_coupon(order_id)
return jsonify({
"user": user_data,
"logistics": logistics_data,
"coupon": coupon_data
})
After:异步并发版(80行)
# app.py - 异步版本
import asyncio
import aiohttp
from flask import Flask, jsonify
app = Flask(__name__)
semaphore = asyncio.Semaphore(100) # 控制最大并发100
async def fetch_json(session, url, timeout=0.5):
"""带超时和信号量的异步GET"""
async with semaphore:
try:
async with session.get(url, timeout=aiohttp.ClientTimeout(total=timeout)) as resp:
return await resp.json()
except (asyncio.TimeoutError, aiohttp.ClientError):
return {"error": "service_unavailable"} # 降级返回
async def fetch_all(order_id):
"""并发获取所有下游数据"""
timeout = aiohttp.ClientTimeout(total=1.0)
async with aiohttp.ClientSession(timeout=timeout) as session:
user_task = fetch_json(session, f"http://user-service/user/{order_id}")
logistics_task = fetch_json(session, f"http://logistics-service/track/{order_id}")
coupon_task = fetch_json(session, f"http://coupon-service/coupon/{order_id}")
results = await asyncio.gather(user_task, logistics_task, coupon_task)
return results
@app.route("/order/")
def get_order(order_id):
# 关键:在同步WSGI中运行协程
user_data, logistics_data, coupon_data = asyncio.run(fetch_all(order_id))
return jsonify({
"user": user_data,
"logistics": logistics_data,
"coupon": coupon_data
})
五、踩坑与优化:三个必须说的坑
坑1:Flask + asyncio.run 的事件循环冲突
第一次测试时,我用loop = asyncio.new_event_loop()手动管理,结果每次请求都报“Event loop is closed”。原因是Flask多线程处理请求时,每个线程都尝试创建自己的事件循环。解决方案很简单:直接asyncio.run(),它会自动创建新循环并关闭旧循环。但注意:如果后续要升级到FastAPI或Starlette,就不要这么写了,直接用asyncio.get_running_loop()。
坑2:aiohttp连接池泄漏
如果每次请求都新建ClientSession,压测时会出现Connection pool is full。正确做法:把ClientSession定义为全局单例,用async with管理生命周期。但全局session需要在应用启动时创建,在请求中复用。我用了模块级变量:
session = None
def get_session():
global session
if session is None:
session = aiohttp.ClientSession()
return session
坑3:超时设置要双重保险
aiohttp.ClientTimeout(total=1.0)是总超时,但TCP连接建立超时默认是5秒。如果下游服务宕机,总超时可能不起作用。必须显式设置connect=0.5:
timeout = aiohttp.ClientTimeout(total=1.0, connect=0.5)
六、效果数据:压测对比
用wrk压测,50并发跑60秒,结果如下:
| 指标 | 同步版本 | 异步版本 | 提升倍数 |
|---|---|---|---|
| QPS | 37 | 312 | 8.4x |
| P99延迟 | 4.8s | 410ms | 11.7x |
| 平均延迟 | 1.2s | 320ms | 3.75x |
| 错误率 | 0% | 0.3%(超时降级) | - |
注意错误率:异步版本有0.3%的请求因为超时降级返回了错误数据。这是预期的——我把超时设为500ms,但下游偶尔有1s以上的毛刺。降级方案是:返回空对象并记录日志,前端根据error字段提示“部分信息加载失败”。
七、总结:什么时候该用asyncio?
不是所有场景都适合asyncio。我的经验是:
- IO密集型(大量网络请求、数据库查询)→ 必须用asyncio,效果立竿见影
- CPU密集型(计算、加密、序列化)→ 用多进程,asyncio反而拖慢速度
- 混合负载 → 用asyncio + concurrent.futures.ThreadPoolExecutor组合
这次改造最值的点在于:只改了IO层,业务逻辑零变化。如果你也在用Flask + Requests做聚合API,建议先试试asyncio + aiohttp。如果追求更高性能,可以直接迁移到FastAPI(原生支持async def),但那是另一个故事了。
最后提醒:压测数据仅供参考,真实环境的网络波动、下游容量都会影响结果。改造前一定要先压测,改造后一定要灰度发布。