一、问题背景:一个「看起来没问题」的查询接口

那是周三下午,运营反馈「订单导出功能卡死」。我查了监控,发现/api/v1/orders/batch这个接口在200并发下:
- 平均RT:305ms(其中数据库查询占270ms)
- QPS:327
- PostgreSQL连接数:飙到150+,接近max_connections=200

当时的代码逻辑很「标准」——Flask同步视图,内部用psycopg2连接PostgreSQL查询订单主表,然后循环调用第三方ERP接口获取物流状态。问题就出在这个循环上:每个订单要串行请求一次ERP,平均耗时80ms。如果有50个订单,光这层就是4秒。

同步模型下,线程池炸了,连接池也炸了。我当时的内心OS:这代码没写错,就是架构错了。

二、环境与版本

改造前先交代环境,方便读者复现:

Python: 3.10.11 (需要3.8+才能完整使用asyncio.run)
Web框架: Flask 2.2.5 (同步,通过gunicorn运行)
HTTP客户端: requests 2.28.2 (同步)
数据库驱动: psycopg2-binary 2.9.6 (同步)
压测工具: wrk 4.2.0
服务器: 4C8G 阿里云ECS,PostgreSQL 14 (max_connections=200)

gunicorn配置: gunicorn -w 8 -k gthread --threads 16 app:app

改造后引入:
- aiohttp 3.8.4 (或httpx 0.24.1,本文用httpx)
- asyncpg 0.27.0
- uvloop 0.17.0 (事件循环替换)

三、方案设计:异步化不是「换框架」,是「换I/O模型」

很多人一谈异步就想到换FastAPI。但你完全可以保留Flask,用asyncio.run()在视图内部跑协程。思路分三层:

  1. 数据库层:psycopg2 → asyncpg。连接池从SQLAlchemy的ThreadPool切换为asyncpg内置Pool,限制max_size=10。
  2. 外部API调用层:requests → httpx.AsyncClient。用asyncio.gather并发请求ERP,而不是for循环串行。
  3. 事件循环:安装uvloop,替换默认的asyncio事件循环。实测uvloop在纯网络I/O场景下能减少约15%开销。

核心设计图(文字版):

同步版本: 请求进来 -> 线程池分配线程 -> psycopg2阻塞查询 -> for循环requests阻塞调用ERP -> 返回
异步版本: 请求进来 -> 主线程asyncio.run -> asyncpg异步查询 -> gather并发调用ERP -> 返回

关键点:异步化要解决的是「等待时间」问题,而不是「计算时间」问题。如果你的逻辑是CPU密集,异步没用。

四、核心实现:before/after代码对比

4.1 同步版本(改造前)

# app_sync.py  — 改造前代码
from flask import Flask, jsonify, request
import psycopg2
import requests
from psycopg2.pool import ThreadedConnectionPool

app = Flask(__name__)
pool = ThreadedConnectionPool(5, 20, host='localhost', dbname='orders')

@app.route('/api/v1/orders/batch', methods=['GET'])
def batch_query():
    order_ids = request.args.get('ids', '').split(',')
    if len(order_ids) > 50:
        return jsonify({'error': 'max 50 ids'}), 400

    conn = pool.getconn()
    try:
        # 1. 查询主表
        cur = conn.cursor()
        cur.execute("SELECT id, amount FROM orders WHERE id = ANY(%s)", (order_ids,))
        rows = cur.fetchall()
    finally:
        pool.putconn(conn)

    # 2. 串行调用ERP — 这里是最痛点
    erp_results = []
    for oid in order_ids:
        resp = requests.get(f'http://erp.example.com/status/{oid}', timeout=3)
        erp_results.append(resp.json())

    return jsonify({'orders': rows, 'erp': erp_results})

if __name__ == '__main__':
    app.run(threaded=True)

这个版本的性能瓶颈一目了然:
- 数据库查询阻塞线程(线程池16个,但200并发全排队)
- ERP调用是「串行」的,50个订单 = 50次RTT

4.2 异步重构版本(改造后)

# app_async.py — 改造后核心代码
import asyncio
import uvloop
from flask import Flask, jsonify, request
import httpx
import asyncpg

app = Flask(__name__)
asyncio.set_event_loop_policy(uvloop.EventLoopPolicy())

# 全局初始化连接池 (在app启动时调用)
async def init_db_pool():
    return await asyncpg.create_pool(
        host='localhost', database='orders',
        min_size=5, max_size=10,  # 关键:限制连接数,避免打爆数据库
        command_timeout=5
    )

async def fetch_erp_status(client: httpx.AsyncClient, oid: str, sem: asyncio.Semaphore):
    """并发获取ERP状态,使用信号量限制并发数"""
    async with sem:  # 限制ERP并发不超过20
        try:
            resp = await client.get(f'http://erp.example.com/status/{oid}', timeout=3)
            return {'oid': oid, 'status': resp.json().get('status')}
        except httpx.TimeoutException:
            return {'oid': oid, 'status': 'timeout'}

async def handle_batch(order_ids):
    # 数据库异步查询
    pool = await init_db_pool()
    async with pool.acquire() as conn:
        rows = await conn.fetch(
            "SELECT id, amount FROM orders WHERE id = ANY($1)", order_ids
        )

    # 并发ERP调用 — 核心优化点
    async with httpx.AsyncClient(timeout=3) as client:
        sem = asyncio.Semaphore(20)  # 防止打爆ERP
        tasks = [fetch_erp_status(client, oid, sem) for oid in order_ids]
        erp_results = await asyncio.gather(*tasks)

    return {'orders': [dict(r) for r in rows], 'erp': erp_results}

@app.route('/api/v1/orders/batch', methods=['GET'])
def batch_query():
    order_ids = request.args.get('ids', '').split(',')
    if len(order_ids) > 50:
        return jsonify({'error': 'max 50 ids'}), 400

    # Flask同步视图内运行事件循环
    result = asyncio.run(handle_batch(order_ids))
    return jsonify(result)

# 重要:需要关闭已有的asyncpg pool,避免连接泄漏
# 生产用uvicorn或hypercorn跑ASGI,这里为演示保留Flask

注意几个细节:
- asyncio.run() 每次创建新事件循环,不能用于高并发生产(每次都要重建连接池)。实际部署请用FastAPI/Quart或AIOHTTP,或者用gunicorn的uvicorn.workers.UvicornWorker。我这里为了对比才用这种「混合模式」。
- asyncio.Semaphore(20)必须有,否则100个订单会同时打ERP,直接给人家打挂。
- asyncpg的fetch返回Record对象,需要dict(r)转JSON。

五、踩坑与优化:三个让我失眠的bug

5.1 坑1:连接池泄漏导致「偶尔卡死」

第一次写完跑测试,前5000请求正常,然后全部超时。排查半天发现是asyncio.run()每次都会新建pool,而init_db_pool()在函数内被反复调用,旧连接没被关闭。最终方案:把pool声明为global,启动时初始化一次,请求时复用。

# 正确做法 —— 在app.py顶层:
_pool = None

async def get_pool():
    global _pool
    if _pool is None:
        _pool = await asyncpg.create_pool(min_size=5, max_size=10)
    return _pool

5.2 坑2:uvloop和asyncpg的SSL冲突

我们数据库开了SSL,uvloop下asyncpg报错ssl.SSLError。查了issue发现是asyncpg 0.27.0的已知bug,升级到0.28.0解决(或者禁用uvloop)。最终选择升级asyncpg,因为uvloop带来的性能收益太香。

5.3 坑3:wrk压测时「连接数看起来没降」

改完压测,看PostgreSQL pg_stat_activity傻眼了——连接数还是150。后来发现是压测客户端自己开了太多连接。调整wrk参数-c 200是指并发连接数,而gunicorn线程模式也有影响。最终用gunicorn -w 4 -k uvicorn.workers.UvicornWorker跑ASGI模式才达到理想效果。

六、效果数据:压测对比(wrk -t8 -c200 -d30s)

指标 同步版本 asyncio版本 提升
QPS 327 1892 4.78倍
平均RT 305ms 52ms 5.87倍
P99 RT 890ms 143ms 6.2倍
PostgreSQL连接数 150+ 10 -93%
错误率 4.5% 0%

关键数据说明:
- PostgreSQL连接数从150降到10,是因为asyncpg的max_size=10强制复用连接。数据库CPU从85%降到32%。
- ERP调用并发从「串行50次」变成「并发20路」,单接口耗时从4秒(5080ms)降到约400ms(50/2080ms),瓶颈变为数据库查询。

七、总结与建议

异步化适合「I/O密集 + 并发高」的场景。如果你遇到类似问题,按这个顺序排查:

  1. 先确认瓶颈在I/O(数据库查询慢、外部API调用慢),而非CPU计算
  2. 从最耗时的依赖开始替换——我的案例里是ERP串行调用,而不是数据库
  3. 控制并发度:Semaphore和连接池max_size是防身符
  4. 别用asyncio.run处理真实请求——它每次新建事件循环,代价太高。生产环境用ASGI框架或worker

最后说点实在的:如果你在维护老Flask项目,没必要全量迁移FastAPI。Flask视图内嵌asyncio.run是过渡方案,但长期来看,建议将高并发接口单独用FastAPI写,或者整体切换到Quart(Flask的异步版本)。

数据不会骗人——当你的QPS从300提升到1800,那种爽感是写同步代码永远体会不到的。