1. 问题背景:一个慢到被投诉的报表接口

上个月运维同事甩给我一个监控截图:“这个/api/v1/report/summary接口,P99延迟2.3秒,上游服务都超时了”。我查了下代码,这是个典型的同步串行调用

# 改造前:同步串行调用,3个HTTP请求顺序执行
def get_summary_report(self, user_id):
    # 调用订单服务
    order_data = requests.get(f"http://order-service/api/orders?uid={user_id}", timeout=3).json()
    # 调用用户服务
    user_data = requests.get(f"http://user-service/api/users?uid={user_id}", timeout=3).json()
    # 调用推荐服务
    rec_data = requests.get(f"http://rec-service/api/recommend?uid={user_id}", timeout=3).json()
    # 业务聚合...
    return self._merge(order_data, user_data, rec_data)

三个服务调用是独立的,但这里串行执行,每次请求固定耗时 = 三个服务延迟之和。当时平均每个服务响应200-300ms,总耗时850ms左右。

2. 环境与版本:明确基线数据

先说下环境,方便你复现对比:
- Python 3.8.10(生产环境,后面会提到这个版本的坑)
- Flask 2.0.1 + gunicorn 20.1.0(4 workers, sync worker class)
- 上游服务:3个内部FastAPI服务,部署在K8s内网
- 压测工具:wrk 4.1.0,单线程100连接压测10秒

压测基线(改造前):

Running 10s test @ http://api-gateway:8080/api/v1/report/summary
  100 threads and 100 connections
  Thread Stats   Avg      Stdev     Max   +/- Stdev
    Latency     1.21s   320.5ms   2.30s    85%
    Req/Sec     120.5     15.3    180.0    75%
  12208 requests in 10.09s, 3.12MB read
Requests/sec:    1210.0

这里注意:单worker QPS约120,多worker后整体吞吐上不去,瓶颈在CPU等待I/O(线程阻塞在requests上)。

3. 方案设计:异步化不是终极方案,要配合连接池

我的思路很直接:用asyncio把3个串行HTTP调用改成并发,同时用aiohttp替代requests。但有几个关键决策点:

  1. 框架选择:Flask是同步WSGI框架,不能直接用async def路由。所以采用双层架构:Flask路由保持同步,内部调用asyncio.run()启动事件循环,或者用loop.run_in_executor把异步逻辑包装成线程执行。
  2. 连接池复用aiohttp.ClientSession必须全局复用,否则每次请求创建新连接池,性能反而更差。
  3. 超时控制:用asyncio.wait_for给每个子任务加超时,避免单个服务慢拖垮整体。
  4. 限流asyncio.Semaphore控制并发连接数,防止打爆上游服务。

架构图(文字版):

Flask (sync) → run_in_executor → asyncio.run()
    → 创建aiohttp.ClientSession(connector=TCPConnector(limit=100))
    → asyncio.gather(3个coroutine, return_exceptions=True)
    → 聚合结果返回

4. 核心实现:asyncio重构代码

先说结论:不要直接用asyncio.run()在Flask路由里,因为每次调用会创建新事件循环,开销巨大。正确做法是全局事件循环 + run_coroutine_threadsafe,或者更简单点:用loop.run_in_executor把整个异步逻辑丢到默认线程池。

4.1 全局Session与辅助函数

# async_client.py
import asyncio
import aiohttp
from aiohttp import TCPConnector

# 全局连接池,limit=100表示最多100个并发连接
_connector = TCPConnector(limit=100, limit_per_host=20, ttl_dns_cache=300)
_session = aiohttp.ClientSession(connector=_connector)

async def fetch_json(session, url, semaphore, timeout=2.0):
    """带限流和超时的异步请求"""
    async with semaphore:
        try:
            async with session.get(url, timeout=aiohttp.ClientTimeout(total=timeout)) as resp:
                if resp.status != 200:
                    return None
                return await resp.json()
        except asyncio.TimeoutError:
            # 记录日志,返回降级数据
            return None
        except Exception as e:
            return None

async def fetch_all(user_id):
    sem = asyncio.Semaphore(50)  # 限制并发50
    urls = [
        f"http://order-service/api/orders?uid={user_id}",
        f"http://user-service/api/users?uid={user_id}",
        f"http://rec-service/api/recommend?uid={user_id}",
    ]
    results = await asyncio.gather(
        *(fetch_json(_session, url, sem) for url in urls),
        return_exceptions=True  # 防止单个异常导致全部失败
    )
    return results  # [order_data, user_data, rec_data] 或 None

4.2 Flask路由集成

# app.py
import asyncio
from flask import Flask, jsonify
from async_client import fetch_all

app = Flask(__name__)
# 全局事件循环,在worker启动时创建
loop = asyncio.new_event_loop()

@app.route('/api/v1/report/summary')
def get_summary():
    user_id = request.args.get('uid')
    if not user_id:
        return jsonify({'error': 'missing uid'}), 400

    # run_coroutine_threadsafe 是线程安全的
    future = asyncio.run_coroutine_threadsafe(fetch_all(user_id), loop)
    try:
        results = future.result(timeout=5)  # 整体超时5s
    except asyncio.TimeoutError:
        return jsonify({'error': 'upstream timeout'}), 504

    # 业务聚合(略)
    return jsonify(merge_data(results))

关键点loop必须在gunicorn worker启动时创建,不能每次请求都建。用run_coroutine_threadsafe把协程提交到事件循环,主线程阻塞等待结果(future.result())。

4.3 gunicorn配置调整

因为异步逻辑跑在事件循环里,worker不需要多线程了,把worker数调成CPU核数×2,worker_class保持同步:

gunicorn -w 8 -k sync -b 0.0.0.0:8080 app:app

5. 踩坑与优化:三个大坑,每个都让我调了半天

5.1 Python 3.8的asyncio.run()坑

一开始我图省事,直接在Flask路由里写asyncio.run(fetch_all(user_id))。结果压测时发现每隔几分钟就报“Event loop is closed”。查了文档:Python 3.8的asyncio.run()每次创建新循环,且ClientSession绑定的是旧循环,下次调用时session关联的传输已经关闭。

解决方案:改用全局事件循环 + run_coroutine_threadsafe。升级Python 3.10后asyncio.run()支持复用,但生产环境一时半会儿升不了,就用全局循环方案。

5.2 aiohttp连接池耗尽

压测到500 QPS时,上游服务开始报“Connection reset by peer”。原因是TCPConnector(limit=100)默认是对所有域名共享100连接,但我们有3个上游服务,每个服务被限制到33个连接,不够用。

优化:改成limit=100, limit_per_host=50(每域名50),同时加上enable_cleanup_closed=True自动清理半开连接。实测后上游重试率从5%降到0.3%。

5.3 Semaphore的坑:放错位置

一开始Semaphore放在fetch_all()里,每次请求都创建新的信号量,等于没限流。应该全局创建,在fetch_json内部使用。上面代码里sem = asyncio.Semaphore(50)fetch_all里创建,其实不对,应该提到模块级别:

_sem = asyncio.Semaphore(50)  # 全局信号量

async def fetch_all(user_id):
    urls = [...]
    results = await asyncio.gather(
        *(fetch_json(_session, url, _sem) for url in urls),  # 用全局_sem
        return_exceptions=True
    )

别小看这个,一开始每个请求新建Semaphore,并发数完全不可控,直接打崩了上游。

6. 效果数据:改造前后对比

压测命令和参数完全一致(wrk -t100 -c100 -d10s),改造后结果:

Running 10s test @ http://api-gateway:8080/api/v1/report/summary
  100 threads and 100 connections
  Thread Stats   Avg      Stdev     Max   +/- Stdev
    Latency   180.2ms    45.1ms   500.0ms   82%
    Req/Sec    1800.5    210.3   2400.0    73%
  18205 requests in 10.09s, 4.65MB read
Requests/sec:    1804.0
指标 改造前 改造后 提升
平均延迟 1.21s 180ms 6.7倍
P99延迟 2.3s 450ms 5.1倍
QPS 120 1800 15倍
错误率 1.2%(超时) 0.1%(降级) 12倍

注意:QPS提升15倍不全是asyncio的功劳,还因为我把worker数从4调到8(用满CPU核)。但老代码即使8个worker,QPS也就200左右,因为每个worker只能处理1个请求,阻塞在I/O上。异步化后单worker就能同时处理几十个请求。

7. 总结与建议

这次改造的核心教训:

  1. 同步代码的并发瓶颈在I/O等待,asyncio把等待时间让给其他协程,吞吐量提升是数量级的。
  2. 别用asyncio.run()在同步框架里,全局事件循环才是正道。
  3. 连接池和信号量是异步编程的命门,不控制迟早出事。
  4. return_exceptions=True必须加,否则一个服务抖动就全挂。

如果你的接口是多个独立HTTP调用聚合,强烈建议试试这个方案。成本不高,收益巨大。但注意:如果业务逻辑里有CPU密集计算(比如大列表排序、正则匹配),别直接用asyncio,配合run_in_executor丢给进程池,否则事件循环会被阻塞。

最后的建议:生产环境升级Python到3.10+,asyncio.run()复用循环的体验好很多,还有官方asyncio.timeout()上下文管理器,比wait_for更优雅。