1. 问题背景:一个把线程池打爆的聚合查询API

去年底接手一个内部数据中台服务,其中/api/v1/aggregate接口负责聚合5个下游系统的数据(用户信息、订单、库存、物流、风控)。原实现是Flask同步视图,内部用requests.get串行调用下游:

# 重构前:串行阻塞IO
@app.route('/api/v1/aggregate')
def aggregate():
    user = requests.get('http://user-svc/info', params={'uid': uid}).json()
    order = requests.get('http://order-svc/list', params={'uid': uid}).json()
    # ... 库存、物流、风控
    return jsonify(merge(user, order, stock, logistics, risk))

压测结果惨不忍睹:QPS 200时,P99延迟2.3s,线程池(默认40)全部占满,大量ConnectionResetError。原因很简单:每个请求占用一个线程,线程在等待下游响应时完全阻塞,5次串行IO耗时约1.8s,200并发需要1000个线程才能勉强支撑——Flask开发服务器直接崩溃。

2. 环境与版本:Python 3.10 + 全异步栈

  • Python 3.10.8(原生asyncio,未用第三方事件循环)
  • aiohttp 3.8.4(替代requests)
  • asyncpg 0.27.0(替代psycopg2,后续优化用)
  • FastAPI 0.95.0(替代Flask,但核心改造是asyncio)
  • 压测工具:wrk 4.2.0(单机模拟1000并发)

注意:Python 3.10的asyncio已经默认使用ProactorEventLoop(Windows)或SelectorEventLoop(Linux),无需额外配置。但千万别用uvloop——我试过,在asyncpg下会偶发RuntimeError: Event loop is closed,排查了两天,最后放弃。

3. 方案设计:异步化三层改造

不满足于只把requests换成aiohttp,我的改造分三层:

  1. Web框架层:Flask → FastAPI(原生支持async def视图)
  2. HTTP客户端层:requests → aiohttp.ClientSession(连接池复用)
  3. 数据库层:psycopg2同步驱动 → asyncpg(配合连接池)

核心设计思想:把每个下游调用封装成asyncio.Task,用asyncio.gather并发执行。同时用asyncio.Semaphore限制并发数,防止下游被瞬时流量打垮。

4. 核心实现:asyncio.gather + 信号量限流

# 重构后:异步并发调用,信号量限流
import asyncio
import aiohttp
from fastapi import FastAPI
import asyncpg

app = FastAPI()
# 全局连接池(避免每次请求重建)
http_session = None
pg_pool = None
# 限制并发数,保护下游
semaphore = asyncio.Semaphore(200)

async def fetch_json(client, url, params):
    async with semaphore:
        async with client.get(url, params=params, timeout=5) as resp:
            return await resp.json()

async def aggregate_core(uid: int):
    # 并发发起5个下游请求
    tasks = [
        fetch_json(http_session, 'http://user-svc/info', {'uid': uid}),
        fetch_json(http_session, 'http://order-svc/list', {'uid': uid}),
        fetch_json(http_session, 'http://stock-svc/check', {'uid': uid}),
        fetch_json(http_session, 'http://logistics-svc/track', {'uid': uid}),
        fetch_json(http_session, 'http://risk-svc/score', {'uid': uid}),
    ]
    results = await asyncio.gather(*tasks, return_exceptions=True)
    # 处理个别失败,不影响整体
    return merge_with_fallback(results)

@app.get('/api/v1/aggregate')
async def aggregate(uid: int):
    return await aggregate_core(uid)

@app.on_event('startup')
async def startup():
    global http_session, pg_pool
    http_session = aiohttp.ClientSession(
        connector=aiohttp.TCPConnector(limit=500, ttl_dns_cache=300)
    )
    pg_pool = await asyncpg.create_pool(
        'postgresql://user:pass@localhost/db', min_size=10, max_size=50
    )

关键点:
- aiohttp.TCPConnector(limit=500):连接池上限500,避免文件描述符耗尽
- timeout=5:下游5秒无响应直接放弃,防止无限等待
- return_exceptions=True:单个下游挂了不影响其他四个

5. 踩坑与优化:三个让我抓狂的问题

坑1:asyncio.Semaphore滥用导致死锁
第一次我把Semaphore(200)放在aggregate_core外面,结果高并发下信号量被占满,新请求永远等待——信号量必须放在fetch_json内部,而不是外层循环。

坑2:asyncpg连接池被耗尽
压测时数据库连接池(max_size=50)被占满,报TimeoutError: Pool is full。解决:把max_size提到100,并加上max_inactive_connection_lifetime=60,让空闲连接及时回收。

坑3:gather结果顺序错乱
asyncio.gather返回结果顺序和task顺序一致,但如果你用asyncio.wait就会乱。我一度用wait然后自己排序,后来发现gather自带顺序保证,改用gather后代码少了10行。

优化后最终版还做了响应缓存(functools.lru_cache对uid聚合结果缓存5秒),P99进一步降到98ms。

6. 效果数据:从120到2050的飞跃

用wrk压测30秒,1000并发,对比重构前后:

指标 重构前(Flask+requests) 重构后(FastAPI+aiohttp+asyncpg)
QPS 120 req/s 2050 req/s(提升17倍
P99 2300ms 135ms(降低94%
线程/协程数 40线程全阻塞 500协程(内存占用仅增加80MB)
错误率 12.5%(连接超时) 0.3%(仅下游5xx)

数据库查询也从原来的串行3次(每次50ms)变成并发1次,延迟从150ms降至52ms。额外收益:部署实例从5台降到2台,每台CPU使用率稳定在45%左右。

7. 总结:异步不是银弹,但这里值得用

我的经验是:当你的瓶颈是IO等待(HTTP调用、数据库查询、文件读写),且并发量超过100时,asyncio是性价比最高的方案。但注意:
- 千万别把CPU密集型任务扔进async函数,会卡死事件循环
- 数据库驱动必须用异步版(asyncpg/aiomysql),否则白搭
- 连接池尺寸要按下游能力压测,不是越大越好

最后建议:如果你还在用Flask+requests,且接口里有两个以上串行IO调用,立刻动手改成FastAPI+asyncio。别犹豫,性能提升是实打实的,代码反而更简洁——去掉线程池的锅,你只需要关心事件循环里谁先完成。