一、背景:一个慢到被投诉的聚合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 关键点说明
-
asyncio.timeoutvsasyncio.wait_for:Python 3.12推荐使用asyncio.timeout上下文管理器,它不会像wait_for那样在取消Task时抛出CancelledError导致难以调试。 -
httpx.AsyncClient连接池:在async with块内复用同一个client,默认连接池大小10,可以避免每次请求都创建新的TCP连接。 -
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,gather的return_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时发现不一致。解决方案:捕获时同时捕获TimeoutError和asyncio.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,异步是降本增效最直接的手段。但有几个点需要注意:
- 别一股脑全改异步:如果你的API计算密集(例如图像处理、加密),异步反而会因为GIL导致性能下降,这时候应该用多进程。
- 超时和限流必须配套:没有Semaphore的异步代码,在突发流量下会把下游打挂。没有timeout的异步,一个慢请求会拖垮整个事件循环。
- 版本要锁定:asyncio在3.11/3.12之间有行为变化,生产环境尽量统一Python版本。
最后,推荐大家用Sanic + httpx的组合替代Flask + requests,对于新项目,这是更现代的选择。但如果是改造老项目,也可以用asyncio.to_thread把同步代码包装成协程,渐进式迁移。
如果你也在优化类似的API,欢迎留言交流。下一篇会写如何用asyncio.Queue做异步任务队列,处理更复杂的编排场景。