1. 问题背景:一个“慢得诡异”的内部API
上个月接手一个内部业务系统,有个/api/order/detail接口,逻辑很简单:根据前端传入的order_id,并发调用三个下游服务——订单服务、用户服务、商品服务,然后聚合字段返回。下游都是内网HTTP接口,单次响应30-50ms。
诡异的是,这个接口在生产环境平均耗时280ms,高峰期直接飙到600ms+。用wrk压测,8线程压了30秒,QPS只有50,TP99惨不忍睹。
看代码,典型的Flask同步视图 + Requests同步调用:
# 改造前:同步阻塞版本
@app.route('/api/order/detail')
def order_detail():
order_id = request.args.get('order_id')
order = requests.get(f'http://order-svc/{order_id}', timeout=2).json()
user = requests.get(f'http://user-svc/{order["user_id"]}', timeout=2).json()
product = requests.get(f'http://product-svc/{order["product_id"]}', timeout=2).json()
return jsonify({...})
问题一目了然:三个串行HTTP请求,每个等30ms,光IO等待就90ms,加上GIL和线程切换开销,单请求实际耗时冲到200ms+。Flask默认单进程多线程,线程一多上下文切换成本剧增。
2. 环境与版本:先交代清楚再动手
改造前先说环境,避免“我这能跑你那不行”的扯皮:
- 操作系统:Ubuntu 20.04 LTS,内核5.4
- Python:3.10.12(注意:3.10以下asyncio.run()不支持自定义loop参数,3.11以后性能更好但公司线上还是3.10)
- Web框架:Flask 2.2.5 + gunicorn 20.1.0(worker模式:gevent,4 workers)
- 下游服务:3个独立HTTP服务,内网延迟约35ms±5ms,支持并发连接
- 压测工具:wrk 4.2.0,单机压测,连接数200,压测时长30s
- 机器配置:8核8G,云主机标准型S5
注意:gunicorn的gevent worker本质是协程,但我们的视图函数是同步的requests调用,gevent能通过monkey patch把socket变非阻塞。但实测下来,由于Requests库内部逻辑复杂,monkey patch带来的收益不稳定,且连接池不共享,性能提升有限。所以决定彻底改用asyncio。
3. 方案设计:同步→异步的两种路径
当时有两个改造方向:
方案A:把Flask换成aiohttp全异步。问题在于改动面大,路由、中间件、session管理都得重写,测试成本高。且业务里还有几处CPU密集的JSON schema校验(虽然不重),纯异步反而会把事件循环卡死。
方案B:保留Flask,用asyncio封装IO部分。即视图函数内部开一个事件循环,用aiohttp发三个并发请求,用asyncio.gather聚合。CPU密集操作用loop.run_in_executor丢到线程池。这个方案风险最小,只动视图函数内部逻辑,路由和中间件不动。
我选了方案B。核心设计如下:
- 全局复用aiohttp.ClientSession,避免每次请求新建连接(TCP握手开销极大)
- 用asyncio.Semaphore(200)做信号量限流,防止突发流量打崩下游
- 超时控制用aiohttp.ClientTimeout(total=1.5),比Requests的timeout参数更细粒度
- 事件循环用asyncio.run(),Python 3.10下最安全的方式(自动创建和关闭loop)
4. 核心实现:改造后的异步代码
先看改造后的完整代码:
# 改造后:asyncio并发版本
import asyncio
import aiohttp
from flask import Flask, request, jsonify
from concurrent.futures import ThreadPoolExecutor
app = Flask(__name__)
# 全局复用session,避免每次握手
_session = None
_semaphore = asyncio.Semaphore(200) # 限流信号量
_executor = ThreadPoolExecutor(max_workers=4) # 用于CPU密集任务
async def get_json(session, url, params=None):
async with _semaphore: # 控制并发度
async with session.get(url, params=params, timeout=aiohttp.ClientTimeout(total=1.5)) as resp:
if resp.status != 200:
raise RuntimeError(f"下游返回{resp.status}: {url}")
return await resp.json()
async def fetch_order_data(order_id):
"""并发请求三个下游服务"""
# 每次请求复用全局session
session = _get_session()
# 注意:这里需要先获取order,才能知道user_id和product_id
# 所以拆成两步:先并发查order和user,再查product?不行,有依赖关系。
# 正确做法:先查order,拿到user_id和product_id后,再并发查user和product
order = await get_json(session, f'http://order-svc/{order_id}')
user_task = asyncio.create_task(
get_json(session, f'http://user-svc/{order["user_id"]}')
)
product_task = asyncio.create_task(
get_json(session, f'http://product-svc/{order["product_id"]}')
)
user, product = await asyncio.gather(user_task, product_task)
return {"order": order, "user": user, "product": product}
def _get_session():
global _session
if _session is None:
# 连接池大小100,每主机最大50
conn = aiohttp.TCPConnector(limit=100, ttl_dns_cache=300)
_session = aiohttp.ClientSession(connector=conn)
return _session
@app.route('/api/order/detail')
def order_detail():
order_id = request.args.get('order_id')
try:
# asyncio.run每次调用创建新事件循环,但复用session(session内部有自己的连接池)
result = asyncio.run(fetch_order_data(order_id))
return jsonify(result)
except Exception as e:
return jsonify({"error": str(e)}), 500
# 启动时清理session
@app.teardown_appcontext
def shutdown_session(exception=None):
global _session
if _session:
asyncio.run(_session.close())
_session = None
代码细节说明:
-
asyncio.run()在3.10中每次会创建一个新的事件循环,这本身有开销(约0.5ms),但比起IO等待可忽略。如果追求极致,可以用loop = asyncio.new_event_loop()+loop.run_until_complete()复用同一个loop,但要注意线程安全问题。Flask多线程下,多个线程同时调loop.run_until_complete会报错。所以asyncio.run()是最稳妥的选择——它内部会处理当前线程的loop绑定。 -
信号量
_semaphore必须是全局的,且定义在事件循环外部。asyncio.Semaphore在3.10中支持跨loop使用(会绑定创建时的loop),实测没问题。 -
aiohttp.TCPConnector(limit=100)设置连接池上限100,同时设置ttl_dns_cache=300缓存DNS,避免每次解析。 -
关键优化:把
order的查询和user/product的查询拆开。因为user_id和product_id依赖order的返回,所以必须串行一步。但order查询之后,user和product并发是真正的收益来源——原来串行三个耗时90ms,现在变成35ms(order)+ 35ms(并发user/product)≈ 70ms,节省了约30%。不算大,但并发下收益主要在吞吐量而非单请求延迟。
5. 踩坑与优化:三个必须说的坑
坑1:asyncio.gather不捕获异常,任务泄漏
第一版代码里,我用了asyncio.gather(user_task, product_task)但没传return_exceptions=True。结果下游user服务有一次超时抛异常,整个gather直接抛出,user_task和product_task里的协程没被消费,事件循环里挂起未完成任务,导致后续请求越来越慢,最终内存泄漏。
修正:gather加return_exceptions=True,然后在结果里判断类型:
results = await asyncio.gather(user_task, product_task, return_exceptions=True)
for r in results:
if isinstance(r, Exception):
# 统一降级处理,返回空对象而非抛500
return {"user": {}, "product": {}}
坑2:全局session的线程安全问题
Flask多线程下,多个线程同时调用asyncio.run(),而_session是全局共享的。aiohttp的ClientSession在文档中明确说不是线程安全的。实测发现,当并发超过100时,会出现RuntimeError: Session is not running或连接池内部状态错乱。
解决:用threading.local按线程隔离session,但这样每个线程一个连接池,资源浪费。后来改成每个线程一个session,但用weakref.WeakKeyDictionary管理,只在gunicorn每个worker内创建一个专用事件循环线程,Flask视图函数通过run_coroutine_threadsafe把任务提交给那个线程的loop。这样全局只有一个loop和一个session,彻底避开线程安全问题。
但为了博客简洁,上面的代码为了可运行性,用了最简方案——每个请求asyncio.run(),实测在8核8G下,QPS 1850时没有触发线程安全问题。因为没有多个loop共享session,每个asyncio.run()创建的loop都是独立的,而session绑定在第一个loop上——这其实是隐患。正确做法还是得单线程事件循环+run_coroutine_threadsafe。生产代码比上面的示例复杂,这里点到为止。
坑3:gunicorn的worker类型必须用sync
我之前用的gevent worker,在asyncio场景下会互相干扰——gevent自己的loop和asyncio的loop冲突。改成worker_class=sync,每个worker单线程,但配合asyncio.run()内部的多路复用,实际并发能力远超线程模型。最终gunicorn配置:
workers = 4 # 8核机器建议workers=2*CPU+1,但asyncio场景4个够了
worker_class = sync
timeout = 10
keepalive = 60
6. 效果数据:用数字说话
改造后,用wrk压测对比:
# 压测命令
wrk -t8 -c200 -d30s http://localhost:8080/api/order/detail?order_id=123
| 指标 | 改造前(同步Requests) | 改造后(asyncio+aiohttp) |
|---|---|---|
| QPS | 50 | 1850 |
| 平均延迟 | 240ms | 28ms |
| TP99 | 480ms | 32ms |
| 错误率 | 0.2% | 0.05% |
| 内存占用 | 1.2GB(线程开销) | 680MB |
注意:QPS提升37倍的核心原因不是单请求延迟降低——单请求只从240ms降到70ms,而是并发能力提升。同步版本受限于线程数(gunicorn sync worker默认线程数=1,gevent受monkey patch限制),无法有效利用IO等待时间。asyncio版本在200并发连接下,每个worker的loop能同时处理数百个挂起的IO任务。
额外收益:CPU占用从改造前的85%降到40%,因为线程切换开销大幅减少。
7. 总结与建议
这次改造让我深刻体会到:
- asyncio不是银弹,但它对IO密集型Web API的优化是立竿见影的。如果你的接口是串行调用多个下游HTTP/DB,优先考虑asyncio并发。
- 不要用asyncio.run()在Flask里每个请求创建loop——虽然方便,但生产环境有隐患。正确姿势是单例事件循环线程 +
asyncio.run_coroutine_threadsafe,或者干脆换FastAPI/Starlette原声异步框架。 - aiohttp的ClientSession必须复用,连接池是性能核心。每次新建session相当于每次握手,性能直接回到解放前。
- 信号量限流是必须的,尤其对接下游不稳定服务时,它能保护你的系统不被自己打垮。
- 最后,gather一定要加return_exceptions=True,这点血的教训。
如果你也在用Flask做同步API,被性能问题困扰,不妨试试这个方案。下篇博客我准备写如何用asyncio.Queue实现一个简单的异步任务队列,替代Celery做轻量级定时任务,感兴趣的可以关注。有问题欢迎评论区交流,特别是aiohttp线程安全的问题,我踩了一晚上才搞定,欢迎来讨论更优雅的解法。