一、问题背景:同步IO把CPU饿死了
先交代一下业务场景。我们有个内部工单系统,前端有一个"用户所有工单+关联操作日志"的聚合报表接口。原来是用Flask写的同步视图,内部依次调用三个下游服务:
- 用户服务(HTTP)拿用户基础信息,耗时约500ms
- 工单服务(HTTP)拿工单列表,耗时约1.2s
- 日志服务(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的
TaskGroup和asyncio.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。但实际落地要考虑三个问题:
- 连接复用:每次请求新建TCP连接,握手开销就要浪费几十毫秒。必须用httpx的
AsyncClient作为连接池。 - 并发控制:下游服务没有限流保护,直接全量并发会打爆它们。用
asyncio.Semaphore限制最大并发数为20。 - 超时控制:原来的同步代码用
requests.get(timeout=3),改成asyncio后必须用asyncio.timeout或asyncio.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。但要注意:
- CPU密集场景别用asyncio,它不会提升计算速度,反而因协程切换更慢。
- 依赖硬编码的下游调用顺序,如果服务间有强依赖(A的结果是B的参数),只能串行,这时候异步收益有限。
- 连接池和并发控制是生命线,没有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
以上,希望对你有用。有问题评论区见。