一、问题背景:一个被IO拖死的查询接口

我们有个内部服务,叫 user-profile-api,对外暴露一个接口:

GET /profile?uid=123456

逻辑很朴素:拿到uid,去三个下游拉数据——Redis查基础画像、MySQL查订单统计、调用推荐服务HTTP接口拿标签,然后聚合返回。

上线初期没问题,日调用量涨到80万之后,监控开始报警:

  • P99延迟 1.8s,P95 也有 900ms
  • 单实例QPS卡在 120 上不去
  • CPU利用率只有 15%,但worker进程全在等IO

典型的IO密集型服务被同步阻塞拖死。当时用的是 Flask 2.3 + Gunicorn(sync worker,8个进程),每个请求串行做三次下游调用,单次RT加起来 200ms 左右,8个进程理论上限也就 40 QPS/进程 × 8 = 320,实际因为GIL和连接池争用只跑到120。

看一眼原始代码就明白问题在哪:

# before: app.py (Flask 2.3.3 + requests 2.31.0)
import requests
import redis
import pymysql
from flask import Flask, jsonify, request

app = Flask(__name__)
rds = redis.Redis(host='redis.internal', port=6379, decode_responses=True)
db = pymysql.connect(host='mysql.internal', user='app', password='***', database='order')

@app.route('/profile')
def profile():
    uid = request.args.get('uid')
    # 1. Redis 基础画像
    base = rds.hgetall(f'profile:{uid}')
    # 2. MySQL 订单统计
    with db.cursor() as cur:
        cur.execute('SELECT COUNT(*), SUM(amount) FROM orders WHERE uid=%s', (uid,))
        order_stat = cur.fetchone()
    # 3. 推荐服务 HTTP
    tags = requests.get(f'http://reco.internal/tags?uid={uid}', timeout=0.5).json()
    return jsonify({'base': base, 'orders': order_stat, 'tags': tags})

三个调用完全独立,却写成串行。就算改成多线程,GIL + 连接池也会让你难受。这就是asyncio的用武之地。

二、环境与版本

  • Python 3.11.6(3.11的asyncio性能比3.8好不少,尤其是Task调度)
  • aiohttp 3.9.1(HTTP客户端)
  • aiomysql 0.2.0(MySQL异步驱动)
  • redis 5.0.1(自带asyncio支持,redis.asyncio
  • fastapi 0.108.0 + uvicorn 0.25.0(替换Flask作为Web层)
  • 压测工具:wrk 4.2.0,wrk -t8 -c200 -d30s

服务器配置:4C8G,CentOS 7,下游服务RT稳定在 60~80ms。

三、方案设计:并发+限流+连接复用

核心思路三条:

  1. 三个下游调用并发执行,用 asyncio.gather 聚合,理论RT从 200ms 降到 max(各调用) ≈ 80ms。
  2. 信号量限流,防止并发打爆下游。每个下游一个 asyncio.Semaphore,Redis/MySQL设50,推荐服务设30。
  3. 连接池全局复用,Redis和MySQL的连接在 lifespan 里初始化,避免每请求建连。aiohttp用 TCPConnector(limit=100, limit_per_host=30)

并发模型大致是这样:

请求 → FastAPI (uvicorn, 1 event loop)
         ├── Task A: redis.hgetall
         ├── Task B: aiomysql query
         └── Task C: aiohttp GET reco
         → asyncio.gather 聚合 → 返回

四、核心实现

4.1 异步依赖初始化

# after: main.py
import asyncio
import aiohttp
import aiomysql
import redis.asyncio as aioredis
from fastapi import FastAPI, Query
from contextlib import asynccontextmanager

# 全局资源
state = {}

@asynccontextmanager
async def lifespan(app: FastAPI):
    # Redis 连接池
    state['redis'] = aioredis.Redis(
        host='redis.internal', port=6379, decode_responses=True,
        max_connections=50, socket_timeout=0.3
    )
    # MySQL 连接池
    state['mysql'] = await aiomysql.create_pool(
        host='mysql.internal', user='app', password='***',
        db='order', minsize=10, maxsize=50, pool_recycle=1800
    )
    # HTTP 连接池
    connector = aiohttp.TCPConnector(limit=100, limit_per_host=30, ttl_dns_cache=300)
    state['http'] = aiohttp.ClientSession(connector=connector, timeout=aiohttp.ClientTimeout(total=0.5))
    # 信号量
    state['sem_redis'] = asyncio.Semaphore(50)
    state['sem_mysql'] = asyncio.Semaphore(50)
    state['sem_reco']  = asyncio.Semaphore(30)
    yield
    # 清理
    await state['http'].close()
    state['mysql'].close()
    await state['mysql'].wait_closed()
    await state['redis'].close()

app = FastAPI(lifespan=lifespan)

4.2 三个下游的异步封装

async def fetch_redis(uid: str):
    async with state['sem_redis']:
        return await state['redis'].hgetall(f'profile:{uid}')

async def fetch_mysql(uid: str):
    async with state['sem_mysql']:
        async with state['mysql'].acquire() as conn:
            async with conn.cursor() as cur:
                await cur.execute(
                    'SELECT COUNT(*), SUM(amount) FROM orders WHERE uid=%s', (uid,)
                )
                return await cur.fetchone()

async def fetch_reco(uid: str):
    async with state['sem_reco']:
        async with state['http'].get(f'http://reco.internal/tags?uid={uid}') as resp:
            return await resp.json()

4.3 聚合接口

@app.get('/profile')
async def profile(uid: str = Query(...)):
    results = await asyncio.gather(
        fetch_redis(uid),
        fetch_mysql(uid),
        fetch_reco(uid),
        return_exceptions=True
    )
    base, order_stat, tags = results
    # 容错:任一失败返回默认值,不阻塞整体
    if isinstance(base, Exception): base = {}
    if isinstance(order_stat, Exception): order_stat = (0, 0)
    if isinstance(tags, Exception): tags = []
    return {'base': base, 'orders': order_stat, 'tags': tags}

启动命令:

uvicorn main:app --host 0.0.0.0 --port 8000 --workers 1 --loop uvloop --http httptools

注意这里 workers=1。因为asyncio本身就是单线程事件循环,一个worker就能吃满一个核的IO并发能力。我实测开4个worker反而因为连接池翻倍、下游压力大导致P99升高。

五、踩坑与优化

坑1:在协程里用了同步的requests

第一版我图省事,fetch_reco 里直接用了 requests.get。结果QPS只从120涨到180,一堆请求还是串行。原因:requests 是阻塞调用,会把整个event loop卡住,asyncio.gather 的并发完全失效。换成 aiohttp 后QPS直接跳到900+。

教训:asyncio代码里不能出现任何同步IO。

坑2:MySQL连接池开太小

一开始 minsize=5, maxsize=10,压测时200并发进来,大量请求卡在 pool.acquire() 上等连接,P99反而变差。调到 minsize=10, maxsize=50 后,等待时间从平均80ms降到5ms以内。

坑3:gather 里没加 return_exceptions

某次推荐服务抖动,一个请求异常导致整个 /profile 抛500。改成 return_exceptions=True 后,单点失败降级为默认值,接口可用性从99.2% 提升到 99.97%。

优化:uvloop + httptools

默认asyncio用的是 selectors 事件循环,换成uvloop后:

  • QPS:1100 → 1280(+16%)
  • P99:210ms → 175ms

安装:pip install uvloop httptools,uvicorn启动加 --loop uvloop --http httptools 即可。

六、效果数据

压测命令统一为 wrk -t8 -c200 -d30s http://api.internal/profile?uid=123456,下游RT稳定在70ms左右。

指标 Before (Flask+Gunicorn 8w) After (FastAPI+uvicorn+uvloop) 提升
QPS 120 1280 10.7x
P50 420ms 95ms 4.4x
P95 900ms 160ms 5.6x
P99 1800ms 175ms 10.3x
CPU利用率 15% 68% -
单机内存 780MB 420MB -46%

线上灰度一周,日均80万调用下:

  • 报错率从 0.8% 降到 0.03%
  • 部署实例从 8台 缩到 2台,成本省了75%
  • 平均延迟从 380ms 降到 110ms

七、总结

这次重构最大的感受是:IO密集型服务的瓶颈从来不是CPU,而是你等IO的方式。asyncio不是什么银弹,但用对了场景,提升是数量级的。

几点经验直接给结论:

  1. 只要接口里有2个以上独立的下游调用,asyncio.gather 就是最直接的优化手段。
  2. 别在协程里混同步库,requestspymysqlredis(同步版)统统换成异步版。
  3. 连接池和信号量的参数必须压测调,不是拍脑袋定的。
  4. uvicorn worker数别开太多,1~2个通常够,多了反而互相抢连接。
  5. return_exceptions=True 是生产环境的标配,别让一个下游抖动拖垮整个接口。

如果你的服务也是这种"调多个下游再聚合"的形态,不妨先profiling看看IO等待占比。超过60%的话,asyncio重构基本稳赚。