一、问题背景:同步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]的异步视图(配合asgirefasync_to_sync),二是把异步逻辑独立成协程,通过asyncio.run()在同步视图里调用。我选了第二种,因为改动最小,且兼容现有中间件。

三、方案设计:asyncio.gather并发替代顺序调用

核心思路:把get_ticket_detail从同步函数改为async def协程函数,内部用aiohttp发起异步HTTP请求。然后在视图函数中:

  1. asyncio.run()启动事件循环
  2. asyncio.gather()并发执行所有工单的协程
  3. 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_threadexecutor参数,限制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倍以上。几点经验:

  1. 同步框架也能用asyncioasyncio.run() + 协程函数,不一定要换FastAPI
  2. Semaphore必须全局:并发控制是异步编程最容易忽略的点
  3. 数据库操作用to_thread丢线程池:别在协程里直接跑同步DB驱动
  4. 连接池和事件循环生命周期要一致:否则遇到Event loop is closed别怪我没提醒

如果后续工单数量继续增长,可以考虑把ClientSession做成进程级复用(用uvicornlifespan),或者引入Redis缓存外部API结果。但就当前业务规模,asyncio这套方案已经足够。