1. 问题背景:一个把CPU吃满的“简单”聚合接口

上个月接手一个内部数据平台,运维告警说有个接口GET /api/user/overview频繁触发CPU 85%+监控。查了代码,逻辑确实“简单”:前端请求一次用户总览,后端需要依次调用4个内部HTTP服务——

  1. 用户基础信息(user-service)
  2. 订单统计(order-service)
  3. 积分流水(points-service)
  4. 最近登录日志(audit-service)

原实现用了requests库,一个接一个地同步调用:

def get_user_overview(user_id):
    user = requests.get(f"http://user-service/users/{user_id}", timeout=2).json()
    orders = requests.get(f"http://order-service/orders?user_id={user_id}", timeout=2).json()
    points = requests.get(f"http://points-service/points?user_id={user_id}", timeout=2).json()
    logs = requests.get(f"http://audit-service/logs?user_id={user_id}", timeout=2).json()
    return {"user": user, "orders": orders, "points": points, "logs": logs}

每个上游服务平均响应400-500ms,串行下来总耗时1.8秒。罪魁祸首不是代码效率,而是IO等待 —— 线程在recv系统调用上被挂起,GIL虽然会释放,但线程切换和内存分配开销在高并发下直接吃满CPU。压测数据(wrk 4.2.0,8线程,200连接,压测60秒):

指标 同步版本
平均延迟 1870ms
P95延迟 2100ms
吞吐量 120 req/s
CPU使用率 85%

2. 环境与版本:为什么选asyncio而不是多线程

  • Python 3.10.8(原生支持asyncio.run()asyncio.timeout()
  • aiohttp 3.8.4(异步HTTP客户端,底层基于asyncio + yarl
  • Flask 2.2.3(主Web框架,通过asyncio.run()桥接协程)
  • wrk 4.2.0(压测工具)

为什么不用concurrent.futures.ThreadPoolExecutor 我试过,线程池大小设为64,吞吐量能到300 req/s,但有两个问题:一是线程创建/切换开销依然存在,CPU在600并发时飙升到70%;二是每个线程维护独立栈,内存占用大(每个线程约8MB,峰值占用2GB+)。asyncio是单线程事件循环,协程切换开销在微秒级,内存占用固定,更适合IO密集场景。

3. 方案设计:三层改造策略

3.1 网络层:requests → aiohttp

同步requests.get()阻塞调用,改为aiohttp.ClientSession.get()异步协程。关键点:ClientSession是重量级对象,必须全局复用(内部维护连接池),不能每次请求都创建。

3.2 并发层:串行 → 协程并发

asyncio.gather()并发发起4个请求。但注意:不能无限并发 —— 上游服务有负载上限,全放开会导致上游雪崩。用asyncio.Semaphore限制并发数,我这里设为32(压测得出,上游服务在50并发时开始报错)。

3.3 容错层:超时与异常隔离

单个上游超时不能拖垮整个接口。使用asyncio.timeout()(Python 3.10+)设置每次请求超时2秒,用return_exceptions=True参数让gather()不因单个任务失败而中断其它任务。

4. 核心实现:完整异步代码

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

app = Flask(__name__)
semaphore = asyncio.Semaphore(32)  # 全局限流,保护上游
session = None  # 全局复用 ClientSession

# 每个上游服务的独立配置
SERVICES = {
    "user": {"url": "http://user-service/users/{user_id}", "timeout": 2.0},
    "orders": {"url": "http://order-service/orders?user_id={user_id}", "timeout": 2.5},
    "points": {"url": "http://points-service/points?user_id={user_id}", "timeout": 1.5},
    "logs": {"url": "http://audit-service/logs?user_id={user_id}", "timeout": 1.0},
}

async def fetch_one(session, name, user_id):
    """单个上游请求,带信号量限流和超时控制"""
    cfg = SERVICES[name]
    url = cfg["url"].format(user_id=user_id)
    async with semaphore:  # 并发数超限时等待
        try:
            # Python 3.10+ 的 asyncio.timeout 替代老式 wait_for
            async with asyncio.timeout(cfg["timeout"]):
                async with session.get(url) as resp:
                    if resp.status != 200:
                        return {name: {"error": f"HTTP {resp.status}"}}
                    return {name: await resp.json()}
        except asyncio.TimeoutError:
            return {name: {"error": "timeout"}}
        except aiohttp.ClientError as e:
            return {name: {"error": f"conn_error: {str(e)[:50]}"}}

async def get_overview_async(user_id):
    """并发聚合4个上游服务"""
    global session
    if session is None:
        # 连接池参数:总连接数100,每个host限制20
        session = aiohttp.ClientSession(
            connector=aiohttp.TCPConnector(limit=100, limit_per_host=20)
        )
    tasks = [fetch_one(session, name, user_id) for name in SERVICES]
    results = await asyncio.gather(*tasks, return_exceptions=True)
    # 合并结果,按服务名聚合
    merged = {}
    for r in results:
        if isinstance(r, dict):
            merged.update(r)
    return merged

@app.route("/api/user/overview")
def user_overview():
    user_id = request.args.get("user_id", type=int)
    if not user_id:
        return jsonify({"error": "missing user_id"}), 400
    # Flask 是同步框架,用 run() 桥接到事件循环
    result = asyncio.run(get_overview_async(user_id))
    return jsonify(result)

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

关键参数说明:
- TCPConnector(limit=100, limit_per_host=20):连接池上限100个,单host并发20。这是压测调出来的值,太小会排队,太大会触发上游连接拒绝。
- asyncio.timeout() 替代老式 asyncio.wait_for():前者支持协程内嵌套超时,后者会取消整个协程导致难以局部处理。
- return_exceptions=True:确保一个服务挂掉不影响另外三个。

5. 踩坑与优化:三个血泪教训

坑1:Flask与asyncio的兼容性问题

直接asyncio.run()在Flask视图函数里跑,生产环境没问题,但如果你用Flask的测试客户端(app.test_client())就会报RuntimeError: asyncio.run() cannot be called from a running event loop。解决:测试时改用pytest-asyncio插件,或者单独写异步测试用例。生产环境Flask跑在Gunicorn + gevent worker下,每个worker是独立进程,不存在嵌套事件循环问题。

坑2:ClientSession初始化时机

我一开始在get_overview_async里每次创建ClientSession,压测发现吞吐量只有200 req/s —— 因为ClientSession初始化要建立TCP连接握手,开销巨大。改成模块级懒加载后,首次请求创建,后续复用,吞吐量直接翻倍。记住:ClientSession是连接池,不是请求对象,必须全局单例

坑3:超时设置不是越小越好

最初把所有服务超时都设为1秒,结果audit-service在高峰期响应1.2秒,导致大量{"error": "timeout"}返回。后来根据上游服务的P99延迟分别设置:user-service 2s,order-service 2.5s,points-service 1.5s,audit-service 1s。超时应该比上游P99略大,而不是统一瞎设

优化1:异步日志采集

logging模块的QueueHandler + 后台线程消费日志,避免日志IO阻塞事件循环。日志量从200行/秒降到不影响性能的5%开销。

优化2:冷启动预热

在Gunicorn启动钩子里,预先创建ClientSession并发送一个健康检查请求,避免第一个真实请求承担连接池初始化开销。

6. 效果数据:同一台机器,同样的wrk压测

指标 同步版本 asyncio版本 提升幅度
平均延迟 1870ms 480ms -74.3%
P95延迟 2100ms 650ms -69.0%
吞吐量 120 req/s 506 req/s +321.7%
CPU使用率 85% 32% -62.4%
内存占用 1.8GB(线程栈) 450MB(协程栈) -75.0%

压测命令:wrk -t8 -c200 -d60s http://127.0.0.1:5000/api/user/overview?user_id=12345

为什么吞吐量能到500+? 因为4个上游请求并发后,单次请求耗时从1.8s降到480ms,相当于原来一个请求的时间窗口内能处理3.7个请求;同时CPU占用率下降,允许更多并发请求在事件循环中轮询。实际压测到800并发时,asyncio版吞吐量开始持平(受限于上游服务能力),但CPU依然只有40%。

7. 总结与踩坑清单

这次重构让我深刻理解了一件事:Python的asyncio不是银弹,但对付IO密集型的Web聚合接口,它是性价比最高的方案。核心收益不是“快”,而是把CPU从无意义的等待中解放出来。

给后来者的建议:

  1. 先用cProfile确认瓶颈是IO等待还是CPU计算。如果是CPU密集型,asyncio没用,该用multiprocessing。
  2. 信号量必须加。不加限流的异步代码,在流量突增时会把上游打挂。
  3. 每个上游独立配置超时。别用全局统一超时,不同服务的响应特性差异很大。
  4. 永远用return_exceptions=True。否则一个服务抖动,整个接口返回500。
  5. 连接池参数要压测调试limitlimit_per_host不是越大越好,取决于上游的accept队列长度。

未来如果接口数量继续膨胀,我会考虑把这段异步逻辑抽离成独立服务,用gRPC或消息队列解耦 —— 但至少现在,asyncio帮我们多撑了3个月的业务增长,没加一台机器。