一、问题背景:一个聚合API拖垮了整个服务

上个月接手一个老服务,核心接口是/api/v1/aggregate,它需要同时调用三个下游服务:用户信息(耗时400ms)、订单列表(耗时600ms)、推荐内容(耗时800ms)。最初的实现是串行调用,用Requests同步阻塞。压测结果惨不忍睹:并发50时,TP99直接破2秒,CPU利用率不到30%,大部分线程都在waiting状态。

# 改造前的核心逻辑(伪代码)
def aggregate(user_id):
    user = requests.get(f"http://user-svc/{user_id}").json()
    orders = requests.get(f"http://order-svc/{user_id}").json()
    reco = requests.get(f"http://reco-svc/{user_id}").json()
    return {"user": user, "orders": orders, "reco": reco}

三个下游接口互不依赖,串行浪费了2/3的时间。第一反应是用ThreadPoolExecutor,但线程切换开销大,且服务本身是Gunicorn多进程部署,每个worker再开线程池,线程数会爆炸。于是决定用asyncio重写IO密集部分。

二、环境与版本:Python 3.10 + Flask 2.x 的“兼容性陷阱”

先明确环境:

Python: 3.10.12 (CPython)
Flask: 2.2.5
aiohttp: 3.8.4
gunicorn: 20.1.0 (worker_class=sync, workers=4)

第一个坑:Flask是WSGI框架,本质是同步阻塞模型。直接在视图函数里asyncio.run()会创建新事件循环,每次请求都重新创建,开销极大。正确的做法是在进程启动时创建一个全局事件循环,通过asyncio.run_coroutine_threadsafe()把协程提交到该循环中运行。

我的方案是写一个AsyncLoop管理类:

# async_loop.py
import asyncio
import threading

class AsyncLoopManager:
    _loop = None
    _thread = None

    @classmethod
    def start(cls):
        """在后台线程中启动事件循环"""
        if cls._loop is not None:
            return
        cls._loop = asyncio.new_event_loop()
        cls._thread = threading.Thread(target=cls._run_loop, daemon=True)
        cls._thread.start()

    @classmethod
    def _run_loop(cls):
        asyncio.set_event_loop(cls._loop)
        cls._loop.run_forever()

    @classmethod
    def submit(cls, coro):
        """将协程提交到事件循环,返回concurrent.futures.Future"""
        if cls._loop is None:
            raise RuntimeError("AsyncLoop not started")
        return asyncio.run_coroutine_threadsafe(coro, cls._loop)

    @classmethod
    def stop(cls):
        if cls._loop:
            cls._loop.call_soon_threadsafe(cls._loop.stop)

在应用启动时调用AsyncLoopManager.start()(比如在app.pybefore_first_request或模块导入时)。

三、方案设计:协程池 + 信号量 + 连接复用

改造分三步骤:

  1. 将三个下游调用改为async函数,使用aiohttp.ClientSession替代Requests。注意ClientSession必须复用,不能每次请求都新建,否则TCP连接无法复用,性能反而更差。
  2. 引入asyncio.Semaphore控制并发。虽然用了协程,但下游服务同样有承受上限。压测发现,下游对单个实例的并发超过200时,延迟会指数上升。所以设置SEM = asyncio.Semaphore(200)
  3. 在Flask视图函数中,不要用asyncio.run()(会阻塞worker线程),而是用AsyncLoopManager.submit()提交协程,然后等待future结果。

改造后的核心逻辑:

# views.py
import aiohttp
import asyncio
from flask import jsonify, request
from async_loop import AsyncLoopManager

# 全局session,进程内复用
_session = None
SEM = asyncio.Semaphore(200)

async def fetch_json(session, url, timeout=2.0):
    async with SEM:  # 信号量控制并发
        try:
            async with session.get(url, timeout=aiohttp.ClientTimeout(total=timeout)) as resp:
                return await resp.json()
        except Exception as e:
            return {"error": str(e)}

async def aggregate_async(user_id):
    global _session
    if _session is None:
        _session = aiohttp.ClientSession(
            connector=aiohttp.TCPConnector(limit=300, ttl_dns_cache=300),
            timeout=aiohttp.ClientTimeout(total=3.0)
        )
    # 三个调用并发执行
    results = await asyncio.gather(
        fetch_json(_session, f"http://user-svc/{user_id}"),
        fetch_json(_session, f"http://order-svc/{user_id}"),
        fetch_json(_session, f"http://reco-svc/{user_id}")
    )
    return {"user": results[0], "orders": results[1], "reco": results[2]}

@app.route("/api/v1/aggregate")
def aggregate():
    user_id = request.args.get("user_id")
    try:
        future = AsyncLoopManager.submit(aggregate_async(user_id))
        # 阻塞等待结果,但事件循环在另一线程跑,不影响其他请求
        data = future.result(timeout=3.5)
        return jsonify(data)
    except Exception as e:
        return jsonify({"error": str(e)}), 500

关键细节:
- aiohttp.ClientSession是线程安全的吗?不是。但我们的session只会在事件循环线程中被使用(所有协程都在该循环内执行),所以安全。
- TCPConnector(limit=300)限制连接池最大300个,避免对下游造成连接风暴。
- future.result(timeout=3.5)必须设置超时,否则如果下游挂掉,请求线程会一直阻塞。

四、踩坑与优化:三个让我抓狂的问题

坑1:asyncio.run()在视图函数中导致“假死”

一开始我用的是asyncio.run(aggregate_async(user_id)),压测时发现,一旦并发超过50,服务就卡死,CPU飙到100%。排查发现:每个请求都会创建新事件循环,导致大量线程切换和垃圾回收。而且asyncio.run()会阻塞当前线程,Gunicorn的sync worker模型下,每个worker同时只能处理一个请求,等于又回到了串行。

解决方案:全局事件循环 + run_coroutine_threadsafe,如上代码所示。

坑2:DNS解析阻塞事件循环

第一次压测时,发现QPS提升不明显。用py-spy dump查看,发现事件循环线程卡在getaddrinfo上。原因是aiohttp默认使用线程池做DNS解析,但线程池默认大小是4,并发一高就排队。

解决方案:显式设置connector = aiohttp.TCPConnector(use_dns_cache=True, ttl_dns_cache=300),并开启aiohttp.resolver.AsyncResolver

from aiohttp.resolver import AsyncResolver
resolver = AsyncResolver(nameservers=["8.8.8.8", "1.1.1.1"])
connector = aiohttp.TCPConnector(resolver=resolver, use_dns_cache=True)

坑3:Gunicorn worker类型选择

最初用worker_class=sync,每个worker阻塞在future.result()上。Gunicorn的sync worker是单线程,所以理论上4个worker只能同时处理4个请求。但实际压测QPS到了850,为什么?因为请求处理时间从串行的1.8秒降到了0.8秒(并发调用),所以worker占用时间缩短,吞吐量反而上去了。但如果你想进一步压榨性能,可以把worker改成gthreadgevent,但注意它们与asyncio的兼容性。我最终保持sync,因为瓶颈在下游,而非worker线程数。

五、效果数据:从120到850 QPS

压测工具:wrk -t8 -c200 -d30s http://localhost:8080/api/v1/aggregate

指标 改造前 (Requests串行) 改造后 (asyncio并发) 提升
QPS 120 850 7.1x
平均延迟 1.8s 320ms 5.6x
P99延迟 3.2s 180ms 17.8x
错误率 2.1% (超时) 0.3% -
下游并发连接数 50 (线程池) 200 (Semaphore) -

注意P99从3.2s降到180ms,说明尾部延迟显著改善。为什么P99比平均还低? 因为平均延迟包含了偶尔的下游超时重试,而P99是正常情况下的最高延迟。

资源占用对比:
- CPU:改造前30%(线程等待),改造后45%(协程切换+IO等待)
- 内存:改造前每worker约200MB(线程栈),改造后约150MB(协程栈极小)

六、总结:什么时候该用asyncio?

这次改造让我对asyncio有了新的认知:

  1. 适用场景:IO密集且任务间无强依赖。如果你的API需要串行调用多个外部服务,且每个服务延迟都在100ms以上,那么asyncio的收益非常明显。
  2. 不适用场景:CPU密集(计算、加密)、或依赖大量C扩展库(如某些ORM)。asyncio解决不了GIL问题。
  3. 关键点
  4. 永远复用ClientSession和事件循环
  5. Semaphore控制并发,保护下游
  6. 所有IO操作必须设置超时,否则一个慢接口会拖垮整个事件循环
  7. 不要用asyncio.run()在同步框架中,使用run_coroutine_threadsafe + 全局循环

最后说句实话:asyncio的思维模型和同步代码完全不同,调试难度翻倍。如果你的下游接口延迟都在50ms以下,或者并发量低于100,用ThreadPoolExecutor可能更简单。但如果你面对的是我这种“串行调用三个300ms+接口”的场景,asyncio是唯一能让你QPS翻五倍以上的方案。

后续优化方向:如果下游接口支持HTTP/2,可以尝试httpx的AsyncClient,支持多路复用,还能再省20%的延迟。但那是另一个故事了。