一、问题背景:一个慢查询API拖垮了整个微服务

去年底我们接手了一个内部数据中台项目,其中一个名为/api/v1/batch-query的接口,接收一组用户ID,并行请求三个下游服务(用户画像、行为标签、信用评分),聚合后返回。

最初版本用requests库发起HTTP请求,在并发数达到50时,接口响应时间急剧恶化。生产环境中,该接口的TP99长期维持在2.3s以上,高峰期甚至导致上游服务雪崩。监控显示,单次请求中,网络IO等待时间占比超过80%。这意味着,99%的CPU时间都在等网络响应。

显然,这是一个典型的IO密集型场景,非常适合用asyncio进行优化。

二、环境与版本

本次优化的运行环境及关键依赖版本如下:

组件 版本
Python 3.10.12
FastAPI 0.104.1
Uvicorn 0.24.0
aiohttp 3.9.1
httpx 0.25.2
uvloop 0.19.0
操作系统 Ubuntu 22.04 LTS
压测工具 wrk 4.2.0

服务部署在4核8G的云服务器上,下游三个服务均为同一机房内的HTTP接口,平均延迟约80ms。

三、方案设计:用asyncio替换同步阻塞

核心思路:将每个下游请求封装为异步协程,通过asyncio.gather并发调度,同时复用TCP连接以减少握手开销

具体设计如下:
1. 使用aiohttp.ClientSession作为连接池,设置连接池大小=200,每主机最大连接数=100,启用keepalive。
2. 引入uvloop替换默认的asyncio事件循环,利用libuv提升事件调度性能。
3. 设置合理的超时:连接超时5s,读取超时10s,避免单个慢请求阻塞整个协程组。
4. 控制并发量:使用asyncio.Semaphore限制同时发起的请求数,防止下游被打爆。

四、核心实现:Before vs After代码

Before(同步版本)

# sync_batch_query.py (优化前)
import requests
from fastapi import FastAPI, HTTPException
from typing import List, Dict

app = FastAPI()
TIMEOUT = (5, 10)  # connect, read timeout

@app.post("/api/v1/batch-query")
async def batch_query(user_ids: List[str]) -> Dict:
    try:
        # 串行请求三个下游服务
        profile = requests.post(
            "http://profile-service/get", 
            json={"ids": user_ids}, 
            timeout=TIMEOUT
        )
        behavior = requests.post(
            "http://behavior-service/get", 
            json={"ids": user_ids}, 
            timeout=TIMEOUT
        )
        credit = requests.post(
            "http://credit-service/get", 
            json={"ids": user_ids}, 
            timeout=TIMEOUT
        )

        return {
            "profiles": profile.json(),
            "behaviors": behavior.json(),
            "credits": credit.json()
        }
    except requests.Timeout:
        raise HTTPException(status_code=504, detail="下游服务超时")

问题分析:虽然FastAPI的handler函数写成了async def,但内部使用的是同步requests库。当三个requests.post被执行时,当前线程会阻塞等待IO,实际上还是串行执行。FastAPI的异步能力在此处完全被浪费。

After(异步版本)

# async_batch_query.py (优化后)
import asyncio
import aiohttp
from fastapi import FastAPI, HTTPException
from typing import List, Dict

app = FastAPI()

# 全局连接池,避免每次请求都新建
CONNECTOR = aiohttp.TCPConnector(
    limit=200,          # 连接池总大小
    limit_per_host=100, # 每主机最大连接数
    ttl_dns_cache=300,  # DNS缓存5分钟
    force_close=False,  # 保持keepalive
)

# 信号量控制并发数,防止下游雪崩
SEMAPHORE = asyncio.Semaphore(50)

async def fetch(session: aiohttp.ClientSession, url: str, payload: Dict) -> Dict:
    async with SEMAPHORE:
        try:
            async with session.post(url, json=payload, timeout=aiohttp.ClientTimeout(
                total=15, connect=5, sock_read=10
            )) as resp:
                resp.raise_for_status()
                return await resp.json()
        except asyncio.TimeoutError:
            raise HTTPException(status_code=504, detail=f"服务超时: {url}")
        except aiohttp.ClientError as e:
            raise HTTPException(status_code=502, detail=str(e))

@app.post("/api/v1/batch-query")
async def batch_query(user_ids: List[str]) -> Dict:
    async with aiohttp.ClientSession(connector=CONNECTOR) as session:
        tasks = [
            fetch(session, "http://profile-service/get", {"ids": user_ids}),
            fetch(session, "http://behavior-service/get", {"ids": user_ids}),
            fetch(session, "http://credit-service/get", {"ids": user_ids}),
        ]
        # 并发执行三个协程
        results = await asyncio.gather(*tasks, return_exceptions=True)

        # 检查是否有异常
        for r in results:
            if isinstance(r, HTTPException):
                raise r

        return {
            "profiles": results[0],
            "behaviors": results[1],
            "credits": results[2]
        }

关键改动
1. aiohttp.ClientSession复用全局连接池,避免每次请求都三次握手。
2. asyncio.gather并发调度三个协程,IO等待时自动切换。
3. asyncio.Semaphore(50)控制同时进行的对外请求数,防止下游服务因突发流量被打挂。
4. 使用raise_for_status()和自定义异常处理,保持错误信息清晰。

五、踩坑与优化

踩坑1:忘了替换事件循环

第一次测试时发现性能只提升了30%,排查后才发现,uvloop需要手动设置:

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

在旧版Python中,默认事件循环是SelectorEventLoop,性能远不如uvloop。使用uvloop后,事件循环的epoll调度效率提升约40%。

踩坑2:连接池耗尽导致连锁失败

初期没有使用Semaphore,结果在压测时,200并发请求瞬间创建了600个到下游的TCP连接(每个请求3个下游),下游服务扛不住,出现大量Connection Refused。加上Semaphore后,将并发数控制在50,下游服务稳定运行。

踩坑3:aiohttp的timeout参数位置

aiohttp的timeout必须通过aiohttp.ClientTimeout对象传入,而不是简单的元组。开始时误写成session.post(url, timeout=(5,10)),结果timeout完全无效,因为aiohttp不识别元组格式。

优化:使用httpx替代aiohttp

如果项目要求更高,可以考虑httpx(0.25.2版本),它同时支持同步和异步API,且API设计更贴近requests。但aiohttp在连接池控制上更精细,适合高并发场景。

六、效果数据:从380到1600 QPS

使用wrk进行压测,命令如下:

wrk -t4 -c200 -d30s --latency http://target-host/api/v1/batch-query

压测结果对比:

指标 同步版本 异步版本 提升倍数
QPS 380 1600 4.2x
TP50 480ms 120ms 4.0x
TP99 2.3s 450ms 5.1x
错误率(5xx) 3.2% 0.1% -

CPU利用率:同步版本在50并发时CPU利用率就达到85%(大量上下文切换),异步版本在200并发时CPU利用率仅60%(协程切换开销远小于线程切换)。

内存占用:同步版本每个请求占用约3MB(线程栈),异步版本每个请求仅占用约50KB(协程栈),内存消耗降低60倍。

七、总结

  1. asyncio不是银弹:它只对IO密集型任务有效。如果你的API涉及大量CPU计算(如加密、图像处理),应该考虑多进程或线程池。

  2. 连接池和并发控制是核心:没有连接池的异步编程,就像没有缓冲区的管道,流量一冲就垮。asyncio.Semaphore是保护下游服务的必备工具。

  3. 版本兼容性很重要:Python 3.10+的asyncio已经比较成熟,但如果使用3.7或3.8,需要注意一些API差异(如asyncio.run底层实现不同)。

  4. 压测数据要说谎:务必在生产环境流量模式下进行压测,单纯用wrk压本地接口,忽略了网络抖动、下游限流等因素,容易得出过于乐观的结论。

最终,这个优化版本上线后,该接口的TP99稳定在500ms以下,再也没有出现过因该接口导致的上游雪崩。异步编程的核心价值,不是让代码更快,而是让等待不再阻塞资源。

如果你们也有类似的阻塞IO噩梦,不妨试试asyncio。记住:别用同步库写异步函数,那只是心理安慰。