一、问题背景:同步IO把CPU饿死了

先交代一下业务场景。我们有个内部工单系统,前端有一个"用户所有工单+关联操作日志"的聚合报表接口。原来是用Flask写的同步视图,内部依次调用三个下游服务:

  1. 用户服务(HTTP)拿用户基础信息,耗时约500ms
  2. 工单服务(HTTP)拿工单列表,耗时约1.2s
  3. 日志服务(HTTP)拿每条工单的操作日志,耗时约1.5s

三个调用串行执行,单次请求总耗时稳定在3.2秒左右。压测时发现一个诡异现象:并发30个请求时,CPU占用率只有30%左右,但QPS只有120,大量线程阻塞在socket.read上。说白了就是线程都在等IO,CPU在摸鱼

我试着把Flask的线程池从默认的10调到50,QPS勉强到180,但内存涨了快300MB,而且GIL导致上下文切换开销剧增。这条路走不通。

二、环境与版本:先交代清楚再动手

改造前的技术栈:

  • Python 3.11.4(官方文档说asyncio在3.10+稳定了,但3.11的TaskGroupasyncio.timeout更顺手)
  • Flask 2.3.3 + gunicorn 21.2.0(worker_class=sync,workers=4)
  • requests 2.31.0(同步HTTP客户端)
  • 下游服务:三个内部微服务,均为HTTP/1.1,平均响应时间见上文

改造后的目标栈:

  • FastAPI 0.104.0(原生支持async def路由)
  • uvicorn 0.24.0(worker_class=asyncio,workers=2,因为asyncio是单线程事件循环,多worker反而会竞争CPU)
  • httpx 0.25.1(异步HTTP客户端,支持连接池)
  • uvloop 0.19.1(加速事件循环)

注意:不要用aiohttp做HTTP客户端,它的连接池管理在0.25版本之前有泄漏问题,我们踩过坑,后面细说。

三、方案设计:并行请求 + 连接池复用

核心思路只有一句话:把三个串行IO变成并行IO。但实际落地要考虑三个问题:

  1. 连接复用:每次请求新建TCP连接,握手开销就要浪费几十毫秒。必须用httpx的AsyncClient作为连接池。
  2. 并发控制:下游服务没有限流保护,直接全量并发会打爆它们。用asyncio.Semaphore限制最大并发数为20。
  3. 超时控制:原来的同步代码用requests.get(timeout=3),改成asyncio后必须用asyncio.timeoutasyncio.wait_for,否则一个下游卡死会导致整个事件循环阻塞。

架构图(文字版):

客户端请求 → FastAPI路由(async def) → 创建httpx.AsyncClient(连接池) 
    → asyncio.gather(三个协程) → 分别请求三个下游服务
    → 聚合数据 → 返回JSON

关键点:AsyncClient要挂在lifespan里,不要每次请求都新建。我们第一次就是这么干的,结果压测时TCP连接数飙升到2000+,被运维警告了。

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

Before(同步Flask版本,节选)

# app.py - 同步版本
from flask import Flask, jsonify
import requests

app = Flask(__name__)

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

def fetch_orders(user_id: str) -> list:
    resp = requests.get(f"http://order-service/api/orders?user_id={user_id}", timeout=3)
    resp.raise_for_status()
    return resp.json()["items"]

def fetch_order_logs(order_ids: list) -> dict:
    resp = requests.post(
        "http://log-service/api/logs/batch",
        json={"order_ids": order_ids},
        timeout=3
    )
    resp.raise_for_status()
    return resp.json()

@app.route("/api/report/")
def report(user_id: str):
    user = fetch_user_info(user_id)          # 串行 500ms
    orders = fetch_orders(user_id)           # 串行 1.2s
    logs = fetch_order_logs([o["id"] for o in orders])  # 串行 1.5s
    return jsonify({"user": user, "orders": orders, "logs": logs})

After(异步FastAPI版本,节选)

# main.py - 异步版本
from fastapi import FastAPI, HTTPException
import asyncio
import httpx
from contextlib import asynccontextmanager

# 全局连接池,供所有请求复用
_client: httpx.AsyncClient | None = None
_semaphore = asyncio.Semaphore(20)  # 限制下游并发

@asynccontextmanager
async def lifespan(app: FastAPI):
    global _client
    # 连接池参数:max_connections=50,timeout=5s
    _client = httpx.AsyncClient(
        timeout=httpx.Timeout(5.0),
        limits=httpx.Limits(max_connections=50, max_keepalive_connections=20)
    )
    yield
    await _client.aclose()

app = FastAPI(lifespan=lifespan)

async def fetch_user_info(user_id: str) -> dict:
    async with _semaphore:
        resp = await _client.get(f"http://user-service/api/users/{user_id}")
        resp.raise_for_status()
        return resp.json()

async def fetch_orders(user_id: str) -> list:
    async with _semaphore:
        resp = await _client.get(f"http://order-service/api/orders?user_id={user_id}")
        resp.raise_for_status()
        return resp.json()["items"]

async def fetch_order_logs(order_ids: list) -> dict:
    async with _semaphore:
        resp = await _client.post(
            "http://log-service/api/logs/batch",
            json={"order_ids": order_ids}
        )
        resp.raise_for_status()
        return resp.json()

@app.get("/api/report/{user_id}")
async def report(user_id: str):
    # 三个协程并行执行
    user, orders, logs = await asyncio.gather(
        fetch_user_info(user_id),
        fetch_orders(user_id),
        fetch_order_logs([])  # 注意:这里需要先拿到orders才能查logs
    )
    # 实际改造中,logs依赖orders,所以分两步:
    # user, orders = await asyncio.gather(fetch_user_info(user_id), fetch_orders(user_id))
    # logs = await fetch_order_logs([o["id"] for o in orders])
    return {"user": user, "orders": orders, "logs": logs}

注意上面代码里有个注释很关键:logs依赖orders的数据,所以不能简单三个协程一起gather。我实际拆成了两步:

# 第一步:并行拿到user和orders
user, orders = await asyncio.gather(
    fetch_user_info(user_id),
    fetch_orders(user_id)
)
# 第二步:拿到orders后,再查logs
logs = await fetch_order_logs([o["id"] for o in orders])

这样总耗时从3.2秒降到 max(500ms, 1200ms) + 1500ms ≈ 2.7秒,其实提升不大。真正的优化在第二步:如果工单数量很多(比如50个),logs接口是支持批量的,我们直接把50个工单ID一次性POST过去,对下游日志服务来说,它内部也是并行查库的。所以第二步的1.5秒其实是批量查询,没法再拆了。

但这里还有一个隐藏优化点:如果日志服务支持按user_id直接查,就可以跟orders并行。可惜我们没权限改下游接口,只能这样。

五、踩坑与优化:三个坑,每一个都掉进去过

坑1:asyncio.run不能用在FastAPI里

一开始我在report函数里写asyncio.run(fetch_user_info(...)),直接报错RuntimeError: asyncio.run() cannot be called from a running event loop。因为FastAPI的async路由本身就在事件循环里跑,你再开一个新的就是嵌套循环。正确做法是直接await,或者用asyncio.create_task

坑2:httpx连接池泄漏

httpx.AsyncClient时,如果不用async with上下文管理器包裹,连接永远不会释放。我们第一版代码把AsyncClient写在函数内:

async def fetch_user_info(user_id):
    async with httpx.AsyncClient() as client:  # 每次请求新建连接池
        resp = await client.get(...)

压测30分钟,ss -s显示TIME_WAIT连接数从几十涨到4000+,最终触发ConnectionResetError。后来改成全局单例连接池,并设置max_keepalive_connections=20,连接复用率从10%提升到95%。

坑3:Semaphore放在函数外是全局的,但要注意事件循环绑定

asyncio.Semaphore(20)如果在模块顶层创建,没问题。但如果放在lifespan里创建,然后传给路由函数,会报got Future attached to a different loop。因为FastAPI启动时会创建新的事件循环,而你在lifespan里创建的Semaphore绑定的是lifespan的事件循环。解决方法是在模块顶层创建全局Semaphore,或者在路由函数内部用asyncio.Lock()包装。

优化:uvloop替换默认事件循环

Python 3.11默认的asyncio事件循环是SelectorEventLoop,纯Python实现。换成uvloop(基于libuv的C扩展),IO调度开销降低30%左右。只需在启动时加一句:

import uvloop
asyncio.set_event_loop_policy(uvloop.EventLoopPolicy())

但要注意:uvloop不支持Windows,我们生产环境是Linux,没问题。

六、效果数据:从120到850,翻了7倍

压测工具:wrk,参数-t4 -c100 -d60s(4线程,100并发,持续60秒)。

指标 改造前(Flask同步) 改造后(FastAPI异步) 提升
QPS 120 850 7.1倍
P50延迟 3.2s 95ms 33倍
P95延迟 2.8s(实际是长尾,P50都3.2s了) 180ms 15.6倍
CPU占用率 30%(闲置) 85%(充分使用) 2.8倍
内存占用 450MB(50线程栈) 180MB(单线程+协程) -60%
下游服务调用次数 每次请求3次HTTP 3次(但连接复用) 相同

性能对比图(文字版)

QPS曲线(100并发下):
before: ▁▁▂▂▂▃▃▂▂▁▁ (峰值180,稳定120)
after:  ▅▆▇███▇▆▅▆▇ (峰值920,稳定850)

注意P95延迟从2.8秒降到180ms,这个提升比QPS更关键。因为同步版本在100并发下,线程池被占满,新请求排队等待,P95直接飙到8秒+,我们当时被业务方投诉"报表加载转圈超过10秒"。

额外收益:部署副本数从4个降到2个(因为单实例吞吐翻倍),云服务器成本每月省了约2000元。

七、总结:异步编程不是银弹,但IO密集场景是真的香

这次改造给我最大的感悟是:异步编程解决的不是"快"的问题,而是"等待"的问题。同步代码里线程在等IO时,CPU资源是浪费的;异步代码把等待让出来,让其他协程用CPU。但要注意:

  1. CPU密集场景别用asyncio,它不会提升计算速度,反而因协程切换更慢。
  2. 依赖硬编码的下游调用顺序,如果服务间有强依赖(A的结果是B的参数),只能串行,这时候异步收益有限。
  3. 连接池和并发控制是生命线,没有Semaphore保护,你的下游服务会被打挂;没有连接池复用,你的TCP连接数会爆炸。

如果让我重新选,我会直接上asyncio + httpx + uvloop这个组合,而不是先试Flask的线程池调优。省下的三天调优时间,够我写两篇博客了。

最后贴一下生产环境的启动命令,方便大家参考:

# uvicorn启动,2个worker(asyncio下多worker反而降低性能,因为事件循环争抢CPU)
uvicorn main:app --host 0.0.0.0 --port 8080 --workers 2 --loop uvloop --http h11

以上,希望对你有用。有问题评论区见。