1. 问题背景:一个慢如蜗牛的订单查询接口

事情要从今年Q2的代码review说起。我们负责的订单服务有一个 /api/orders/batch 接口,用于批量查询订单详情(每次最多20个订单)。逻辑很简单:接收订单ID列表 → 循环调用下游三个服务(用户信息、商品信息、物流信息) → 聚合返回。

压测时发现:订单量一旦超过5个,接口响应时间就线性飙升。20个订单的场景下,平均响应时间高达1.6s,P99更是到了3.2s。而这三个下游服务的响应时间其实都只有30-50ms——问题出在同步串行调用上。

伪代码大概是这样的:

# 同步版本(问题代码)
def get_orders_detail(order_ids):
    results = []
    for oid in order_ids:
        user = requests.get(f"http://user-svc/{oid}").json()
        product = requests.get(f"http://product-svc/{oid}").json()
        logistics = requests.get(f"http://logistics-svc/{oid}").json()
        results.append({...})
    return results

20个订单 × 3个服务 = 60次串行HTTP调用,每次等待30ms就是1800ms。这还没算序列化和网络开销。

2. 环境与版本:我们踩过的坑都跟版本有关

先交代当时的运行环境:
- Python 3.9.12(3.9以下asyncio.run有坑,后面讲)
- Flask 2.2.3(同步框架,需要结合ASGI或线程池)
- aiohttp 3.8.4(用于异步HTTP客户端)
- databases 0.7.0(异步数据库驱动,基于SQLAlchemy Core)
- Redis 6.2 + aioredis 2.0.1(后来发现aioredis被合并到redis-py 4.x)
- 服务器:4核8G的阿里云ECS,CentOS 7.9

重要提示:如果你的Python版本低于3.8,asyncio的很多特性(如asyncio.Runner)不可用。我们最初在3.7上尝试,发现asyncio.run()不能嵌套调用,后来升级到3.9才解决。

3. 方案设计:从同步到异步的三个关键改造

核心思路:将Flask的同步处理函数改为协程+异步IO。具体分三步:

  1. 网络IO异步化:用aiohttp替换requests,实现HTTP调用的并发等待
  2. 数据库IO异步化:用databases库连接PostgreSQL,异步执行查询
  3. 并发控制:用asyncio.Semaphore限制并发数,防止下游服务被打挂

架构图简化为:Flask(通过Gunicorn + uvicorn worker运行)→ 协程调度 → aiohttp异步请求 + databases异步查询 → 聚合返回

这里要特别注意:Flask本身不支持异步,所以我们用Gunicorn的uvicorn worker来运行Flask应用。配置如下:

gunicorn -k uvicorn.workers.UvicornWorker -w 4 --timeout 120 app:app

每个worker内部运行事件循环,处理异步请求。

4. 核心实现:代码对比与关键细节

4.1 异步HTTP调用(核心改造)

改造前的同步代码(已简化):

import requests

def fetch_order_detail(order_id):
    user = requests.get(f"http://user-svc/{order_id}", timeout=2).json()
    product = requests.get(f"http://product-svc/{order_id}", timeout=2).json()
    logistics = requests.get(f"http://logistics-svc/{order_id}", timeout=2).json()
    return {"user": user, "product": product, "logistics": logistics}

改造后的异步代码:

import asyncio
import aiohttp
from aiohttp import ClientTimeout, TCPConnector

# 全局连接池复用,避免每次创建新连接
connector = TCPConnector(limit=100, limit_per_host=30, ttl_dns_cache=300)
timeout = ClientTimeout(total=10)

async def fetch_order_detail_async(session, order_id, sem):
    async with sem:  # 并发控制:同时最多20个协程执行HTTP请求
        async with session.get(f"http://user-svc/{order_id}", timeout=timeout) as resp:
            user = await resp.json()
        async with session.get(f"http://product-svc/{order_id}", timeout=timeout) as resp:
            product = await resp.json()
        async with session.get(f"http://logistics-svc/{order_id}", timeout=timeout) as resp:
            logistics = await resp.json()
        return {"user": user, "product": product, "logistics": logistics}

async def batch_fetch(order_ids):
    sem = asyncio.Semaphore(20)  # 控制并发度为20
    async with aiohttp.ClientSession(connector=connector) as session:
        tasks = [fetch_order_detail_async(session, oid, sem) for oid in order_ids]
        results = await asyncio.gather(*tasks, return_exceptions=True)
        # 处理异常:如果某个任务失败,记录日志并返回空
        return [r if not isinstance(r, Exception) else None for r in results]

关键细节
- TCPConnector(limit=100): 限制总连接数,防止文件描述符耗尽
- limit_per_host=30: 对每个下游服务的连接数限制,避免压垮一个服务
- ttl_dns_cache=300: DNS缓存5分钟,减少DNS解析开销
- ClientTimeout(total=10): 全局超时10秒,防止个别慢请求拖垮整体
- asyncio.Semaphore(20): 控制并发度20,这是根据下游服务容量测出来的最佳值

4.2 异步数据库查询

同步数据库操作(已简化):

import psycopg2

def get_order_from_db(order_id):
    conn = psycopg2.connect(dsn="postgresql://...")
    cur = conn.cursor()
    cur.execute("SELECT * FROM orders WHERE id=%s", (order_id,))
    return cur.fetchone()

改为异步:

import databases
from databases import Database

database = Database("postgresql+asyncpg://user:pass@host/db", min_size=5, max_size=20)

async def get_order_from_db_async(order_id):
    query = "SELECT * FROM orders WHERE id = :oid"
    return await database.fetch_one(query=query, values={"oid": order_id})

# 在应用启动时连接,关闭时断开
async def startup():
    await database.connect()

async def shutdown():
    await database.disconnect()

注意databases库底层使用asyncpg,连接池大小max_size=20需要与并发度匹配。我们一开始设了100,结果把PostgreSQL的连接数打满了,后来调小到20配合Semaphore才稳定。

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

坑1:Flask异步视图的兼容性问题
Flask 2.0+支持异步视图函数,但前提是使用ASGI服务器(如uvicorn)。我们最初用flask run调试,发现异步函数根本不执行,返回空值。后来查文档才知道,Flask的dev server是Werkzeug,不支持异步。解决方案:开发环境也用uvicorn app:app启动,或者用flask run --with-threads(但不推荐)。

坑2:aiohttp连接池泄漏
上线第一天,下游服务反馈说我们的请求突然变少了。查日志发现aiohttp报错Too many open connections。原因是我们在每个请求里都创建了新的ClientSession,没有复用。修复:将ClientSession实例化为全局变量,加上connector限制连接数。

坑3:asyncio.gather的异常处理
最初代码直接返回await asyncio.gather(*tasks),结果一个订单查询失败抛异常,整个batch都挂了。修复:加上return_exceptions=True,手动处理异常并返回None。同时增加重试逻辑:

retry_count = 3
for attempt in range(retry_count):
    try:
        return await fetch_order_detail_async(...)
    except Exception as e:
        if attempt == retry_count - 1:
            logger.error(f"Order {order_id} failed after {retry_count} retries: {e}")
            return None
        await asyncio.sleep(0.1 * (2 ** attempt))  # 指数退避

6. 效果数据:性能提升的量化对比

压测工具:locust 2.15.1,模拟100个并发用户,持续压测5分钟。

指标 同步版本 异步版本(无并发控制) 异步版本(Semaphore=20)
QPS 152 3200 2850
平均延迟 820ms 48ms 62ms
P50延迟 680ms 35ms 45ms
P99延迟 2.1s 180ms 120ms
错误率 0.3% 2.1%(连接超时) 0.05%

解读
- 无并发控制的异步版本QPS最高(3200),但错误率也高(2.1%),因为下游服务扛不住突发的并发请求,导致大量连接超时。
- 加上Semaphore=20后,QPS略降(2850),但错误率降到0.05%,P99延迟反而更低(120ms vs 180ms),因为避免了拥塞崩溃。
- 对比同步版本,QPS提升了18倍,P99延迟降低了17.5倍。

资源消耗对比(4核8G服务器):
- 同步版本:CPU 65%,内存 320MB
- 异步版本:CPU 45%,内存 280MB
异步版本不仅性能更好,资源占用反而更低,因为减少了线程切换开销。

7. 总结:异步编程的适用边界

这次改造让我深刻理解了asyncio的适用场景:
- IO密集型任务(HTTP、DB、文件读写)效果显著,计算密集型任务效果有限
- 并发度需要控制,不是越大越好。我们最终定在20,是根据下游服务的TP99和连接池大小计算出来的
- 错误处理要严谨,异步的异常传播更加隐蔽,务必加上return_exceptions和重试机制
- 版本兼容性:Python 3.8+、aiohttp 3.8+、databases 0.7+,低于这些版本容易踩坑

最后给个建议:如果你的Flask应用遇到IO瓶颈,别急着上celery或消息队列,先试试asyncio——成本最低,效果往往出人意料。前提是做好连接池管理和并发控制。

(全文完)