一、问题背景:当你的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 defawait的用法,以及如何在一个同步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.runasyncio.run每次创建事件循环的开销约1ms,且不易管理全局资源。如果必须用Flask,考虑gunicorngevent worker而不是asyncio。
  • 性能数字依赖上游:我们的上游服务P95是1.2秒,如果你的上游是2秒,异步后P95约2秒,不会更慢——但吞吐量会提升到相近水平。

最后,异步编程的核心是让等待不阻塞。在IO密集场景下,它比多线程更轻量、比多进程更省内存。但请记住:没有银弹,只有适合场景的选择。如果你的API还在串行等待,不妨试试asyncio,但记得控制好并发度和超时。