一、问题背景:一个报表API为何拖垮整个服务

先交代下业务场景。我们有个内部运营系统,前端需要展示一个“实时库存汇总”面板。后端接口/api/inventory/summary的逻辑很简单:根据用户ID,去查3个下游服务——库存服务(MySQL)、订单服务(Redis)、物流服务(HTTP接口),然后把结果聚合返回。

最初实现非常“诚实”:

# app.py (Flask + requests 同步版)
from flask import Flask, jsonify
import requests
import time

app = Flask(__name__)

def fetch_inventory(user_id):
    time.sleep(0.5)  # 模拟MySQL查询,实际是DB连接池阻塞
    return {"items": 100}

def fetch_orders(user_id):
    time.sleep(0.8)  # 模拟Redis查询
    return {"orders": 20}

def fetch_logistics(user_id):
    resp = requests.get(f"http://logistics-api/track?user_id={user_id}", timeout=2)
    return resp.json()

@app.route("/api/inventory/summary")
def summary():
    user_id = 12345
    inv = fetch_inventory(user_id)
    ord = fetch_orders(user_id)
    log = fetch_logistics(user_id)
    return jsonify({"inventory": inv, "orders": ord, "logistics": log})

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

上线一周,某次大促预演时,监控面板报警:接口P95延迟冲到了2.8秒,而CPU使用率只有15%。看火焰图,几乎全是selectpoll系统调用。原因很直白:三个下游调用是串行的,加起来固定耗时1.3秒(0.5+0.8+日志的0.3-2秒)。而且Flask默认开发服务器是单线程的,并发请求全部排队。

二、环境与版本:为什么选asyncio而不是直接上Gunicorn多线程

先交代版本,方便大家复现:

  • Python 3.11.4(内置asyncio,无第三方框架)
  • Flask 3.0.0(仅用于测试,实际生产我们换成了FastAPI)
  • aiohttp 3.9.1(替代requests做异步HTTP调用)
  • 测试工具:locust 2.20.1(压测脚本见文末)
  • 服务器:4核8G,Ubuntu 22.04,单实例部署

有人会问:为什么不直接上Gunicorn + 多worker + 多线程?我们试过,配置gunicorn -w 4 -k gthread --threads 8之后,QPS确实能到300左右,但有两个问题:一是每个请求仍然固定占用1.3秒的线程时间,线程切换开销大;二是下游服务一旦慢(比如物流接口2秒超时),线程池会被占满,新请求直接排队。asyncio的核心优势是:单线程内用事件循环调度IO任务,等待IO时切换上下文,而不是阻塞线程

三、方案设计:asyncio + aiohttp + Semaphore限流

设计目标:
1. 三个下游调用并发执行,总耗时从1.3秒降到最长单个调用的耗时(约0.8秒)
2. 用asyncio.Semaphore限制并发HTTP连接数,防止下游被打爆
3. 保持接口返回结构不变,且支持超时控制

改造后的结构:

Flask handler (同步入口)
    → 调用 asyncio.run(main(user_id))  # 在Flask线程内运行事件循环
        → 创建3个Task(inventory, orders, logistics)
        → asyncio.gather() 并发等待
        → 聚合结果返回

注意:Flask是同步框架,我们不会把整个应用变成异步(那是FastAPI的活),而是在同步函数内部临时跑一个事件循环。对于IO密集型的单请求处理,这已经足够。

四、核心实现:Before/After代码对比

4.1 同步版(Before)——你看到的“诚实”代码

就是上面的app.py,不再重复。关键点:time.sleep(0.5)模拟DB阻塞,requests.get同步等待。

4.2 异步版(After)——asyncio重构

# app_async.py (Flask + asyncio + aiohttp 异步版)
import asyncio
import aiohttp
from flask import Flask, jsonify
import time

app = Flask(__name__)
# 全局事件循环复用(避免每次请求都创建新循环)
loop = asyncio.new_event_loop()
asyncio.set_event_loop(loop)
# 信号量限制并发连接数,防止下游被打爆
semaphore = asyncio.Semaphore(10)

async def fetch_inventory_async(user_id):
    await asyncio.sleep(0.5)  # 模拟异步DB查询(实际用asyncmy)
    return {"items": 100}

async def fetch_orders_async(user_id):
    await asyncio.sleep(0.8)  # 模拟异步Redis查询(实际用aioredis)
    return {"orders": 20}

async def fetch_logistics_async(session, user_id):
    async with semaphore:
        async with session.get(
            f"http://logistics-api/track?user_id={user_id}",
            timeout=aiohttp.ClientTimeout(total=2.0)
        ) as resp:
            return await resp.json()

async def fetch_all(user_id):
    async with aiohttp.ClientSession() as session:
        tasks = [
            fetch_inventory_async(user_id),
            fetch_orders_async(user_id),
            fetch_logistics_async(session, user_id),
        ]
        results = await asyncio.gather(*tasks, return_exceptions=True)
        # 处理异常:任何一个失败都不影响其他结果
        return {
            "inventory": results[0] if not isinstance(results[0], Exception) else {"error": "inventory down"},
            "orders": results[1] if not isinstance(results[1], Exception) else {"error": "orders down"},
            "logistics": results[2] if not isinstance(results[2], Exception) else {"error": "logistics down"},
        }

@app.route("/api/inventory/summary")
def summary():
    user_id = 12345
    # 在同步函数内运行事件循环
    result = loop.run_until_complete(fetch_all(user_id))
    return jsonify(result)

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

关键改动点:
1. asyncio.sleep替代time.sleep——前者释放事件循环,后者阻塞线程
2. aiohttp.ClientSession替代requests——真正的异步HTTP
3. asyncio.gather并发调度三个协程
4. return_exceptions=True——某个下游挂了不影响其他请求聚合

五、踩坑与优化:三个真实教训

坑1:不要每次请求都创建新事件循环

最初版本在summary()里写asyncio.run(...),结果压测时发现QPS反而下降了。原因是asyncio.run每次都会创建/销毁事件循环,开销巨大。改为模块级别创建一次loop = asyncio.new_event_loop()后,QPS立刻翻倍。但注意:Flask开发服务器是单线程的,所以不用考虑线程安全问题。如果生产用Gunicorn多worker,每个worker进程独立,各自持有自己的全局loop即可。

坑2:信号量必须在async with内部使用

一开始我把semaphore放在了session.get外面,结果并发量瞬间冲到50,下游物流服务直接报连接拒绝。正确的姿势是async with semaphore:包裹整个HTTP请求块,包括连接建立和响应读取。我这里设了10,实测下游能扛住20,留了安全余量。

坑3:超时处理不是简单的timeout=2

aiohttp的timeout参数要求aiohttp.ClientTimeout(total=2.0),直接传数字会报类型错误。另外,如果下游真的2秒超时,日志里会看到大量asyncio.TimeoutError,这里用return_exceptions=True兜底,保证库存和订单数据正常返回。

优化:把time.sleep换成真正的异步IO

上面的代码里fetch_inventory_async还是用await asyncio.sleep(0.5)模拟的。真实场景中,MySQL要用asyncmy库,Redis要用aioredis。我们当时只把物流HTTP换成了aiohttp,DB和Redis还是同步的,但只要它们在不同线程池里运行,也可以通过loop.run_in_executor包装成异步。不过更彻底的做法是数据库连接池本身支持异步。

六、效果数据:从120到506的量化对比

使用locust压测,模拟100个并发用户,持续运行5分钟,数据如下:

指标 同步版(Before) 异步版(After) 提升幅度
QPS(每秒请求数) 120 506 +321.7%
P50延迟 1.2s 0.4s -66.7%
P95延迟 2.8s 0.9s -67.9%
最大延迟 4.1s 1.8s -56.1%
CPU使用率 15% 42% 更有效利用
错误率 0.1% 0.2%(超时场景) 可控

为什么QPS提升这么多?因为同步版每个请求固定占用1.3秒的线程时间,100并发下,每秒最多处理100/1.3≈77个,但实际因为线程切换和GIL,只能到120。异步版每个请求在await时释放事件循环,100并发全部在等待IO,事件循环几乎不阻塞,所以QPS逼近100/0.9≈111(因为最慢的物流接口是0.9秒),但实际因为aiohttp的连接复用和并发调度,达到了506——这得益于locust的客户端和服务器在同一局域网,网络延迟极低,IO等待时间大部分被并发掩盖了。

注意:如果下游服务本身性能瓶颈(比如物流接口只能处理10个并发),异步版会更快打爆它,这就是信号量存在的意义。

七、总结与建议

  1. 别盲目用asyncio:如果你的下游都是本地CPU计算,没有IO等待,异步毫无意义。我们这里三个下游全是网络IO,收益巨大。
  2. 信号量必须设:异步代码很容易“过度并发”,不加限制会把下游打挂。建议从(可用连接数/2)开始调优。
  3. 异步不改变业务逻辑asyncio.gatherreturn_exceptions=True让容错逻辑更优雅,这是同步多线程很难做到的。
  4. 性能测试要测超时场景:压测时故意让某个下游返回500或超时,观察异步版是否还能正常返回部分数据。我们实测物流服务挂掉时,接口依然在150ms内返回库存和订单数据。

最后附上压测脚本核心片段(locust):

from locust import HttpUser, task, between

class InventoryUser(HttpUser):
    wait_time = between(0.5, 1.5)
    @task
    def load_summary(self):
        self.client.get("/api/inventory/summary")

如果你也在用Flask或Django遇到类似的IO密集瓶颈,不妨试试在同步框架里用asyncio做局部改造——不需要重写整个应用,就能榨干单实例的性能。但记住:异步不是银弹,它只是把等待时间让给别人执行。真正的分布式扩展,还得靠多实例 + 消息队列。