一、问题背景

去年底我们做了一个内部运营后台,其中一个接口 /api/user/orders 需要聚合三个下游服务:用户中心(获取用户等级)、订单服务(拉取最近订单)、风控服务(检查用户是否被限制)。这三个服务都是HTTP接口,响应时间分别在80ms、200ms、150ms左右。

最初的实现非常直接——用Flask写了个同步视图函数,顺序调用三个服务,然后拼装返回。逻辑上没问题,但上线后运营同学抱怨“点一下要等半天”。我用ab压测了一下,单机QPS只有120,P99延迟1.2秒。更糟的是,当某个下游服务变慢时,整个接口会被拖死。

我决定不改架构、不引入新服务,只用Python的asyncio把这三个串行调用改成并发。结果QPS提升到1800,P99降到85ms。下面把过程和代码完整分享出来。

二、环境与版本

  • Python 3.11.4(3.8+都支持,但3.11的asyncio性能更好)
  • Flask 2.3.2(仅作为Web框架,不依赖其异步能力)
  • aiohttp 3.8.5(用于异步HTTP客户端)
  • 压测工具:wrk 4.2.0,参数 wrk -t4 -c100 -d30s --latency
  • 部署:单机 4核8G,Ubuntu 22.04,Gunicorn 20.1.0 + gevent 22.10.2(before用gevent,after用asyncio)

注意:我没有用Flask的async视图,因为Flask本身对async支持有限。我的做法是在Flask同步视图里通过 asyncio.run() 调用一个异步函数。虽然这会创建新事件循环,但配合Gunicorn的gevent worker,整体效果依然很好。更优雅的方案是直接用FastAPI或Quart,但本文聚焦asyncio改造,不换框架。

三、方案设计

核心思路:把三个下游调用从串行改为并发。用 asyncio.gather 同时发起三个HTTP请求,总耗时从 t1+t2+t3 降到 max(t1,t2,t3)

但有几个细节要处理:
1. 每个下游调用需要独立的超时和异常处理,不能因为一个失败就全盘失败。
2. 需要复用aiohttp的ClientSession,避免每次请求都新建连接。
3. 事件循环不能每次请求都创建和销毁,否则开销很大。

我的设计是:
- 在Flask应用启动时创建一个全局的aiohttp ClientSession,绑定到一个后台事件循环。
- 视图函数通过 asyncio.run_coroutine_threadsafe 把协程提交到该循环,并等待结果。
- 这样既复用了连接池,又避免了频繁创建事件循环。

四、核心实现

4.1 Before:同步串行版本

# app_sync.py
import requests
from flask import Flask, jsonify, request

app = Flask(__name__)

USER_SVC = "http://user-svc:8001"
ORDER_SVC = "http://order-svc:8002"
RISK_SVC = "http://risk-svc:8003"

def fetch_user(uid):
    r = requests.get(f"{USER_SVC}/user/{uid}", timeout=1.0)
    return r.json()

def fetch_orders(uid):
    r = requests.get(f"{ORDER_SVC}/orders/{uid}", timeout=1.0)
    return r.json()

def fetch_risk(uid):
    r = requests.get(f"{RISK_SVC}/check/{uid}", timeout=1.0)
    return r.json()

@app.route("/api/user/orders")
def get_user_orders():
    uid = request.args.get("uid")
    user = fetch_user(uid)
    orders = fetch_orders(uid)
    risk = fetch_risk(uid)
    return jsonify({
        "user": user,
        "orders": orders,
        "risk": risk
    })

压测命令:wrk -t4 -c100 -d30s --latency http://127.0.0.1:5000/api/user/orders?uid=123

结果:QPS 120,P99 1.2s,平均延迟 830ms。

4.2 After:asyncio并发版本

# app_async.py
import asyncio
import aiohttp
from flask import Flask, jsonify, request

app = Flask(__name__)

USER_SVC = "http://user-svc:8001"
ORDER_SVC = "http://order-svc:8002"
RISK_SVC = "http://risk-svc:8003"

# 全局事件循环和session
loop = asyncio.new_event_loop()
session = None

async def init_session():
    global session
    connector = aiohttp.TCPConnector(
        limit=200,           # 总连接池大小
        limit_per_host=50,   # 每个下游最多50个连接
        ttl_dns_cache=300,   # DNS缓存5分钟
        keepalive_timeout=30 # 长连接保持30秒
    )
    timeout = aiohttp.ClientTimeout(total=1.0, connect=0.3)
    session = aiohttp.ClientSession(connector=connector, timeout=timeout)

async def fetch_user(uid):
    try:
        async with session.get(f"{USER_SVC}/user/{uid}") as resp:
            return await resp.json()
    except Exception as e:
        return {"error": str(e)}

async def fetch_orders(uid):
    try:
        async with session.get(f"{ORDER_SVC}/orders/{uid}") as resp:
            return await resp.json()
    except Exception as e:
        return {"error": str(e)}

async def fetch_risk(uid):
    try:
        async with session.get(f"{RISK_SVC}/check/{uid}") as resp:
            return await resp.json()
    except Exception as e:
        return {"error": str(e)}

async def fetch_all(uid):
    user, orders, risk = await asyncio.gather(
        fetch_user(uid),
        fetch_orders(uid),
        fetch_risk(uid),
        return_exceptions=True
    )
    return {"user": user, "orders": orders, "risk": risk}

@app.route("/api/user/orders")
def get_user_orders():
    uid = request.args.get("uid")
    future = asyncio.run_coroutine_threadsafe(fetch_all(uid), loop)
    result = future.result(timeout=2.0)
    return jsonify(result)

def start_background_loop():
    asyncio.set_event_loop(loop)
    loop.run_until_complete(init_session())
    loop.run_forever()

# 在Gunicorn的post_fork钩子中启动后台循环
# gunicorn.conf.py
# def post_fork(server, worker):
#     import threading
#     from app_async import start_background_loop
#     t = threading.Thread(target=start_background_loop, daemon=True)
#     t.start()

启动命令:gunicorn -c gunicorn.conf.py -w 4 -k gevent app_async:app

压测结果:QPS 1800,P99 85ms,平均延迟 52ms。

五、踩坑与优化

坑1:事件循环不能每次请求都创建。 最初我在视图里直接 asyncio.run(fetch_all(uid)),QPS只有400。因为每次请求都创建和销毁事件循环,还要重新建立TCP连接。改成全局循环+ClientSession后,QPS翻了4倍。

坑2:aiohttp的TCPConnector默认limit是100,但limit_per_host是0(无限制)。 这会导致某个下游服务被过多连接打挂。我显式设置 limit=200, limit_per_host=50,既保证并发,又保护下游。

坑3:超时设置要分层。 我一开始只设了 total=1.0,结果DNS解析慢的时候整个请求卡住。后来加上 connect=0.3,把连接阶段单独限制,整体更稳定。

优化点:return_exceptions=True 加上,这样某个下游失败不会导致整个接口500,而是返回部分数据+错误信息。运营后台对部分失败是可以接受的。

六、效果数据

指标 Before(同步) After(asyncio) 提升
QPS 120 1800 15倍
P99延迟 1.2s 85ms 14倍
平均延迟 830ms 52ms 16倍
CPU使用率 35% 68% 提高但可接受
内存占用 420MB 510MB 略增

压测条件:4核8G,Gunicorn 4 workers,wrk 4线程100连接30秒。下游服务用mock模拟固定延迟(80ms/200ms/150ms)。

注意:After版本CPU使用率更高,因为并发处理更多请求。但单机QPS从120到1800,单位请求的CPU成本其实下降了。

七、总结

asyncio不是银弹,但在IO密集型场景下,它用很小的改造代价换来了数量级的性能提升。我的经验是:

  1. 不要为了异步而异步。如果下游是CPU密集型,asyncio帮不了你。
  2. 事件循环和连接池一定要复用,否则性能还不如同步。
  3. 超时和异常处理要分层,保证部分失败不影响整体。
  4. 如果新项目,直接上FastAPI;老项目用本文的 run_coroutine_threadsafe 方案也能平滑迁移。

完整代码我放在GitHub gist上(链接略),有需要的同学可以自取。如果你也在做类似改造,欢迎留言交流。