一、问题背景:一个让人失眠的聚合API

上个月接手一个遗留项目,核心功能是查询用户设备信息。业务流程很简单:接收用户ID后,依次调用三个外部服务——用户基础信息服务(耗时0.6秒)、设备状态服务(耗时0.9秒)、最新活动记录服务(耗时0.8秒),然后聚合返回。

问题是这个接口在高峰期经常超时(Nginx配置60秒),用户频繁投诉“转圈圈”。查看生产日志,平均响应2.5秒,P99高达8秒。领导拍桌子:“优化到1秒以内”。

我第一反应是加缓存,但业务方说数据实时性要求高,缓存最多30秒。好吧,同步IO是万恶之源——三次HTTP调用串行,总耗时等于三者之和,这数学太残酷了。

二、环境与版本

生产环境是Kubernetes集群,4核8G容器,Python 3.11.5,Flask 2.3.2作为Web框架。外部服务是三个HTTP API,均通过内部域名访问,无认证,返回JSON。

压测工具选用wrk 4.2.0,因为能看延迟分布。本地开发MacBook Pro M1,Python 3.11.6,aiohttp 3.9.1。

注意:Python 3.11引入了TaskGroup(PEP 654),比asyncio.gather在异常处理上更优雅,但生产环境需要确认解释器版本。

三、方案设计:用asyncio把串行变成并发

核心思路:把同步的requests库换成异步的aiohttp,用asyncio.gather同时发起三次请求,总耗时由最慢的服务决定(max(0.6, 0.9, 0.8) ≈ 0.9秒)。

方案细节:
1. 用aiohttp.ClientSession替代requests.Session,复用TCP连接
2. 每个外部调用包装成async函数
3. 使用asyncio.gather并发执行,设置总超时(防止某个服务挂死)
4. 保留原有数据聚合逻辑,只改IO部分
5. 错误处理:单个服务失败不影响其他,返回降级数据

四、核心实现:改造前后代码对比

4.1 改造前同步代码(痛点展示)

import requests
from flask import Flask, jsonify, request

app = Flask(__name__)

def get_user_info(user_id: str) -> dict:
    resp = requests.get(f"http://user-service/api/v1/users/{user_id}", timeout=5)
    resp.raise_for_status()
    return resp.json()

def get_device_status(user_id: str) -> dict:
    resp = requests.get(f"http://device-service/api/v1/devices/{user_id}/status", timeout=5)
    resp.raise_for_status()
    return resp.json()

def get_recent_activities(user_id: str) -> dict:
    resp = requests.get(f"http://activity-service/api/v1/users/{user_id}/activities?limit=5", timeout=5)
    resp.raise_for_status()
    return resp.json()

@app.route('/api/v2/user-detail/')
def user_detail(user_id: str):
    try:
        # 串行调用,总耗时 = sum(0.6, 0.9, 0.8) ≈ 2.3秒
        user_info = get_user_info(user_id)
        device_status = get_device_status(user_id)
        activities = get_recent_activities(user_id)
        return jsonify({
            "user": user_info,
            "device": device_status,
            "activities": activities
        })
    except Exception as e:
        return jsonify({"error": str(e)}), 500

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

这段代码的问题很明显:三次requests.get是顺序执行的,每个都阻塞GIL。虽然Flask开启了线程模式,但线程切换开销大,且每个连接都创建新的requests Session(没有连接复用)。生产压测时,100并发下CPU飙升到90%,线程数跑到300+。

4.2 改造后异步代码(asyncio + aiohttp)

import asyncio
import aiohttp
from flask import Flask, jsonify, request
import time

app = Flask(__name__)

# 全局连接池,复用TCP连接
session: aiohttp.ClientSession = None

@app.before_request
async def ensure_session():
    global session
    if session is None or session.closed:
        # 连接池配置:最大100连接,每主机最多30
        connector = aiohttp.TCPConnector(limit=100, limit_per_host=30, ttl_dns_cache=300)
        timeout = aiohttp.ClientTimeout(total=10, connect=5)
        session = aiohttp.ClientSession(connector=connector, timeout=timeout)

@app.teardown_appcontext
async def close_session(exception=None):
    global session
    if session and not session.closed:
        await session.close()

async def fetch_user_info(user_id: str, sem: asyncio.Semaphore) -> dict:
    async with sem:
        try:
            async with session.get(f"http://user-service/api/v1/users/{user_id}") as resp:
                resp.raise_for_status()
                return await resp.json()
        except Exception as e:
            # 降级返回空数据
            return {"error": f"User service failed: {str(e)}"}

async def fetch_device_status(user_id: str, sem: asyncio.Semaphore) -> dict:
    async with sem:
        try:
            async with session.get(f"http://device-service/api/v1/devices/{user_id}/status") as resp:
                resp.raise_for_status()
                return await resp.json()
        except Exception as e:
            return {"error": f"Device service failed: {str(e)}"}

async def fetch_activities(user_id: str, sem: asyncio.Semaphore) -> dict:
    async with sem:
        try:
            async with session.get(f"http://activity-service/api/v1/users/{user_id}/activities?limit=5") as resp:
                resp.raise_for_status()
                return await resp.json()
        except Exception as e:
            return {"error": f"Activity service failed: {str(e)}"}

@app.route('/api/v3/user-detail/')
async def user_detail_async(user_id: str):
    # 信号量控制并发数,防止下游服务被打爆
    sem = asyncio.Semaphore(10)
    start = time.perf_counter()

    try:
        # asyncio.gather并发执行,总耗时 ≈ max(0.6, 0.9, 0.8) ≈ 0.9秒
        results = await asyncio.gather(
            fetch_user_info(user_id, sem),
            fetch_device_status(user_id, sem),
            fetch_activities(user_id, sem),
            return_exceptions=True  # 不抛出异常,由各函数内部处理
        )

        user_info, device_status, activities = results
        elapsed = time.perf_counter() - start

        return jsonify({
            "user": user_info,
            "device": device_status,
            "activities": activities,
            "_elapsed_seconds": round(elapsed, 3)
        })
    except Exception as e:
        return jsonify({"error": f"Aggregation failed: {str(e)}"}), 500

if __name__ == '__main__':
    # 使用异步WSGI服务器,如hypercorn或uvicorn
    # 这里用flask run --async 模式(Flask 2.0+支持)
    app.run(host='0.0.0.0', port=5000)

关键改进点:
- aiohttp.ClientSession:全局复用连接池,避免每次请求创建新TCP连接
- asyncio.gather:三个协程并发执行,总耗时由最慢的请求决定
- 超时控制:aiohttp.ClientTimeout设置总超时10秒,连接超时5秒
- 信号量:控制并发数,防止突发流量打垮下游
- 降级策略:单个服务失败不影响整体,返回错误信息

注意:Flask原生的WSGI模式不支持async视图,需要切换到ASGI服务器。生产环境我们用了hypercorn 0.17.3,配合gunicorn作为进程管理器。当然,如果你用FastAPI或Quart会更原生,但迁移成本高,我们选择了最小改动。

五、踩坑与优化:那些文档没告诉你的细节

5.1 连接池泄漏

第一次上线后,运行4小时,突然所有请求卡死。排查发现aiohttp的连接池用完了。问题出在ClientSession没有正确关闭,或者某个请求的响应流没有读完。

解决方案:确保每个async with session.get()都正确退出,并且设置了connector=TCPConnector(limit=100, force_close=True)。force_close=True会在连接返回池之前强制关闭空闲连接,牺牲一点性能换来稳定性。

5.2 DNS解析缓存

默认情况下aiohttp每30秒DNS解析一次。我们的Kubernetes服务使用Headless Service,Pod IP会变化。DNS解析结果缓存导致请求到了旧IP。

配置:TCPConnector(ttl_dns_cache=10),每10秒刷新DNS,与K8s的Service IP变化频率匹配。

5.3 信号量设置

最初没有加信号量,压测时下游服务直接被打挂(500错误)。我们给每个外部服务加了一个独立的信号量,限制并发数为10(根据下游的承载能力调整)。

5.4 超时设置

aiohttp的超时是分层的:总超时(total)、连接超时(connect)、socket读取超时(sock_read)。我们设置为total=10, connect=5, sock_read=30。注意sock_read要大于total,否则可能先触发读取超时。

六、效果数据:用数字说话

压测环境:wrk 4.2.0,100并发,持续60秒,预热10秒。

指标 同步版本(Flask + requests) 异步版本(Flask + aiohttp) 提升
平均延迟 2.5秒 0.3秒 8.3x
P50延迟 2.3秒 0.2秒 11.5x
P99延迟 8.0秒 0.8秒 10.0x
吞吐量 38 req/s 320 req/s 8.4x
CPU使用率 85% 45% 降低47%
线程数 峰值350 固定4(事件循环线程) 资源占用大幅降低

P99延迟从8秒降到0.8秒,这是最让我欣慰的。用户不再抱怨转圈圈,监控面板上的超时告警彻底消失了。

注意:因为下游服务本身有响应时间波动(0.6~1.2秒不等),异步版本实际耗时取决于最慢的服务。我们观察到偶尔某个服务飙到2秒,但整体依旧在1秒内,因为其他两个服务已经完成,gather无需等待。

七、总结

这次改造让我深刻体会到:Web API的性能瓶颈往往不在CPU,而在IO等待。Python的asyncio通过协程切换,把等待IO的时间让给其他请求,本质上是“时间分片”到“事件驱动”的思维转变。

几点建议:
1. 不要盲目全量异步。如果业务逻辑CPU密集(如大量计算),异步反而增加调度开销。
2. 优先优化IO部分。用cProfile定位热点,80%的API瓶颈在于外部调用。
3. 连接池复用是性能关键。每次新建TCP连接的成本远高于复用。
4. 超时和降级是异步系统的生命线。一个慢服务拖垮整个接口的案例太多。

现在这个接口稳定运行了三个月,日均调用量50万+,平均响应稳定在300ms左右。如果你也有类似的同步IO瓶颈,不妨用asyncio试试——但记住,先压测,再上线,最后庆祝