一、问题背景:一个拖垮整个服务的"简单接口"
上个月接手一个老项目,Flask 1.1.2 + Gunicorn 20.1.0,部署在K8s里。运维反馈说有个 /api/v1/order/summary 接口经常触发存活探针失败,因为响应时间超过3秒直接被kill。
看代码后发现,这个接口的逻辑很直白:前端请求过来后,它需要先后调用内部的用户服务(获取用户等级)、订单服务(获取订单列表)、营销服务(计算优惠)。三个调用之间没有数据依赖,但代码是串行写的:
# 改造前:同步串行调用
@app.route('/api/v1/order/summary')
def get_summary():
user_level = requests.get('http://user-svc/level', timeout=2).json()
order_list = requests.get('http://order-svc/list', timeout=2).json()
discount = requests.get('http://marketing-svc/discount', timeout=2).json()
return jsonify({'level': user_level, 'orders': order_list, 'discount': discount})
三个内部服务平均延迟都在800ms~1s左右,串起来就是2.5~3秒。这还没完,Gunicorn配的是同步worker,workers=4,每个worker同一时刻只能处理一个请求。一旦并发上来,所有worker都被阻塞在这个接口上,其他接口全部排队。
压测数据:用wrk -t4 -c100 -d30s 打这个接口,QPS只有120,P95延迟3100ms,P99直接4000ms+。CPU利用率虽然只有60%,但线程全部卡在requests.get的IO等待上——典型的同步IO阻塞。
二、环境与版本:老项目改造的典型配置
在动手前先交代一下环境,这很重要因为后面遇到的坑跟版本强相关:
- Python: 3.8.10(项目锁定,不能升级到3.10+)
- Flask: 1.1.2(框架版本老,没有原生async支持)
- aiohttp: 3.8.4(用于替代requests做异步HTTP调用)
- Gunicorn: 20.1.0(最终配合uvicorn worker使用)
- 部署:K8s单Pod,4核8G,三个内部服务平均响应时间约800ms
这里有个关键决策:不换框架。虽然FastAPI天生支持async,但项目里还有几百个同步接口,全量迁移成本太高。所以采用Flask + asyncio.run() 的兼容方案,只改造这一个热点接口。
三、方案设计:asyncio.gather + aiohttp连接池
核心思路很简单:三个内部调用互相独立,用asyncio.gather并发执行,总耗时从三个之和变成最慢的一个。
但有几个细节必须处理:
-
aiohttp.ClientSession必须复用。如果每次请求创建新session,TCP握手开销会让性能比requests还差。session内部维护连接池,默认100个连接,超时回收时间10秒。
-
信号量限流。内部服务不经打,如果这个接口被刷到1000QPS,每个请求并发3个内部调用,就是3000QPS打到下游。必须用
asyncio.Semaphore限制最大并发数,我这边设了200。 -
超时控制。requests的timeout参数是每次调用的,aiohttp里需要
ClientTimeout(total=2)。 -
Gunicorn worker类型。同步worker不能跑async代码,必须换成
uvicorn.workers.UvicornWorker。但它要求WSGI应用是ASGI兼容的,Flask不行。所以最终方案是:保留Gunicorn管理进程,但用meinheld或gunicorn[gevent]?不,我用的是单独开一个aiohttp的web服务?这样太复杂。
最终选了最简单可靠的方案:在Flask视图函数里用asyncio.run()启动一个临时事件循环,内部跑gather。代价是每次请求会创建/销毁事件循环(有开销,但实测下来比串行节省的时间多得多)。
四、核心实现:改造后的异步代码
import asyncio
import aiohttp
import time
from flask import Flask, jsonify
app = Flask(__name__)
# 全局session,进程内复用连接池
_session = None
async def get_session():
global _session
if _session is None or _session.closed:
timeout = aiohttp.ClientTimeout(total=2)
_session = aiohttp.ClientSession(timeout=timeout)
return _session
async def fetch_json(session, url):
"""带信号量和重试的GET请求"""
async with semaphore:
try:
async with session.get(url) as resp:
if resp.status == 200:
return await resp.json()
else:
return {'error': f'status_{resp.status}'}
except asyncio.TimeoutError:
return {'error': 'timeout'}
except aiohttp.ClientError as e:
return {'error': str(e)}
# 全局信号量,限制并发200
semaphore = asyncio.Semaphore(200)
async def async_summary():
session = await get_session()
urls = [
'http://user-svc/level',
'http://order-svc/list',
'http://marketing-svc/discount'
]
# 并发执行三个请求
results = await asyncio.gather(
*(fetch_json(session, url) for url in urls),
return_exceptions=False
)
return {
'level': results[0],
'orders': results[1],
'discount': results[2]
}
@app.route('/api/v1/order/summary')
def get_summary_async():
start = time.perf_counter()
# 每次请求创建新的事件循环
data = asyncio.run(async_summary())
print(f'async total: {(time.perf_counter() - start)*1000:.1f}ms')
return jsonify(data)
注意asyncio.run()每次会创建新事件循环,这有个代价:Semaphore在async_summary内部创建,每次都是全新的,所以信号量限流失效了。正确做法是把Semaphore和Session都定义在模块级别,让它们跨事件循环复用。
上面代码有个隐藏bug:Semaphore绑定到事件循环,如果在循环外创建会报错。需要这样改:
# 正确做法:在第一次进入循环时创建
semaphore = None
async def get_semaphore():
global semaphore
if semaphore is None:
semaphore = asyncio.Semaphore(200)
return semaphore
然后在async_summary里调用await get_semaphore()。
五、踩坑与优化:两个让我抓狂的问题
坑1:aiohttp连接池耗尽(ClientOSError: [Errno 24] Too many open files)
改造后压测到500QPS时,日志开始刷这个错。原因:ClientSession默认连接池大小是100,但asyncio.run()每次创建新循环,旧循环里的session被丢弃但没有正确关闭,导致文件描述符泄漏。
解决:_session全局复用还不够,必须在进程退出时显式关闭。另外把连接池大小调大:
connector = aiohttp.TCPConnector(
limit=100, # 每主机连接数
limit_per_host=30, # 每主机并发限制
force_close=False, # 保持keep-alive
enable_cleanup_closed=True # 关键!清理半关闭连接
)
_session = aiohttp.ClientSession(connector=connector)
坑2:Gunicorn同步worker跑asyncio.run()会导致事件循环冲突
如果你直接在同步worker里用asyncio.run(),当并发请求进来时,每个worker同时执行asyncio.run(),内部会创建多个事件循环。这在CPython里没问题,但如果有稍微老一点的第三方库绑定了loop,会报RuntimeError: Event loop is closed。
解决:最终我把Gunicorn worker换成gthread(线程模式),每个worker内部用线程池处理请求,每个线程自己跑asyncio.run()。配置如下:
gunicorn -w 2 --threads 8 -k gthread -b 0.0.0.0:8000 app:app
实测2个worker × 8线程,每线程一个独立事件循环,完美避开了loop冲突。
六、效果数据:延迟降88%,QPS翻7倍
改完后再跑wrk压测,同样的参数wrk -t4 -c100 -d30s http://localhost:8000/api/v1/order/summary:
| 指标 | 改造前 | 改造后 | 提升 |
|---|---|---|---|
| QPS | 120 | 850 | 608% |
| P50延迟 | 780ms | 380ms | 51%↓ |
| P95延迟 | 3100ms | 420ms | 86%↓ |
| P99延迟 | 4200ms | 550ms | 87%↓ |
| CPU峰值 | 65% | 50% | 15%↓ |
为什么P50只从780降到380?因为下游三个服务本身平均就有800ms延迟,并行后理论极限是800ms。但实测P50只有380ms,说明内部服务响应比想象中快(平均200ms),串行时额外等待时间主要是网络往返和线程切换。
额外收益:由于不再阻塞worker线程,其他接口的延迟也降了。比如/api/v1/user/info原本偶尔会超时(因为worker被占满),现在P99从2000ms降到200ms。
七、总结:什么时候值得用asyncio
这次改造的核心收益不是因为async比同步快,而是消除了串行等待。如果三个调用有依赖关系,异步也救不了你。另外注意几点:
asyncio.run()每次创建事件循环有性能开销(实测约0.5ms),对于高QPS接口建议用asyncio.new_event_loop()复用循环,但要注意线程安全问题。- 信号量限流必须全局复用,否则等于没限。
- aiohttp的session生命周期管理是最大的坑,建议封装成模块级单例,配合
atexit注册关闭钩子。 - 如果项目允许升级,FastAPI + httpx.AsyncClient是更干净的方案,但老项目用Flask+asyncio.run()过渡完全可行。
最后提醒:异步改造不是银弹,如果你只有一个IO调用或者调用有依赖,同步代码反而更简单。性能优化第一原则永远是——先测量,再动手。