一、问题背景:一个“看起来没问题”的同步接口

公司内部有一个订单查询接口,逻辑很简单:根据用户ID查订单列表,然后补全商品名称和店铺信息。代码写得很规矩,Flask 2.2 + SQLAlchemy 1.4 + PyMySQL,没毛病。但线上监控显示,这个接口在高峰期P99延迟超过2.5s,QPS一过100就开始超时报警。

看代码挺干净的:

# before.py
@app.route('/api/orders', methods=['GET'])
def get_orders():
    user_id = request.args.get('user_id')
    # 1. 查用户订单
    orders = db.session.execute(
        text("SELECT id, product_id, store_id, amount FROM orders WHERE user_id = :uid"),
        {"uid": user_id}
    ).fetchall()
    # 2. 循环查商品和店铺(这里最要命)
    result = []
    for o in orders:
        product = db.session.execute(
            text("SELECT name FROM products WHERE id = :pid"),
            {"pid": o.product_id}
        ).fetchone()
        store = db.session.execute(
            text("SELECT name FROM stores WHERE id = :sid"),
            {"sid": o.store_id}
        ).fetchone()
        result.append({
            "order_id": o.id,
            "product_name": product[0] if product else "N/A",
            "store_name": store[0] if store else "N/A",
            "amount": str(o.amount)
        })
    return jsonify(result)

问题一目了然:N+1查询。一个用户平均有8个订单,每个订单要额外查两次数据库,总共17次同步DB往返。每次DB查询走PyMySQL,默认单连接,串行执行,一次往返平均5-8ms(本地网络),最坏情况光DB就吃掉100ms+。线程池默认8线程,QPS一上来就排队。

二、选型:为什么不换FastAPI,而是用asyncio改造

同事提议直接换FastAPI + async SQLAlchemy。但我考虑了两点:

  1. 现有代码3000多行,全部迁移成本高,风险大。
  2. 我们真正的问题不是框架慢,而是同步阻塞IO。Flask的视图函数可以改成协程,用asyncio跑异步DB驱动,效果一样。

最终方案:保留Flask框架,用flask[async] + asyncio.run()包装视图 + aiomysql替代PyMySQL + uvloop替换默认事件循环

环境版本:

  • Python 3.10.6
  • Flask 2.2.5 (支持async视图)
  • aiomysql 0.2.0
  • uvloop 0.17.0
  • MySQL 8.0.28(连接池大小20)

三、核心实现:改造后的异步接口

3.1 连接池初始化(避免每次请求新建连接)

# async_pool.py
import asyncio
import aiomysql

POOL = None

async def init_pool():
    global POOL
    POOL = await aiomysql.create_pool(
        host='127.0.0.1',
        port=3306,
        user='app_user',
        password='secret',
        db='shop',
        minsize=5,
        maxsize=20,
        autocommit=True,
        echo=False
    )

def get_pool():
    if POOL is None:
        raise RuntimeError("Pool not initialized")
    return POOL

3.2 异步视图函数

注意:aiomysql的游标默认是元组,用DictCursor方便取列名。核心改动是把循环里的同步DB调用改为await,并且asyncio.gather并发处理订单的商品和店铺查询(这比顺序await快得多)。

# after.py
import asyncio
import aiomysql
from flask import jsonify, request
from async_pool import get_pool

async def fetch_one(sql, args):
    pool = get_pool()
    async with pool.acquire() as conn:
        async with conn.cursor(aiomysql.DictCursor) as cur:
            await cur.execute(sql, args)
            row = await cur.fetchone()
            return row

async def fetch_order_details(order):
    """并发查询商品和店铺信息"""
    product_sql = "SELECT name FROM products WHERE id = %s"
    store_sql = "SELECT name FROM stores WHERE id = %s"
    # 关键:gather并发执行两个查询
    product, store = await asyncio.gather(
        fetch_one(product_sql, (order['product_id'],)),
        fetch_one(store_sql, (order['store_id'],))
    )
    return {
        'order_id': order['id'],
        'product_name': product['name'] if product else 'N/A',
        'store_name': store['name'] if store else 'N/A',
        'amount': str(order['amount'])
    }

@app.route('/api/orders', methods=['GET'])
async def get_orders_async():
    user_id = request.args.get('user_id')
    pool = get_pool()
    async with pool.acquire() as conn:
        async with conn.cursor(aiomysql.DictCursor) as cur:
            await cur.execute(
                "SELECT id, product_id, store_id, amount FROM orders WHERE user_id = %s",
                (user_id,)
            )
            orders = await cur.fetchall()

    if not orders:
        return jsonify([])

    # 并发处理所有订单的详情查询
    tasks = [fetch_order_details(o) for o in orders]
    results = await asyncio.gather(*tasks)
    return jsonify(results)

# 启动时初始化连接池
if __name__ == '__main__':
    asyncio.run(init_pool())
    app.run(host='0.0.0.0', port=5000, threaded=False)  # 注意:关掉线程模式

关键改动点:

  • async def 视图函数:Flask 2.2原生支持协程视图,内部会自动跑在asyncio事件循环上。
  • asyncio.gather 并发:把原来串行的8次商品/店铺查询变成并发请求,DB往返时间从8×7ms ≈ 56ms 降到单次最慢7ms。
  • 连接池aiomysql.create_pool维护20个连接,避免每个请求新建连接的开销(原来PyMySQL每次都要握手认证,约5ms)。
  • threaded=False:Flask默认threaded=True会起线程跑每个请求,但协程视图不需要线程,开了反而浪费上下文切换。

四、踩坑与优化:三个让我半夜加班的问题

坑1:信号量限制并发查询

上线后第一次压测,QPS到300时数据库连接池直接爆了,报错Too many connections。原因:gather会把所有订单的查询同时发出,每个订单2个查询,20个连接根本不够用。解决办法:加信号量限制并发度。

from asyncio import Semaphore
# 全局信号量,限制最多10个并发查询
sem = Semaphore(10)

async def fetch_one(sql, args):
    async with sem:  # 加锁
        pool = get_pool()
        async with pool.acquire() as conn:
            async with conn.cursor(aiomysql.DictCursor) as cur:
                await cur.execute(sql, args)
                return await cur.fetchone()

坑2:uvloop和aiomysql的兼容性

我一开始用了uvloop.install(),结果aiomysql在连接池回收时出bug(RuntimeError: Event loop is closed)。查了github issue,发现uvloop 0.17.0和aiomysql 0.2.0有兼容问题。解决方案:不全局install uvloop,只在特定场景用。或者干脆不用,asyncio默认事件循环在3.10已经够快了。最后我选择不用uvloop,因为压测数据差异只有8%,不值得冒风险。

坑3:连接池泄漏

调试时发现内存缓慢上涨,排查后是pool.acquire()之后忘了release()。注意aiomysql的上下文管理器会自动释放,但如果你手动acquire()就必须手动release()。我代码里用了async with所以没问题,但同事写的旧代码没注意。

五、效果数据:wrk压测对比

硬件:16C16G CentOS 7,MySQL在同一台机器(本地socket连接)。

压测命令:wrk -t8 -c200 -d30s http://localhost:5000/api/orders?user_id=12345

指标 改造前(同步) 改造后(异步) 提升
QPS(平均值) 120 528 340%
P99延迟 2.8s 410ms 85%↓
P50延迟 450ms 180ms 60%↓
错误率(超时5s) 3.2% 0%
MySQL连接数占用 8(线程池) 峰值15(连接池) 更可控

数据解读:为什么QPS提升这么多?主要不是asyncio魔法,而是消除了N+1查询的串行等待。原来一个请求要等17次DB往返(约120ms),现在并发后只需2次往返(约14ms)。连接池复用也省掉了每次新建TCP连接的开销。asyncio本身的贡献在于:当有200个并发连接时,同步线程池要排队,而协程可以同时处理所有IO等待。

六、总结与建议

这个重构花了2天时间,收益明显。但有几个忠告:

  1. 别盲目追求asyncio。如果你的接口只有一个DB查询,没有循环调DB,同步代码完全够用。asyncio的价值在于大量IO等待的场景。

  2. 连接池一定要加信号量。否则高并发下连接数会失控。

  3. Flask的async视图有坑。Flask 2.2的async视图是跑在一个独立事件循环里的,你没法在视图里直接调用asyncio.run()(会报错)。正确做法是让视图本身是async def。如果遇到需要同步代码的第三方库,用asyncio.to_thread包装。

  4. 压测要看P99。平均延迟好看没用,P99才是用户感受。我们这次P99从2.8s降到410ms,用户投诉直接没了。

后悔没早做。这个接口上线半年,服务器从2台加到4台,现在改完后2台就够了,省下的机器成本够吃一年下午茶。


后续计划:下一步打算把orders表加上user_id联合索引(现在只查了id主键),预计还能再降30%延迟。等索引优化完再水一篇。