一、问题背景:同步IO阻塞了整个Worker
上个月接手一个内部运维工单系统,有个接口 /api/tickets/detail,前端反应特别慢。看了下代码,逻辑很简单:查数据库拿工单列表,然后对每个工单去调用外部监控系统的API获取状态。
核心代码长这样:
# 改造前:同步顺序调用
def get_ticket_detail(ticket_id):
db_info = query_db(ticket_id) # 30ms
monitor_info = call_monitor_api(ticket_id) # 500-800ms
return merge(db_info, monitor_info)
@app.route('/api/tickets')
def tickets():
ticket_ids = get_all_ticket_ids() # 50ms
result = [get_ticket_detail(tid) for tid in ticket_ids]
return jsonify(result)
问题很明显:每个工单都要串行等待外部API响应。20个工单就是20 × 600ms = 12秒。虽然实际场景中工单数量平均5个,但高峰期能到30个,接口直接超时。
线上用gunicorn开了4个worker,每个worker处理一个请求时都在傻等IO。压测结果:平均耗时2.8秒,QPS 320,CPU利用率只有15%——大量时间在等IO,资源全浪费了。
二、环境与版本
改造前先把环境说清楚,避免版本差异踩坑:
- Python 3.10.12(asyncio在3.10后成熟很多,3.7之前的写法别参考)
- Flask 2.3.3(用
flask[async]扩展包支持异步视图) - aiohttp 3.9.1(替代requests做异步HTTP调用)
- gunicorn 21.2.0(worker模式换成
uvicorn或者用asyncio模式) - 压测工具:wrk 4.2.0,单线程200连接压测10秒
注意:Flask本身是WSGI同步框架,不能直接跑asyncio。两个方案:一是用flask[async]的异步视图(配合asgiref的async_to_sync),二是把异步逻辑独立成协程,通过asyncio.run()在同步视图里调用。我选了第二种,因为改动最小,且兼容现有中间件。
三、方案设计:asyncio.gather并发替代顺序调用
核心思路:把get_ticket_detail从同步函数改为async def协程函数,内部用aiohttp发起异步HTTP请求。然后在视图函数中:
- 用
asyncio.run()启动事件循环 - 用
asyncio.gather()并发执行所有工单的协程 - 用
asyncio.Semaphore(5)控制并发数,防止把外部监控系统打崩
为什么用Semaphore?因为外部监控API是第三方服务,以前同步请求每秒最多10个;如果异步后200个连接同时打过去,对方的限流策略会直接返回503。所以必须限制并发。
四、核心实现:改造后的完整代码
先看改造后的协程函数:
# 改造后:异步协程 + 信号量限流
import asyncio
import aiohttp
from asyncio import Semaphore
# 全局信号量,限制并发5个
_sem = Semaphore(5)
async def fetch_monitor_info(session, ticket_id):
"""异步获取监控状态,带超时和重试"""
url = f"http://monitor-api.internal/status/{ticket_id}"
async with _sem: # 并发控制
for attempt in range(3):
try:
async with session.get(url, timeout=aiohttp.ClientTimeout(total=2)) as resp:
if resp.status == 200:
return await resp.json()
except (aiohttp.ClientError, asyncio.TimeoutError) as e:
if attempt == 2:
return {"ticket_id": ticket_id, "status": "unknown", "error": str(e)}
await asyncio.sleep(0.2 * (attempt + 1)) # 退避重试
return None
async def get_ticket_detail_async(session, ticket_id):
"""并发的工单详情聚合"""
db_info = await asyncio.to_thread(query_db, ticket_id) # 数据库是同步的,丢线程池
monitor_info = await fetch_monitor_info(session, ticket_id)
return {**db_info, "monitor": monitor_info}
def get_all_tickets_async(ticket_ids):
"""同步入口,内部启动事件循环"""
async def runner():
async with aiohttp.ClientSession() as session:
tasks = [get_ticket_detail_async(session, tid) for tid in ticket_ids]
results = await asyncio.gather(*tasks, return_exceptions=True)
return results
return asyncio.run(runner())
视图函数变成:
@app.route('/api/tickets')
def tickets():
ticket_ids = get_all_ticket_ids()
results = get_all_tickets_async(ticket_ids)
return jsonify(results)
关键点:
- asyncio.to_thread 把同步的数据库查询丢到线程池,不阻塞事件循环
- aiohttp.ClientSession 复用连接池,避免每次新建TCP连接
- return_exceptions=True 防止单个工单异常导致整体失败
- 超时设为2秒,重试3次,避免某个接口卡死整个请求
五、踩坑与优化:三个真实问题
坑1:Event Loop is closed 异常
第一次跑的时候,gunicorn报错 RuntimeError: Event loop is closed。原因:asyncio.run()每次都会创建并关闭一个新事件循环,但aiohttp的ClientSession内部持有连接池,跨事件循环使用会出问题。解决方案:把ClientSession的创建放在runner()内部,保证session和loop同生命周期。代码里已经是这样了,但初版是把session存成全局变量,踩了这个坑。
坑2:数据库连接池被打满
因为asyncio.to_thread会把同步DB操作丢进线程池,默认线程池大小是min(32, cpu+4)。如果并发工单数超过线程池大小,任务会排队。我在gunicorn配置里加了--threads 16,并显式设置asyncio.to_thread的executor参数,限制DB并发为8个。
坑3:外部API的限流
不加Semaphore时,压测100个并发请求,外部监控API直接返回429。加了Semaphore(5)后,稳定在5个并发,对方服务无压力。但要注意:Semaphore必须在事件循环内部创建,且全局复用。我一开始写在函数内部,导致每个请求都新建信号量,完全失效。
六、效果数据:wrk压测对比
压测条件:wrk -t4 -c200 -d10s http://localhost:8000/api/tickets,模拟5个工单ID。
| 指标 | 改造前(同步) | 改造后(asyncio) | 提升 |
|---|---|---|---|
| 平均耗时 | 2.8s | 420ms | 6.7倍 |
| P99耗时 | 3.5s | 780ms | 4.5倍 |
| QPS | 320 | 2100 | 6.6倍 |
| CPU利用率 | 15% | 68% | 4.5倍 |
| 外部API并发 | 1(串行) | 5(Semaphore限流) | 稳定 |
数据说明:QPS提升主要来自IO等待时间被并发利用。原来一个worker等600ms外部响应时CPU空转,现在5个协程同时等,单位时间能处理更多请求。另外,超时重试机制让接口在外部API抖动时依然稳定,P99从3.5s降到780ms是关键收益。
七、总结
asyncio不是银弹,但对付IO密集型接口非常有效。这次改造只动了核心逻辑,数据库和外部API调用方式没变,但性能提升6倍以上。几点经验:
- 同步框架也能用asyncio:
asyncio.run()+ 协程函数,不一定要换FastAPI - Semaphore必须全局:并发控制是异步编程最容易忽略的点
- 数据库操作用
to_thread丢线程池:别在协程里直接跑同步DB驱动 - 连接池和事件循环生命周期要一致:否则遇到
Event loop is closed别怪我没提醒
如果后续工单数量继续增长,可以考虑把ClientSession做成进程级复用(用uvicorn的lifespan),或者引入Redis缓存外部API结果。但就当前业务规模,asyncio这套方案已经足够。