一、问题背景

去年接手了一个聚合查询服务,逻辑很简单:客户端传一个商品ID进来,服务端要同时去查三个下游接口——库存服务、价格服务、促销服务,把结果拼起来返回。

原来的实现是Flask + requests,典型写法:

# before: app.py
from flask import Flask, jsonify
import requests

app = Flask(__name__)

INVENTORY_URL = "http://inventory.internal/api/stock"
PRICE_URL = "http://price.internal/api/price"
PROMO_URL = "http://promo.internal/api/promo"

@app.route("/product/")
def get_product(pid):
    stock = requests.get(INVENTORY_URL, params={"pid": pid}, timeout=2).json()
    price = requests.get(PRICE_URL, params={"pid": pid}, timeout=2).json()
    promo = requests.get(PROMO_URL, params={"pid": pid}, timeout=2).json()
    return jsonify({
        "pid": pid,
        "stock": stock["count"],
        "price": price["amount"],
        "promo": promo.get("tag"),
    })

这段代码的问题非常明显:三个下游请求是串行的。每个接口平均响应60~80ms,三个加起来就是200ms打底。下游稍微抖动一下,单个请求轻松上500ms。

线上表现:4核8G的机器,gunicorn开4个worker,每个worker 8个线程,单机QPS稳定在120左右,P99延迟1.2秒。大促前压测,8台机器才勉强扛住1000 QPS,成本高得离谱。

其实这三个请求彼此完全独立,没有依赖关系。串行执行纯属浪费。哪怕用线程池并发,也能把200ms压到80ms。但线程模型有GIL和上下文切换的开销,高并发下线程数会爆炸。真正合适的方案是asyncio——I/O密集型场景,协程才是正解。

二、环境与版本

  • Python 3.11.6(3.11的asyncio在异常处理和task创建上比3.8快不少,官方benchmark显示TaskGroup场景有约20%提升)
  • aiohttp 3.9.3
  • FastAPI 0.109.2 + uvicorn 0.27.0(既然都上asyncio了,框架也一起换掉)
  • gunicorn 21.2.0(只用它做进程管理,worker_class用uvicorn.workers.UvicornWorker)
  • 压测工具:wrk 4.2.0,脚本模式,8线程200连接,持续60秒

如果你还在Python 3.8,建议至少升到3.10,asyncio.timeout()这个上下文管理器是3.11才有的,3.10之前只能手动wait_for,写起来啰嗦。

三、方案设计

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

具体设计上考虑了几点:

  1. 框架选型:Flask是WSGI,天生同步,硬要在里面跑asyncio得用asyncio.run()包一层,但每个请求都创建一个新event loop,开销巨大。直接换FastAPI,原生ASGI,event loop全局复用。
  2. HTTP客户端:requests是同步阻塞的,换成aiohttp。关键是ClientSession要复用,不能每个请求都新建,否则连接池形同虚设。
  3. 并发控制:三个请求用asyncio.gather并发,但要加return_exceptions=True,避免一个下游挂了把整个请求带崩。
  4. 超时:整体超时2秒,单个下游超时1.5秒,超时的那个降级返回默认值,不影响其他两个。
  5. 连接池:ClientSession的connector限制总连接数100,每个host限制20,避免打爆下游。

四、核心实现

先看重构后的核心代码:

# after: main.py
import asyncio
import aiohttp
from fastapi import FastAPI
from contextlib import asynccontextmanager

INVENTORY_URL = "http://inventory.internal/api/stock"
PRICE_URL = "http://price.internal/api/price"
PROMO_URL = "http://promo.internal/api/promo"

session: aiohttp.ClientSession | None = None

@asynccontextmanager
async def lifespan(app: FastAPI):
    global session
    connector = aiohttp.TCPConnector(
        limit=100,            # 总连接池上限
        limit_per_host=20,    # 单host上限
        ttl_dns_cache=300,    # DNS缓存5分钟
        enable_cleanup_closed=True,
    )
    timeout = aiohttp.ClientTimeout(total=1.5, connect=0.5)
    session = aiohttp.ClientSession(connector=connector, timeout=timeout)
    yield
    await session.close()

app = FastAPI(lifespan=lifespan)


async def fetch_json(url: str, pid: str) -> dict:
    try:
        async with session.get(url, params={"pid": pid}) as resp:
            resp.raise_for_status()
            return await resp.json()
    except (aiohttp.ClientError, asyncio.TimeoutError) as e:
        # 记录日志,返回降级值
        return {"_error": str(e)}


@app.get("/product/{pid}")
async def get_product(pid: str):
    async with asyncio.timeout(2.0):  # 整体兜底超时,Python 3.11+
        stock_task = asyncio.create_task(fetch_json(INVENTORY_URL, pid))
        price_task = asyncio.create_task(fetch_json(PRICE_URL, pid))
        promo_task = asyncio.create_task(fetch_json(PROMO_URL, pid))

        stock, price, promo = await asyncio.gather(
            stock_task, price_task, promo_task
        )

    return {
        "pid": pid,
        "stock": stock.get("count", 0) if "_error" not in stock else -1,
        "price": price.get("amount", 0) if "_error" not in price else -1,
        "promo": promo.get("tag") if "_error" not in promo else None,
    }

启动命令:

gunicorn main:app \
  -w 4 \
  -k uvicorn.workers.UvicornWorker \
  --bind 0.0.0.0:8000 \
  --timeout 30 \
  --max-requests 10000 \
  --max-requests-jitter 1000

这里几个点值得展开说:

为什么用create_task而不是直接把协程扔给gather 其实gather可以直接接收协程,效果一样。但显式create_task能让三个请求在进入gather之前就已经开始调度,语义更清楚,调试时也方便单独cancel某个task。

为什么timeout分两层? 单请求1.5秒是给下游的,整体2秒是给自己的。如果三个下游都超时,1.5秒就能返回;如果某个下游响应慢但没超时,整体2秒兜底,防止请求堆积。

ClientSession为什么放lifespan? 每个请求新建session会导致TCP连接无法复用,三次握手开销直接吃掉并发收益。实测新建session的版本QPS只有复用版本的三分之一。

五、踩坑与优化

坑1:event loop阻塞。 一开始我在异步函数里用了json.loads处理大响应,还有一处用了requests做埋点上报,结果QPS上不去。asyncio是单线程event loop,任何同步阻塞调用都会卡住整个loop。后来全部换成orjsonaiohttp,QPS直接翻倍。

坑2:连接池太小。 最初limit_per_host设的10,压测到800 QPS就上不去了,日志里全是Connection pool is full, discarding connection。改成20后,单机跑到1800 QPS稳定。

坑3:DNS解析。 内网服务用的是域名,aiohttp默认每次请求都走一次DNS(虽然是getaddrinfo,但在高并发下也是瓶颈)。加上ttl_dns_cache=300后,P99降了约30ms。

坑4:gunicorn worker数。 一开始按CPU核数配了4个worker,后来发现asyncio场景下worker不是越多越好。4核机器上,4个worker跑满CPU,但每个worker的event loop都有大量空闲。实测3个worker + 每个worker更高并发,整体吞吐更高。最终用4 worker是考虑到故障隔离。

坑5:asyncio.gather的异常传播。 默认return_exceptions=False,任何一个task抛异常,gather会立即抛出,但其他task不会被取消,会继续跑完。这会导致"请求已经返回了,后台还在发请求"的诡异现象。要么用return_exceptions=True,要么在except里手动cancel。我选了前者,配合fetch_json内部的try/except,双保险。

六、效果数据

压测条件:4核8G,8台机器压到2台,wrk 8线程200连接,60秒。

指标 Before (Flask+requests) After (FastAPI+aiohttp) 提升
单机QPS 120 1800 15x
P50延迟 210ms 65ms 3.2x
P99延迟 1200ms 180ms 6.7x
单机CPU 85% 62% -
部署机器数 8 2 -75%
大促峰值承载 1000 QPS 3600 QPS 3.6x

几个细节:

  • P99从1.2秒到180ms,主要是消除了串行等待。理论上三个60ms的请求并发后应该是60ms,实测65ms,多出来的5ms是协程调度和结果聚合的开销,可以接受。
  • CPU从85%降到62%,因为线程上下文切换没了。原来8个线程抢GIL,现在单event loop跑协程,切换成本极低。
  • 机器从8台减到2台,一年省下来的云成本够发好几个月的年终奖(笑)。

一个反直觉的点:单机QPS到1800之后,继续加连接数收益递减。200连接时QPS 1800,400连接时QPS只到1950,P99反而涨到250ms。原因是下游三个服务本身有承载上限,我们的并发把压力传导过去了。所以最终生产环境我限了limit_per_host=15,宁可自己慢一点,也不能把下游打挂。

七、总结

这次重构的核心其实就一句话:I/O密集型场景,串行改并发,同步改异步。技术上没什么黑魔法,但落地时细节很多——event loop不能阻塞、连接池要复用、超时要分层、gather的异常语义要搞清楚,每一个踩错都会让性能打对折。

几点经验:

  1. 别在WSGI里硬塞asyncio,该换框架就换。FastAPI的学习成本一天就能补上,收益是十几倍QPS。
  2. ClientSession、连接池、DNS缓存是aiohttp性能的三驾马车,缺一不可。
  3. 超时一定要分层,单请求超时 + 整体超时,配合降级逻辑,下游抖动时你的服务才不会雪崩。
  4. 压测要压到瓶颈,找到拐点再回退10%,留出安全边际。QPS不是越高越好,稳定才是。

最后提醒一句:asyncio不是银弹。如果你的服务是CPU密集型(比如图像处理、大量计算),上asyncio只会更慢,该用多进程就用多进程。选对场景,才是真正的优化。