1. 问题背景:一个慢到被运维找上门的API

我们有个内部服务叫/api/v1/orders/batch,前端一次要拉取50个订单的详情。最初实现就是标准的Flask同步视图,内部串行调用下游订单服务的五个接口(详情、商品、物流、优惠、用户)。每次请求平均耗时1.8秒,压测时QPS卡在180左右。

运维同事丢来一张监控截图:CPU才用了30%,但线程数飙到400+,大量请求堆积在等待下游IO。这就是典型的IO密集型瓶颈——线程都在等网络响应,GIL锁导致CPU无法充分利用,而线程切换开销又大。

2. 环境与版本:先交代清楚实验条件

  • Python 3.10.12(asyncio在3.10才支持asyncio.timeout,3.8的旧项目建议用async_timeout包)
  • Flask 2.3.3(WSGI同步框架,需要配合asgiref或直接换aiohttp)
  • aiohttp 3.9.1(异步HTTP客户端)
  • gunicorn 21.2.0(生产部署,worker模式后面细说)
  • 压测工具:wrk 4.2.0,参数-t8 -c200 -d30s

测试环境:4核8G Docker容器,下游服务用fastapi模拟,每个接口延迟200ms。

3. 方案设计:为什么选asyncio而不是ThreadPoolExecutor

先看原始代码(简化版):

# before_sync.py
import requests
from flask import Flask, jsonify

app = Flask(__name__)

def fetch_order_detail(order_id):
    # 模拟串行调用5个下游接口,每个200ms
    result = {}
    result['order'] = requests.get(f'http://svc/order/{order_id}').json()
    result['product'] = requests.get(f'http://svc/product/{order_id}').json()
    result['logistics'] = requests.get(f'http://svc/logistics/{order_id}').json()
    result['promo'] = requests.get(f'http://svc/promo/{order_id}').json()
    result['user'] = requests.get(f'http://svc/user/{order_id}').json()
    return result

@app.route('/api/v1/orders/batch', methods=['POST'])
def batch():
    order_ids = request.json['order_ids']  # 50个ID
    result = [fetch_order_detail(oid) for oid in order_ids]
    return jsonify(result)

这里有两个层次的串行:订单之间串行(50个订单循环),每个订单内部串行(5个下游接口)。

有人会问:用ThreadPoolExecutorrequests不行吗?我试过,能提到500QPS,但有两个问题:1)线程池大小难调,调大了内存暴涨(每个线程默认8MB栈);2)下游服务一旦变慢,线程池排队,整体延迟迅速恶化。

asyncio的优势在于单线程内协程切换,不需要操作系统线程上下文切换,一个进程能支撑数万连接。更重要的是,aiohttp原生复用连接池,避免每次请求都重新握手。

4. 核心实现:全异步化改造

改造分两步:第一步,把requests换成aiohttp,实现fetch_order_detail的异步版本;第二步,把Flask换成aiohttp.web(因为Flask不支持异步视图,虽然3.0支持了async def视图,但底层WSGI还是同步的)。

# after_async.py
import asyncio
import aiohttp
from aiohttp import web

# 全局连接池:限制最大连接数,避免打爆下游
connector = aiohttp.TCPConnector(limit=100, ttl_dns_cache=300)
session = aiohttp.ClientSession(connector=connector)

async def fetch_one(session, url, timeout=0.5):
    """单个下游请求,带超时控制"""
    try:
        async with session.get(url, timeout=aiohttp.ClientTimeout(total=timeout)) as resp:
            return await resp.json()
    except asyncio.TimeoutError:
        return {'error': 'timeout'}

async def fetch_order_detail(session, order_id):
    # 并发请求5个下游,不相互等待
    urls = {
        'order': f'http://svc/order/{order_id}',
        'product': f'http://svc/product/{order_id}',
        'logistics': f'http://svc/logistics/{order_id}',
        'promo': f'http://svc/promo/{order_id}',
        'user': f'http://svc/user/{order_id}',
    }
    tasks = [fetch_one(session, url) for url in urls.values()]
    results = await asyncio.gather(*tasks, return_exceptions=True)
    return dict(zip(urls.keys(), results))

async def batch_handler(request):
    data = await request.json()
    order_ids = data['order_ids']  # 50个
    # 控制并发度:同时最多处理20个订单,避免下游过载
    sem = asyncio.Semaphore(20)

    async def bounded_fetch(oid):
        async with sem:
            return await fetch_order_detail(session, oid)

    # 全部订单并发执行
    results = await asyncio.gather(*(bounded_fetch(oid) for oid in order_ids))
    return web.json_response(results)

app = web.Application()
app.router.add_post('/api/v1/orders/batch', batch_handler)

if __name__ == '__main__':
    web.run_app(app, host='0.0.0.0', port=8000, access_log=None)

关键改动点:
1. asyncio.gather并发所有订单:50个订单同时启动,每个订单内部5个请求也并发,理论最大并发=505=250个连接。
2.
Semaphore(20)限流:防止一次性250个请求把下游打满,实测下游连接数超过100就会开始报错。
3.
aiohttp.ClientTimeout(total=0.5)*:每个下游请求超时500ms,整体最坏情况约1秒(因为5个并发,取最慢的),比原来1.8秒好。

部署命令也变了,不能用flask run

# 原先:gunicorn -w 4 before_sync:app
# 现在:aiohttp应用直接用python起,或配合uvloop
python after_async.py
# 生产建议:python -m aiohttp.web -H 0.0.0.0 -P 8000 after_async:app

5. 踩坑与优化:三个坑差点让我放弃

坑1:Event Loop被同步代码阻塞
最初我在fetch_order_detail里不小心调用了time.sleep(0.1)模拟延迟,结果整个事件循环卡住,QPS直接归零。排查半天发现是同事在中间件里用了同步的redis.get记住:asyncio代码里绝不能有同步阻塞调用,一切IO都要走协程(await)或者丢给线程池(loop.run_in_executor)。

坑2:连接池耗尽导致连环超时
aiohttp.ClientSession默认连接池大小是100。当并发订单数从10调到50时,下游连接数瞬间到250,aiohttp内部排队等待可用连接,导致超时。解决方式就是上面代码里的connector参数,限制limit=100,同时用Semaphore控制应用层并发。

坑3:Python 3.10的asyncio.timeout vs 老项目
3.10之前用asyncio.wait_for,但它在超时后不会取消任务,导致任务继续运行浪费资源。3.11的asyncio.timeout是上下文管理器,推荐用这个。我因为生产环境是3.10.12,用了asyncio.timeout,效果不错。

6. 效果数据:对比压测结果

用wrk压测30秒,200并发,结果如下:

指标 同步版 (Flask + requests) 异步版 (aiohttp)
QPS 182 1547
平均延迟 1080ms 640ms
P99延迟 2300ms 280ms
最大延迟 4100ms 950ms
内存占用 420MB (4 worker × 105MB) 230MB (1进程)
CPU使用 30% (瓶颈在线程切换) 85% (有效利用)

QPS提升约8.5倍,P99从2.3秒降到280ms(8.2倍提升)。最意外的是内存从420MB降到230MB——因为去掉了4个gunicorn worker,单进程事件循环处理所有请求。

补充一个真实场景的观察:同步版在下游服务抖动时(延迟从200ms涨到800ms),QPS会跌到50以下;异步版由于超时控制(500ms),最坏情况P99仍在800ms内,不会无限堆积。

7. 总结与后续优化方向

这次重构让我确信:在Python里做高并发IO密集型服务,asyncio是首选方案。但要注意几点:
- 不是所有项目都适合改造,如果业务逻辑里同步库占大头(比如重度依赖requestspsycopg2同步驱动),改造成本高,收益有限。
- 异步化后要重点监控连接池状态事件循环延迟(可以用loop.slow_callback_duration设置警告阈值)。
- 后续可以做的优化:用uvloop替代默认事件循环(再提升15%左右)、HTTP/2多路复用(减少连接数)、缓存下游响应(async_lru装饰器)。

最后提一句,如果团队刚接触asyncio,建议先从一个非核心API试点,踩过坑再推广。我这套代码已经在生产跑了两个月,稳定在1200-1500 QPS,运维终于不用半夜找我了。