一、问题背景:当你的API在等待中“躺平”
去年年底我负责一个数据聚合服务,接口需要同时从用户服务、订单服务和风控服务拉取数据,然后拼装返回给前端。最初版本用Flask的同步视图函数,内部直接调用Requests库:
@app.route('/api/v1/aggregate')
def aggregate():
user = requests.get('http://user-service/info', timeout=3).json()
order = requests.get('http://order-service/list', timeout=3).json()
risk = requests.get('http://risk-service/check', timeout=3).json()
return {'user': user, 'order': order, 'risk': risk}
上线第一周一切正常。随着流量增长,问题开始暴露——三个上游接口的P95响应时间分别是600ms、800ms、1200ms。由于是串行调用,单个请求最坏情况下需要等待2.6秒。压测显示:单机4核8G配置下,QPS达到120时CPU只有30%使用率,但P95延迟已经飙到9.8秒——因为线程池(Flask默认线程池大小由max_workers决定,默认值为min(32, cpu+4))中的线程全部阻塞在requests.get上,新请求只能排队。
关键问题在于:线程在等待IO时,CPU是空闲的,但GIL和线程切换成本让这种空闲无法转化为吞吐量。当时我脑子里闪过两个选择:改用Tornado或者Sanic?不,团队对Flask的生态和中间件依赖太深,迁移成本不可接受。于是决定:保留Flask路由,但内部用asyncio把三个阻塞请求变成并发协程。
二、环境与版本:Python 3.10.8 + Flask 2.2.2
改造前确认了依赖版本,这是避免踩坑的第一步:
Python 3.10.8(内置asyncio,无需额外安装)
Flask 2.2.2
aiohttp 3.8.4
gunicorn 20.1.0(生产部署)
为什么不用最新的Python 3.12?因为公司的基础镜像还停留在3.10,而且3.10中asyncio.get_event_loop()的deprecation警告已经给出,但asyncio.run()完全够用。另外,生产环境用gunicorn多worker部署(4个worker,每个worker同步模式),每个worker内部跑一个asyncio事件循环。
三、方案设计:协程不是银弹,需要控制并发度
核心思路是:用asyncio + aiohttp替代Requests,将三个独立IO请求并发执行。但直接asyncio.gather会有风险——如果上游服务突然抖动,100个并发请求瞬间打满连接池,反而拖垮自己。因此必须加asyncio.Semaphore限制并发数。
设计如下:
- 每个worker进程内,创建一个全局的
aiohttp.ClientSession(复用TCP连接池,避免每次请求都握手) - 用
Semaphore(10)控制单个请求内并发上限(因为只有3个请求,这里其实不饱和,但为了通用性保留) - 设置总超时5秒,每个子请求超时3秒,防止某个上游挂死拖垮整体
- 保持返回数据结构与同步版本完全一致,前端无感知
四、核心实现:从requests.get到async def fetch
改造后的关键代码,注意async def和await的用法,以及如何在一个同步Flask视图函数中运行协程:
import asyncio
import aiohttp
from flask import Flask, jsonify
app = Flask(__name__)
# 全局session,复用连接池
_session = None
_semaphore = asyncio.Semaphore(10)
async def get_session():
global _session
if _session is None:
_session = aiohttp.ClientSession(
timeout=aiohttp.ClientTimeout(total=5, connect=3)
)
return _session
async def fetch_json(url):
async with _semaphore:
session = await get_session()
async with session.get(url) as resp:
return await resp.json()
async def fetch_all():
# 并发发起3个请求,返回顺序与gather参数顺序一致
results = await asyncio.gather(
fetch_json('http://user-service/info'),
fetch_json('http://order-service/list'),
fetch_json('http://risk-service/check'),
return_exceptions=True # 防止某个异常导致整体崩溃
)
# 处理异常情况
for i, r in enumerate(results):
if isinstance(r, Exception):
results[i] = {'error': str(r)}
return results
@app.route('/api/v1/aggregate')
def aggregate():
# 在同步函数中运行异步代码
user, order, risk = asyncio.run(fetch_all())
return jsonify({'user': user, 'order': order, 'risk': risk})
注意asyncio.run()每次创建一个新的事件循环,这意味着每次请求都会重建session? 不,_session是全局变量,即使事件循环关闭,ClientSession对象依然存在,下次asyncio.run()时会复用。但这里有个隐患:ClientSession内部绑定的事件循环已经关闭,会导致“Event loop is closed”错误。解决方案是在fetch_all中检查session是否属于当前循环,或者用小技巧——将session创建延迟到事件循环内部,并用asyncio.Lock保护。
修正版:
_session = None
_session_lock = asyncio.Lock()
async def get_session():
global _session
async with _session_lock:
if _session is None or _session.closed:
_session = aiohttp.ClientSession()
return _session
asyncio.Lock是协程级别的锁,安全。但如果两个请求同时进入asyncio.run(),它们处于不同事件循环,_session_lock会不会跨循环冲突?会。这就是为什么不推荐asyncio.run()频繁调用的原因。更优雅的做法是用asyncio.new_event_loop()在模块加载时创建并set,但生产环境中我们最终选择了将Flask视图改为异步视图——使用flask[async]扩展,或者干脆用quart。但为了篇幅,这里展示的是可用方案,实际生产我最终用了quart替换了Flask,但那完全是另一篇文章了。
五、踩坑与优化:三次血泪教训
坑1:asyncio.run()与全局session的循环绑定。上面已提到,解决方案是用quart替代Flask(支持原生异步视图),或者每次请求手动await session.close()再新建——但那样连接池就废了。最终我选择在app.before_request中创建session,app.teardown_request中关闭,确保session生命周期与请求绑定,循环一致。
坑2:超时设置不生效。aiohttp.ClientTimeout(total=5)是总超时,但如果你用asyncio.gather并发三个请求,每个请求自己的total=5是独立计时,总时间可能达到5秒而非5秒内完成。要控制整体响应时间,需要在await asyncio.wait_for(fetch_all(), timeout=5)外层再包一层。我最终用了asyncio.wait_for确保接口P99不超过5秒。
坑3:内存泄漏。aiohttp的ClientSession如果不关闭,每请求会泄漏一个socket。压测10万次后,ss -s显示TIME_WAIT连接激增。解决:async with上下文管理器确保session自动关闭,同时设置connector=aiohttp.TCPConnector(limit=100, ttl_dns_cache=300)限制连接池大小和DNS缓存时间。
最终优化后的代码结构(quart版本核心片段):
from quart import Quart, jsonify
import aiohttp
app = Quart(__name__)
@app.route('/api/v1/aggregate')
async def aggregate():
async with aiohttp.ClientSession() as session:
async def fetch(url):
async with session.get(url, timeout=aiohttp.ClientTimeout(total=3)) as resp:
return await resp.json()
# 并发
user, order, risk = await asyncio.gather(
fetch('http://user/info'),
fetch('http://order/list'),
fetch('http://risk/check'),
return_exceptions=True
)
return jsonify({...})
六、效果数据:从9.8秒到1.2秒,CPU终于被用满
改造完成后,用locust进行10分钟压测,模拟200并发用户,对比数据如下(4核8G云主机,gunicorn 4 worker):
| 指标 | 同步Flask+Requests | 异步Quart+aiohttp |
|---|---|---|
| P95延迟 | 9.8秒 | 1.2秒 |
| P99延迟 | 14.9秒 | 2.3秒 |
| 吞吐量 | 120 QPS | 960 QPS |
| CPU利用率 | 30% | 85% |
| 内存占用 | 180MB | 210MB |
注意:P95从9.8降到1.2,不是简单的8倍,而是因为同步版本中线程阻塞导致队列排队,实际等待时间 = 上游延迟 + 排队时间。异步版本中,IO等待被事件循环让渡给其他请求,单个请求的延迟完全由上游决定(最慢的1.2秒)。吞吐量提升8倍是因为单worker能处理的并发请求数从max_workers=4(默认)提升到数千协程,CPU利用率从30%升到85%,说明资源被真正利用。
七、总结与建议
- 不是所有API都适合异步:如果你的内部没有IO等待(纯CPU计算),异步反而因协程切换损失性能。用
cProfile先分析瓶颈。 - asyncio不是线程替代品:它解决的是IO并发,不是CPU并行。多核利用仍需多进程(gunicorn多worker)。
- 生产环境优先用
quart而非flask+asyncio.run:asyncio.run每次创建事件循环的开销约1ms,且不易管理全局资源。如果必须用Flask,考虑gunicorn的geventworker而不是asyncio。 - 性能数字依赖上游:我们的上游服务P95是1.2秒,如果你的上游是2秒,异步后P95约2秒,不会更慢——但吞吐量会提升到相近水平。
最后,异步编程的核心是让等待不阻塞。在IO密集场景下,它比多线程更轻量、比多进程更省内存。但请记住:没有银弹,只有适合场景的选择。如果你的API还在串行等待,不妨试试asyncio,但记得控制好并发度和超时。