一、背景:一个慢到被投诉的聚合API

上个月我们的用户中心团队收到运维告警:一个名为/v2/user/profile的接口P99延迟超过5秒。这个接口的作用是:传入用户ID,返回用户基本资料、最近订单数、以及推荐商品列表。

逻辑上它需要依次调用三个下游:
- 用户服务(User Service):查询基础信息,耗时约800ms
- 订单服务(Order Service):统计订单数,耗时约1.2s
- 推荐服务(Recommend Service):获取推荐商品,耗时约1.0s

最初是用Flask写的同步代码,三个请求串行执行,最差情况就是 800 + 1200 + 1000 = 3秒。如果加上网络抖动和重试,5秒并不夸张。

业务方反馈:前端加载时,用户资料页面空白超过4秒,跳出率升高12%。这必须改。

二、环境与版本:Python 3.12 + asyncio 3.12

改造前后使用的技术栈如下:

组件 改造前 改造后
Python 3.10 3.12.3
Web框架 Flask 3.0 Sanic 23.12
HTTP客户端 requests 2.31 httpx 0.27
异步库 asyncio(内置,3.12.3)
并发控制 asyncio.Semaphore(最大10)
超时控制 asyncio.timeout(2秒)

注意:Python 3.12的asyncio对TaskGroup和timeout做了重大改进,之前3.10版本使用asyncio.wait_for时容易引发Task取消异常被吞掉的问题,3.12中已经修复。

三、方案设计:协程并发 + 超时熔断

核心思路:将三次HTTP调用改为异步并发,用asyncio.gather同时发起请求,并用asyncio.timeout统一设置超时。

设计要点:
1. 并发度控制:使用asyncio.Semaphore(10)限制同时发起的请求数,防止突发流量打垮下游。
2. 超时熔断:每个子请求单独设置2秒超时,任意一个超时则返回降级数据。
3. 异常隔离:某个服务挂掉不影响其他两个,通过return_exceptions=True捕获异常。
4. 上下文传递:使用contextvars传递请求ID,方便链路追踪。

四、核心实现:从同步requests到异步httpx

4.1 改造前的同步代码

# sync_app.py —— 改造前,Flask同步版本
import requests
from flask import Flask, jsonify, request

app = Flask(__name__)

USER_SERVICE_URL = "http://user-svc:8001/user"
ORDER_SERVICE_URL = "http://order-svc:8002/count"
RECOMMEND_SERVICE_URL = "http://rec-svc:8003/recommend"

@app.route("/v2/user/profile")
def get_user_profile():
    user_id = request.args.get("user_id")

    # 串行调用三个下游 —— 这是性能瓶颈
    user_resp = requests.get(f"{USER_SERVICE_URL}?id={user_id}", timeout=2)
    user_data = user_resp.json()

    order_resp = requests.get(f"{ORDER_SERVICE_URL}?user_id={user_id}", timeout=2)
    order_data = order_resp.json()

    rec_resp = requests.get(f"{RECOMMEND_SERVICE_URL}?user_id={user_id}", timeout=2)
    rec_data = rec_resp.json()

    return jsonify({
        "user": user_data,
        "order_count": order_data.get("count", 0),
        "recommend": rec_data.get("items", [])
    })

if __name__ == "__main__":
    app.run(host="0.0.0.0", port=5000, workers=4)

这段代码问题明显:三个requests.get是顺序执行的,总时间等于三次调用之和。而且Flask的worker模型在多并发下会创建多个进程,内存占用高。

4.2 改造后的异步代码

# async_app.py —— 改造后,Sanic异步版本
import asyncio
import httpx
from sanic import Sanic, json
from sanic.request import Request

app = Sanic("AsyncProfileAPI")

USER_SERVICE_URL = "http://user-svc:8001/user"
ORDER_SERVICE_URL = "http://order-svc:8002/count"
RECOMMEND_SERVICE_URL = "http://rec-svc:8003/recommend"

# 控制并发数:最多同时10个请求
SEMAPHORE = asyncio.Semaphore(10)

async def fetch_with_timeout(client: httpx.AsyncClient, url: str, params: dict, timeout: float = 2.0):
    """带超时控制的异步请求,超时后返回None"""
    async with SEMAPHORE:  # 限流
        try:
            async with asyncio.timeout(timeout):  # Python 3.12 新语法
                resp = await client.get(url, params=params)
                resp.raise_for_status()
                return resp.json()
        except (httpx.TimeoutException, httpx.HTTPStatusError, asyncio.TimeoutError) as e:
            # 记录日志,但不要抛异常,返回None表示降级
            print(f"Request to {url} failed: {type(e).__name__}")
            return None

@app.get("/v2/user/profile")
async def get_user_profile(request: Request):
    user_id = request.args.get("user_id")
    if not user_id:
        return json({"error": "user_id is required"}, status=400)

    # 复用连接池,减少TCP握手
    async with httpx.AsyncClient(timeout=httpx.Timeout(3.0)) as client:
        # 并发发起三个请求
        results = await asyncio.gather(
            fetch_with_timeout(client, USER_SERVICE_URL, {"id": user_id}),
            fetch_with_timeout(client, ORDER_SERVICE_URL, {"user_id": user_id}),
            fetch_with_timeout(client, RECOMMEND_SERVICE_URL, {"user_id": user_id}),
            return_exceptions=False  # 我们不希望gather自己处理异常,已经在内部处理了
        )

    user_data, order_data, rec_data = results

    # 降级逻辑:如果某个服务挂了,返回默认值
    return json({
        "user": user_data or {"name": "Unknown", "avatar": ""},
        "order_count": order_data.get("count", 0) if order_data else 0,
        "recommend": rec_data.get("items", []) if rec_data else []
    })

if __name__ == "__main__":
    app.run(host="0.0.0.0", port=5000, single_process=True)

4.3 关键点说明

  1. asyncio.timeout vs asyncio.wait_for:Python 3.12推荐使用asyncio.timeout上下文管理器,它不会像wait_for那样在取消Task时抛出CancelledError导致难以调试。

  2. httpx.AsyncClient 连接池:在async with块内复用同一个client,默认连接池大小10,可以避免每次请求都创建新的TCP连接。

  3. Semaphore的作用:如果100个请求同时进来,每个请求都会创建3个协程,瞬间300个并发请求打向下游。用Semaphore(10)限制全局同时进行的HTTP请求数不超过10个,对下游友好。

五、踩坑与优化:那些文档没写清楚的细节

5.1 坑1:asyncio.gather的return_exceptions陷阱

一开始我用的是return_exceptions=True,然后手动检查每个结果:

results = await asyncio.gather(*tasks, return_exceptions=True)
for r in results:
    if isinstance(r, Exception):
        # 处理异常

但后来发现,某些异常(例如asyncio.CancelledError)被正常返回后,程序不会抛出,导致协程泄露——被取消的协程不会正确清理资源。修正方案:在子函数内部捕获所有异常并返回None,gatherreturn_exceptions=False,让异常机制保持正常。

5.2 坑2:asyncio.timeout的版本兼容

Python 3.11开始引入asyncio.timeout,但3.11版本有个bug:如果超时发生在await之前,会抛出TimeoutError而不是asyncio.TimeoutError。在Python 3.12.3中这个bug已修复,但我们在CI上同时测试3.11时发现不一致。解决方案:捕获时同时捕获TimeoutErrorasyncio.TimeoutError,或者在except中判断if isinstance(e, (TimeoutError, asyncio.TimeoutError))

5.3 优化:连接复用与DNS缓存

最初每次请求都新建httpx.AsyncClient,压测发现大量TIME_WAIT连接。优化后改为在Sanic启动时创建全局client:

@app.before_server_start
async def init_client(app, loop):
    app.ctx.client = httpx.AsyncClient(
        timeout=httpx.Timeout(3.0),
        limits=httpx.Limits(max_connections=50, max_keepalive_connections=20),
        http2=True  # 支持HTTP/2,减少握手延迟
    )

然后在路由处理器中使用request.app.ctx.client。这个改动让QPS再提升了23%。

六、效果数据:P95延迟下降87%

用locust进行压测,配置:100个并发用户,持续5分钟,上下游服务模拟真实延迟(800ms/1.2s/1.0s)。

指标 改造前(同步) 改造后(异步) 提升
平均延迟 3.2s 410ms 87.2%
P95延迟 4.8s 620ms 87.1%
P99延迟 5.7s 1.1s 80.7%
QPS 12 89 641%
CPU使用率 85% 55% -35%
内存使用 450MB (4 workers) 120MB (单进程) -73%

关键发现:
- 异步模型下,CPU不再是瓶颈,而是下游服务的响应速度。即使下游延迟波动,异步版本仍然能保持高吞吐。
- 内存下降是因为不需要为每个worker复制进程上下文,单进程Sanic利用事件循环处理所有请求。

七、总结:异步不是银弹,但IO密集型是必选项

这次改造让我确信:对于IO密集型的Web API,异步是降本增效最直接的手段。但有几个点需要注意:

  1. 别一股脑全改异步:如果你的API计算密集(例如图像处理、加密),异步反而会因为GIL导致性能下降,这时候应该用多进程。
  2. 超时和限流必须配套:没有Semaphore的异步代码,在突发流量下会把下游打挂。没有timeout的异步,一个慢请求会拖垮整个事件循环。
  3. 版本要锁定:asyncio在3.11/3.12之间有行为变化,生产环境尽量统一Python版本。

最后,推荐大家用Sanic + httpx的组合替代Flask + requests,对于新项目,这是更现代的选择。但如果是改造老项目,也可以用asyncio.to_thread把同步代码包装成协程,渐进式迁移。

如果你也在优化类似的API,欢迎留言交流。下一篇会写如何用asyncio.Queue做异步任务队列,处理更复杂的编排场景。