一、问题背景:一个被同步IO拖死的接口

去年接手了一个内部服务,功能很简单:一个聚合查询接口,接收一个商品ID,然后并发去查三个下游服务——库存服务、价格服务、评价服务,把结果合并返回。

代码大概长这样:

import requests
from flask import Flask, jsonify, request

app = Flask(__name__)

INVENTORY_URL = "http://inventory-svc:8001/stock"
PRICE_URL = "http://price-svc:8002/price"
REVIEW_URL = "http://review-svc:8003/review"

@app.route("/product/")
def get_product(pid):
    stock = requests.get(f"{INVENTORY_URL}/{pid}", timeout=2).json()
    price = requests.get(f"{PRICE_URL}/{pid}", timeout=2).json()
    review = requests.get(f"{REVIEW_URL}/{pid}", timeout=2).json()
    return jsonify({"stock": stock, "price": price, "review": review})

三个下游平均响应时间分别是40ms、50ms、60ms。因为是串行调用,光下游耗时就是150ms,加上Flask本身的处理,单请求平均200ms左右。

上线后问题很快暴露:高峰期QPS一过100,接口就开始排队,P99延迟飙到850ms,超时告警不断。用ab压测,4核8G的机器单进程QPS只有120,CPU利用率不到15%——典型的IO密集型瓶颈,CPU全在等网络。

二、环境与版本

  • Python: 3.11.6
  • asyncio: 标准库自带
  • aiohttp: 3.9.1
  • Flask: 3.0.0(保留作为WSGI入口,后面会说明)
  • 压测工具: wrk 4.2.0
  • 机器: 4核8G,CentOS 7.9

选aiohttp而不是httpx,是因为当时httpx的异步连接池在高并发下表现不如aiohttp稳定(httpx 0.25版本连接复用有bug)。这个选择后面看是对的。

三、方案设计

核心思路就一句话:把串行的三次下游调用改成并发

三个下游之间没有依赖关系,完全可以同时发出去。串行150ms,并发后取决于最慢的那个,理论60ms。

但这里有个坑:Flask是同步WSGI框架,直接在里面用asyncio.run()会阻塞worker,而且每次请求都新建事件循环,开销巨大。所以方案分两层:

  1. 框架层:把Flask换成支持ASGI的框架。我选了FastAPI(0.109.0),因为它原生支持async def路由。
  2. IO层:用aiohttp的ClientSession,全局复用一个session,利用连接池。

架构变成:

客户端 → FastAPI (uvicorn) → asyncio.gather(
                                  aiohttp GET inventory,
                                  aiohttp GET price,
                                  aiohttp GET review
                              ) → 合并返回

关键点:ClientSession必须全局复用,不能每次请求创建。每次创建session意味着每次都要重新建TCP连接、TLS握手,性能比同步还差。

四、核心实现

先看改造后的完整代码:

import asyncio
import aiohttp
from fastapi import FastAPI
from contextlib import asynccontextmanager

INVENTORY_URL = "http://inventory-svc:8001/stock"
PRICE_URL = "http://price-svc:8002/price"
REVIEW_URL = "http://review-svc:8003/review"

# 全局session,应用生命周期内复用
session: aiohttp.ClientSession = None

@asynccontextmanager
async def lifespan(app: FastAPI):
    global session
    # 连接池配置:总连接100,每个host最多50
    connector = aiohttp.TCPConnector(
        limit=100,
        limit_per_host=50,
        ttl_dns_cache=300,
        enable_cleanup_closed=True,
    )
    timeout = aiohttp.ClientTimeout(total=2, connect=0.5)
    session = aiohttp.ClientSession(connector=connector, timeout=timeout)
    yield
    await session.close()

app = FastAPI(lifespan=lifespan)

async def fetch(url: str) -> dict:
    async with session.get(url) as resp:
        resp.raise_for_status()
        return await resp.json()

@app.get("/product/{pid}")
async def get_product(pid: str):
    # 三个请求并发发出
    results = await asyncio.gather(
        fetch(f"{INVENTORY_URL}/{pid}"),
        fetch(f"{PRICE_URL}/{pid}"),
        fetch(f"{REVIEW_URL}/{pid}"),
        return_exceptions=True,
    )
    stock, price, review = results
    # 单个下游失败不影响整体,降级处理
    return {
        "stock": stock if not isinstance(stock, Exception) else None,
        "price": price if not isinstance(price, Exception) else None,
        "review": review if not isinstance(review, Exception) else None,
    }

启动命令:

uvicorn main:app --host 0.0.0.0 --port 8000 --workers 4 --loop uvloop

注意--loop uvloop,uvloop是asyncio事件循环的Cython实现,比标准库快2-4倍。这个参数在高并发下差别很明显。

asyncio.gatherreturn_exceptions=True是个关键细节。默认情况下,只要有一个协程抛异常,gather会立刻抛出,其他协程的结果就丢了。加上这个参数后,异常会作为结果返回,我们可以做降级——库存查不到不影响返回价格。

五、踩坑与优化

坑1:ClientSession创建时机

第一版我图省事,在路由函数里创建session:

async def get_product(pid: str):
    async with aiohttp.ClientSession() as session:  # 错误示范
        ...

压测QPS只有400,比预期低一个数量级。原因是每次请求都新建连接池,TCP连接无法复用,大量时间花在握手和TIME_WAIT上。改成全局session后,QPS直接翻倍。

坑2:连接池limit设太小

一开始connector的limit用的默认值100,但limit_per_host也是默认的0(不限制)。实际压测发现单host连接数飙到几百,下游服务开始拒绝连接。后来显式设置limit_per_host=50,配合下游的承载能力调优。

坑3:DNS解析阻塞

aiohttp默认每次请求都做DNS解析,虽然走的是线程池,但在高QPS下线程池会成为瓶颈。加上ttl_dns_cache=300后,DNS结果缓存5分钟,这块开销基本消失。

坑4:uvicorn worker数量

一开始只开1个worker,QPS卡在900上不去。因为虽然IO是异步的,但JSON序列化、路由匹配这些还是占CPU。改成4个worker(等于CPU核数)后,QPS跑到2800。注意worker多了session是每个worker独立的,连接池要按worker数分摊考虑。

坑5:超时设置

aiohttp的ClientTimeout如果不设,默认是5分钟。下游一旦卡住,连接就被占着不放。设置total=2, connect=0.5后,快速失败快速释放,P99明显改善。

六、效果数据

用wrk压测,命令:

wrk -t4 -c200 -d30s --latency http://localhost:8000/product/12345

对比数据:

指标 改造前(Flask+requests) 改造后(FastAPI+aiohttp)
QPS 120 2800
P50延迟 180ms 42ms
P99延迟 850ms 95ms
CPU利用率 15% 68%
内存占用 180MB 240MB

QPS提升23倍,P99降低89%。内存略涨是因为连接池和4个worker的开销,可以接受。

几个观察:

  1. CPU利用率从15%涨到68%,说明瓶颈终于从IO转移到了CPU,这是健康的。
  2. 如果继续加worker到8个,QPS还能涨到3500左右,但P99开始抖动,因为上下文切换开销上来了。4核机器4个worker是甜点。
  3. 单请求延迟从200ms降到45ms,基本就是最慢下游(60ms)加上框架开销,符合预期。

七、总结

这次改造的核心其实就三点:

  1. 并发替代串行:asyncio.gather把三个串行IO变成并发,理论延迟从sum变成max。
  2. 连接复用:全局ClientSession + 合理的连接池配置,避免重复握手。
  3. 运行时优化:uvloop + 多worker,把CPU也用起来。

但asyncio不是银弹。如果你的接口是CPU密集型(比如大量计算、图像处理),异步化没有任何帮助,反而会因为事件循环的调度开销变慢。异步只解决IO等待问题。

另外,异步代码的调试比同步麻烦得多,堆栈信息不直观,异常容易静默丢失。建议在gather里统一加return_exceptions,配合日志把异常打出来。

最后,如果你的服务还在用requests串行调下游,且QPS上不去、CPU利用率低,那基本可以确定是IO瓶颈,值得试试asyncio。改造成本不高,收益很明显。