Back to Journal
02 / Entry· 20 min read

FastAPI 并发实战:进程、线程、协程在一个框架里的分工

从 epoll/kqueue 一路拆到 libuv、uvloop、httptools、Uvicorn、Starlette 和 Pydantic Core,解释 FastAPI 为什么快,再用阻塞、异步串行超时、线程池、多进程与 Agent API 说明它的边界

🔊 系统朗读

上一篇把进程、线程、协程的机理拆了一遍:GIL 为什么让多线程算不快、await 为什么是让出控制权、事件循环靠什么扛住上万连接。

机理懂了,放到真实框架里长什么样?FastAPI 是最好的观察样本——它把三种并发手段全用上了:每个 worker 是一个进程,进程里一个事件循环跑协程,协程处理不了的同步代码扔进线程池。这篇就用 FastAPI 把上一篇的知识全部兑现,所有数据都是真实压测跑出来的(FastAPI 0.141.1 + Uvicorn 0.52.1 + Python 3.12,10 核 Mac)。

FastAPI 为什么快?先把“快”说清楚

“FastAPI 很快”经常被说成一句口号,但这里至少混着三件事:

所谓的快真正含义不代表什么
协议路径开销低收包、HTTP 解析、路由和校验这条热路径尽量少在 Python 里兜圈业务函数不需要时间
IO 并发吞吐高一个线程可以挂起大量正在等数据库、LLM、RAG 的请求单次下游调用会变快
开发速度快类型注解同时生成校验、OpenAPI 和交互文档运行时零成本

真正让它在基准测试和 IO 型服务里表现好的,是前两项:底层热路径由成熟的 C/Rust 组件承担,上层再用事件循环避免“一个连接占一个线程”

所以先给结论:FastAPI 没有发明新的网络算法。它更像一个总装厂,把操作系统、libuv、uvloop、httptools、Uvicorn、ASGI、Starlette 和 Pydantic Core 接成了一条短路径。每一层都只做自己擅长的事。

FastAPI 的家底:从 Python 一路站到操作系统内核

把完整栈展开,不是三层,而是下面这些层:

text
你的路径函数 / Agent 编排
          ▲
FastAPI   │ 依赖图、参数提取、OpenAPI、请求/响应模型
          ├──────────────> Pydantic Core(Rust:校验与序列化核心)
          ▲
Starlette │ 路由、中间件、Request/Response、WebSocket
          ▲
   ASGI   │ scope + receive + send 消息协议
          ▲
Uvicorn   │ 连接管理、协议适配、流控、创建请求 Task
          ├─ HTTP:httptools(C 解析器)/ h11(Python 回退)
          └─ Loop:uvloop(Cython + libuv)/ asyncio(标准回退)
                         ▲
                       libuv(C)
                         ▲
Linux epoll / macOS kqueue / 其他平台对应的内核机制
                         ▲
                   非阻塞 socket

这里有个重要的准确性问题:uvloop 和 httptools 不是任何环境里都必然启用。Uvicorn 的 auto 模式会在安装了 uvloop 时优先选择它,否则回退到标准 asyncio;HTTP 解析也是 httptools 可用就优先使用,否则回退 h11。安装 uvicorn[standard] 会带上这些常用加速依赖:

bash
python -m pip install "uvicorn[standard]"
uvicorn app:app --loop auto --http auto

也可以显式指定 --loop uvloop --http httptools,让依赖缺失时直接暴露问题,而不是静默回退。具体选择逻辑可以看 Uvicorn 事件循环文档httptools 项目

下面继续往下挖:这些“巨人”各自快在哪里?

第一层巨人:内核的 IO 多路复用

最朴素的服务器模型是一个连接配一个线程。线程调用阻塞式 recv() 后,就睡在那里等数据。连接少时很直观;连接多时,哪怕绝大多数客户端都在思考、打字或等下游,每条连接仍要占线程栈和调度资源,操作系统还要在大量线程之间做上下文切换。

FastAPI 这条栈走的是另一条路:socket 被设成非阻塞,线程不在某个连接上干等,而是把一批 socket 交给内核:

text
“这 10000 个 socket 里,哪个能读/能写了再通知我。”

Linux 用 epoll,macOS/BSD 用 kqueue。内核掌握网卡和 socket 缓冲区的真实状态,可以只把已经就绪的事件交回来,不需要 Python 一遍遍扫描所有连接。libuv 的设计文档描述的就是这套单线程、非阻塞网络 IO 模型。

一次等待的实际过程是:

text
协程执行 await socket.read()
  -> 当前 Task 挂起
  -> socket fd 注册到事件循环
  -> 线程进入 kqueue/epoll 等待
  -> 内核发现 fd 可读
  -> 事件循环把对应 Task 放回 ready queue
  -> 协程从 await 后面继续

这一步省掉的不是网络传输时间,而是陪着网络一起等的线程。10000 个连接大部分都没数据时,一个线程只处理那一小撮真正就绪的连接。这是高并发的地基。

第二层巨人:libuv 和 uvloop

操作系统接口在不同平台上并不一样。libuv 用 C 把 epoll、kqueue、定时器、异步通知、socket 和进程信号封装成统一事件循环;它长期服务于 Node.js 生态,核心价值是成熟、跨平台和热路径足够靠近系统调用。

uvloop 再用 Cython 实现 Python 的 asyncio.AbstractEventLoop 接口,把 asyncio 的 socket 监听、transport、timer 和 callback 调度接到 libuv;Task/Future 仍然遵守 asyncio 的协作式调度模型。对业务代码来说:

python
await client.get(url)

写法和语义完全没有改变;换掉的是下面负责“登记 fd、轮询内核、管理 timer、唤醒 callback”的发动机。相比大量逻辑都在 Python 层运行,uvloop 把更多调度、缓冲区和 transport 工作压到 C/Cython + libuv 路径中,减少 Python 函数调用、对象周转和解释器开销。

因此 uvloop 提升的是事件循环调度和网络 transport 的效率,不是让协程自动并行,更不会让一个 8 秒的 RAG 请求变成 2 秒。它是 asyncio 的高性能替换实现,不是另一套 async 语义。

在这篇文章的 macOS 压测环境里,如果 uvloop 已启用,最底下实际走的是:

text
asyncio API -> uvloop -> libuv -> kqueue -> macOS 内核

第三层巨人:httptools 把 HTTP 字节解析留在 C 层

内核通知“socket 有数据”以后,拿到的只是字节:

http
POST /agents/run HTTP/1.1\r\n
Host: example.com\r\n
Content-Type: application/json\r\n
Content-Length: 37\r\n
\r\n
{"message":"帮我查一下 FastAPI"}

服务器要从连续字节流里识别请求行、URL、Header、Content-Length 和 Body 边界。这是每个 HTTP 请求都要走的高频路径。httptools 是高性能 HTTP 解析器的 Python 绑定,核心状态机在 C 层增量消费字节,并在 URL、Header、Body、消息完成等节点回调 Uvicorn,而不是用 Python 字符串操作一段段拆。

这会降低协议解析的 CPU 开销和临时 Python 对象数量。但它只负责 HTTP 信封,不会校验 Body 里的业务 JSON;message 是不是字符串、session_id 是否缺失,是后面 FastAPI/Pydantic 的工作。

如果 httptools 没安装,Uvicorn 会使用 h11 这个纯 Python HTTP/1.1 状态机。功能仍然正确,只是执行路径不同。所以“FastAPI 使用 C 解析 HTTP”应该带上前提,不能写成无条件事实。

第四层巨人:Uvicorn 把网络协议翻译成 ASGI

uvloop 和 httptools 解决了“怎么高效收字节、拆 HTTP”,但它们不知道什么是 FastAPI 路由。Uvicorn 负责中间的协议适配:

  1. 接受连接,维护 transport 和 HTTP 协议状态。
  2. 把 method、path、headers 等组装成 ASGI scope
  3. 为请求创建 Task,调用上层 ASGI 应用。
  4. 把请求 Body 变成 http.request 消息交给 receive()
  5. 接收应用通过 send() 发来的 http.response.start/body
  6. 编码成 HTTP 响应并写回 socket。

ASGI 应用最小可以小到这样:

python
async def app(scope, receive, send):
    request = await receive()
    assert request["type"] == "http.request"

    await send({
        "type": "http.response.start",
        "status": 200,
        "headers": [(b"content-type", b"text/plain")],
    })
    await send({
        "type": "http.response.body",
        "body": b"hello",
    })

scope 描述连接,receivesend 是两个异步消息通道。服务器只认协议,框架只认 ASGI 消息,两边因此可以独立演进。这就是 ASGI 规范真正解决的问题。

Uvicorn 还做了一个常被忽略的性能保护:背压。如果请求 Body 缓冲达到高水位,它会暂停继续读取;应用消费后再恢复。响应写缓冲太高时,send 也会等待缓冲降下来。背压不一定让理想基准数字更漂亮,但能防止慢客户端或大请求把内存无限吃掉,具体行为见 Uvicorn 流控文档。稳定地快,比短时间跑得快更重要。

第五层巨人:Starlette 保持 Web 层足够薄

Uvicorn 把请求交给上层后,先接住它的其实是 Starlette。FastAPI 的路由、中间件、Request/Response、WebSocket、静态文件和 BackgroundTasks 等能力,大量建立在 Starlette 上。

Starlette 的关键是遵守同一套 ASGI callable 形式。中间件不是神秘钩子,而是一层包一层的异步函数:

text
Uvicorn
  -> 错误处理中间件
    -> CORS / 鉴权 / 日志中间件
      -> Router
        -> Endpoint

每一层都接收 scope, receive, send,做完自己的事再调用下一层。没有为了每个异步请求额外创建专属线程,也没有把网络层重新包装一遍。Starlette 本身功能少、路径短,所以适合作为 FastAPI 的 Web 地基。

也正因为如此,裸 Starlette 通常会比 FastAPI 再少一点开销:它没有 FastAPI 的依赖解析和 Pydantic 模型校验。FastAPI 所谓“快”,不是比自己的地基还快,而是在增加类型安全、依赖注入和文档能力以后,仍然保持较低的框架开销。

第六层巨人:Pydantic Core 把类型校验压到 Rust

请求终于到了 FastAPI。它主要做三件事:

  1. 根据路由签名,从 path、query、header、cookie 和 body 提取参数。
  2. 解析依赖图,执行鉴权、数据库会话等依赖。
  3. 用请求/响应模型完成校验、转换和序列化。

FastAPI 不会在每次请求里重新反射整份函数签名。注册 APIRoute 时,它已经根据类型注解构建依赖关系和字段元数据,请求到来后按这张图执行。OpenAPI schema 也不是每个业务请求都重新生成。

Pydantic v2 又把大部分类型验证和序列化原语下沉到独立的 pydantic-core。Python 侧把 BaseModel 和类型注解编译成 CoreSchema,Rust 核心按这份 schema 执行字符串、数字、日期、容器和嵌套模型等校验,避免用 Python 递归解释每一个字段。架构细节可以看 Pydantic 内部架构

这里也不能神化:校验永远有成本。一个不做模型校验、直接返回预编码 bytes 的 Starlette 接口当然可能更快;复杂依赖、深层模型和大响应也会增加耗时。FastAPI 的价值是用一笔可控成本换来明确契约、统一错误和自动文档,而 Pydantic Core 把这笔成本尽量压低。

每一层到底省掉了什么

负责什么主要省掉的开销
epoll / kqueue通知哪些 socket 已就绪为每个等待连接保留线程、反复扫描连接
libuv / uvloop事件循环、transport、timer、callback大量 Python 层调度与协议对象周转
httptools增量解析 HTTP/1.1 字节Python 层字符串拆分和状态机开销
UvicornHTTP/WebSocket 与 ASGI 互转、流控保持协议适配层精简,并用背压限制缓冲区增长
Starlette路由、中间件、请求响应抽象厚重 Web 框架路径和不必要功能
Pydantic Core类型验证与序列化核心Python 逐字段递归校验
FastAPI依赖图、参数系统、OpenAPI 整合每次请求重新分析签名和手写胶水代码

现在再说“站在巨人的肩膀上”,就不是一句套话了:内核负责发现事件,C/Cython 负责搬运与调度,C 解析 HTTP,Rust 校验数据,Python 主要保留业务编排。

一个请求的完整旅程:从网卡到路径函数

把这些层串起来,一次请求的真实旅程是这样的:

text
1. 网卡收到数据,内核 TCP 栈把字节放进 socket 接收缓冲区
2. kqueue/epoll 标记这个 socket 可读,唤醒正在等待的事件循环
3. libuv/uvloop 触发对应 transport/protocol 回调
4. httptools 增量解析请求行、Header 和 Body 边界
5. Uvicorn 组装 ASGI scope 和 http.request 消息,创建请求 Task
6. Starlette 依次经过中间件和 Router,找到目标 Endpoint
7. FastAPI 解析依赖和参数,Pydantic Core 校验输入
8. 路径函数开始执行
9. 路径函数 await 数据库、LLM 或 RAG,当前 Task 挂起,线程回到第 2 步
10. 下游 fd 就绪,Task 恢复;响应模型序列化后沿 ASGI send 原路写回 socket

第 9 步才是 IO 并发吞吐高的核心。假设同一线程上有三个请求:

text
时间 ──────────────────────────────────────────────>
请求 A  执行 ── await LLM ────────────────── 恢复 ── 返回
请求 B           执行 ── await DB ── 恢复 ── 返回
请求 C                    执行 ── await RAG ───── 恢复
线程    [A]      [B]      [C]          [B]          [A][C]

线程没有同时执行 A、B、C 的 Python 代码,它只是在每个任务等待 IO 时去推进另一个任务。切换发生在明确的 await 点,由用户态 Task 调度完成,不需要操作系统在大量请求线程之间抢占切换。

这解释了 FastAPI 的两个性能来源:

  • 每次真正运行时,路径短:网络、HTTP 解析和模型校验中的重活尽量落在 C/Rust 层。
  • 每次不得不等待时,不占着线程:Task 挂起,事件循环立刻服务其他已就绪请求。

但也解释了它的边界:

  • 一个 FSRAG 调用本身需要 8 秒,FastAPI 不会把它变成 2 秒。
  • for query in queries: await search(query) 仍然是异步串行,当前请求耗时仍会累加。
  • async def 里执行同步阻塞库,Task 没机会让出,整条事件循环都会停。
  • 大体量 JSON、复杂模型校验、图片处理、本地推理等 CPU 工作仍然要消耗真实 CPU 时间。

所以更准确的说法不是“FastAPI 让每个请求都更快”,而是:

FastAPI 用较低的协议和框架开销处理请求,并在 IO 等待占主导时,让一个线程高效承载大量并发;它优化的是资源利用率和吞吐,不会消灭业务本身的耗时。

核心决策:def 还是 async def

FastAPI 最精妙、也最容易踩坑的设计来了:路径函数两种写法都接受,但走的是两条完全不同的执行路径

Starlette 的路由层会检查你的路径函数:

text
def 路径函数       -> 扔进线程池执行(不占用事件循环)
async def 路径函数 -> 直接在事件循环里 await(协程)

也就是说,FastAPI 替你把上一篇的选型问题做了一层封装:同步代码自动获得线程池,异步代码享受协程。但前提是——你得写对。

实测:三种写法,三种命运

起一个演示服务,三个端点都模拟 0.5 秒的 IO,写法各不相同:

python
import asyncio
import time
from fastapi import FastAPI

app = FastAPI()

@app.get("/async-right")
async def async_right():
    """正确:async def + 异步等待"""
    await asyncio.sleep(0.5)
    return {"ok": True}

@app.get("/async-wrong")
async def async_wrong():
    """错误:async def 里调阻塞函数"""
    time.sleep(0.5)   # 拖死整个事件循环
    return {"ok": True}

@app.get("/def-blocking")
def def_blocking():
    """正确:def 路径函数,FastAPI 自动扔进线程池"""
    time.sleep(0.5)
    return {"ok": True}

然后并发 8 个请求打同一个端点,看总耗时:

端点8 并发耗时执行位置
/async-right(async + await)0.53s事件循环线程(协程并发)
/async-wrong(async + 阻塞)4.05s事件循环线程(全部串行)
/def-blocking(def + 阻塞)0.52sAnyIO worker thread(线程池)

数据和上一篇的实验严丝合缝:

  • async defawait,8 个请求在事件循环里交替推进,总耗时等于一个请求的耗时
  • async def 里写 time.sleep,协程没有任何让出点,事件循环被按死,8 个请求排队串行——0.5s × 8 = 4.05s,一秒不多
  • def 写法虽然里面是阻塞调用,但 FastAPI 把它扔进了线程池,8 个请求分到不同线程并行等待

事故现场:一个错误接口拖垮全站

单个端点串行还只是慢。真实服务里更可怕的是连坐——事件循环是全局只有一个的,一个接口把它堵住,所有接口陪葬。实测:同时发 1 个 /async-wrong 和 4 个 /async-right

text
/async-wrong   : 0.51s    <- 罪魁祸首
/async-right   : 1.01s    <- 无辜接口被拖到两倍耗时
/async-right   : 1.01s
/async-right   : 1.01s
/async-right   : 1.01s

4 个写法完全正确的接口,仅仅因为和一个错误接口共享事件循环,响应时间全部翻倍。如果 wrong 接口里是 10 秒的阻塞调用,这 10 秒内整个服务对所有用户不可用。

这就是 FastAPI 生产事故的第一大来源:async def 里调了同步阻塞库——requests、同步版数据库驱动、某些云的官方 SDK、open().read(),甚至一条看似无害的 time.sleep。上一篇说的"一个阻塞调用拖死所有任务",在 FastAPI 里就是字面意义上的全站瘫痪。

线程池有多大?

def 路径函数去的那个线程池不是无限的。FastAPI 用的是 AnyIO 的线程池,默认容量 40。用 50 个并发请求打 /def-blocking 验证:

text
40 并发: 0.59s   <- 刚好一批
50 并发: 1.04s   <- 40 个先跑,剩下 10 个排队等第二批

50 个请求果然分成了两批。这意味着如果大量请求都走 def 路径函数,40 就是这个进程的处理能力天花板,多余的请求要排队。容量可以通过 AnyIO 的 capacity limiter 调整,但默认值在大部分场景够用——真不够用的信号,往往是你把太多本该异步的活塞给了线程池。

决策树

text
路径函数里要做什么?

├─ 有 IO,且用的是异步库(httpx、asyncpg、aiomysql、redis.asyncio)
│    └─ async def,一路 await 到底
│
├─ 有 IO,但库只有同步版(requests、psycopg2、部分云厂商 SDK)
│    └─ 用 def(FastAPI 自动进线程池)
│       或 async def + await asyncio.to_thread(阻塞函数)
│
├─ 纯 CPU 计算
│    ├─ 几毫秒内能算完:async def 直接写,别折腾
│    └─ 算得久:谁都不合适,交给任务队列或多进程
│
└─ 没有任何 IO,只是读内存/拼数据
     └─ async def,零开销

铁律只有一条,再强调一遍:async def 里不允许出现任何阻塞调用。拿不准一个库是不是异步的,就看它的文档里有没有要求你 await——同步函数在协程里没有让出点,没有例外。

全程都是 async,为什么最后还是超时?

还有一种事故比阻塞事件循环更隐蔽:代码里的每个 IO 都正确地写了 await,整个接口仍然超时。下面是一段真实 Agent 检索逻辑的核心结构,一次请求要拿多条 message 去调用 RAG:

python
async with aiohttp.ClientSession() as session:
    for query in queries:
        async with session.post(
            FSRAG_API_URL,
            json=build_payload(query),
            timeout=aiohttp.ClientTimeout(total=30),
        ) as resp:
            data = await resp.json()

这段代码确实是非阻塞的:等待 FSRAG 时,事件循环还能处理其他用户的请求,服务不会像 time.sleep 那样全站卡死。但对当前这个请求来说,for + await 仍然是严格串行:

text
query 1:  ├──── 8s ────┤
query 2:                 ├──── 9s ─────┤
query 3:                                ├──── 7s ───┤
当前请求: ├──────────────── 24s ────────────────────┤

await 的意思是“我先让出事件循环,等结果回来再从这里继续”,不是“把后面的循环也同时启动”。如果 5 个 query 平均各要 8 秒,这个函数就要约 40 秒;外层 FastAPI、Nginx 或调用方只给了 30 秒,它当然会超时,哪怕每一行都是异步代码。

这里还容易误读 ClientTimeout(total=30)。按照 aiohttp 的定义,这个 total 覆盖的是一次 session.post 的完整生命周期,包括连接、重定向、读取响应体,不是整个 queries 批次。串行执行 N 次时,批次最坏耗时可以接近:

text
T_serial ≈ t1 + t2 + ... + tN
最坏预算 ≈ N × 30s

而且 except aiohttp.ClientError 不应该被当成完整的超时处理。总超时应单独捕获 asyncio.TimeoutError,否则日志里可能只看到外层 500/504,看不到究竟是哪条 query 超时。

query 相互独立:改成有界并发

如果这些 query 互不依赖,就可以同时启动,但不要裸奔式地一次发几百个。下面是在原逻辑上做的最小重构:

python
import asyncio
import json

import aiohttp

PER_QUERY_TIMEOUT = aiohttp.ClientTimeout(
    total=20,
    connect=3,
    sock_read=15,
)
BATCH_TIMEOUT_SECONDS = 50


async def _call_fsrag(
    dataset_id: str,
    messages: list[str],
    rag_key: str,
    session: aiohttp.ClientSession,
    slots: asyncio.Semaphore,
) -> str:
    queries = [m.strip() for m in messages if m and m.strip()]
    headers = {
        "Authorization": f"Bearer {rag_key}",
        "Content-Type": "application/json",
    }

    async def fetch_one(query: str) -> list[dict]:
        payload = {
            "req": {
                "question": query,
                "dataset_ids": [dataset_id],
                "document_ids": [],
                "top_enhance": False,
                "rerank_id": "cohere.rerank-v3-5:0",
                "highlight": True,
            }
        }

        try:
            # slots 在整个 worker 内共享,避免多个 HTTP 请求一起压垮 RAG。
            async with slots:
                async with session.post(
                    FSRAG_API_URL,
                    headers=headers,
                    json=payload,
                    timeout=PER_QUERY_TIMEOUT,
                ) as resp:
                    if resp.status >= 400:
                        logger.error(
                            "FSRAG HTTP error: status=%s query=%s",
                            resp.status,
                            query,
                        )
                        return []

                    data = await resp.json()
        except asyncio.TimeoutError:
            logger.warning("FSRAG query timeout: query=%s", query)
            return []
        except aiohttp.ClientError as exc:
            logger.error("FSRAG request failed: query=%s error=%s", query, exc)
            return []

        if data.get("code") != 10000:
            logger.warning("FSRAG bad response: code=%s", data.get("code"))
            return []

        return data.get("data", {}).get("chunks", [])

    try:
        # gather 才真正让多条 query 同时推进;总 deadline 兜住整个批次。
        async with asyncio.timeout(BATCH_TIMEOUT_SECONDS):
            groups = await asyncio.gather(
                *(fetch_one(query) for query in queries)
            )
    except TimeoutError:
        logger.error("FSRAG batch timeout: queries=%d", len(queries))
        return json.dumps(
            {"success": False, "error": "FSRAG batch timeout"},
            ensure_ascii=False,
        )

    # gather 保持输入顺序。并发任务只返回结果,统一在这里去重,
    # 不让多个任务同时修改 seen/context 这两个共享容器。
    context: list[dict] = []
    seen: set[str] = set()

    for chunks in groups:
        for chunk in chunks:
            content = chunk.get("content", "").strip()
            if not content or content in seen:
                continue
            seen.add(content)
            context.append(
                {
                    "content": content,
                    "recommend_url": chunk.get("document_keyword", ""),
                    "similarity": chunk.get("similarity", 0),
                }
            )

    return json.dumps(
        {"success": True, "context": context},
        ensure_ascii=False,
    )

假设 8 条 query、并发上限 C=4、单条都耗时 8 秒,原来的串行耗时约 64 秒;改完后分两批,理想耗时约 16 秒:

text
T_concurrent ≈ ceil(N / C) × 每批最慢请求耗时

这里 SemaphoreClientSession 都作为参数传入,是有意的。原代码在一次 _call_fsrag 内确实复用了 session,但如果每个 FastAPI 请求都会调用一次这个函数,仍然会反复创建连接池;而把 Semaphore(4) 写在函数内部,也只能限制单个批次,挡不住 20 个用户各发 4 条请求。生产里应该在 lifespan 中各创建一次,放到 app.state,让同一个 worker 的所有请求共同复用连接池和并发额度。

并发也不是万能药

这次改造只解决“多条独立 query 被串行累计”的问题,另外两种情况不会因此消失:

  • 单条 query 本身就超过 30 秒:并发后它还是超过 30 秒。应该优化 RAG、缩小检索范围,或把运行改成长任务,而不是继续增大并发数。
  • query 之间有前后依赖:后一条必须使用前一条结果时就不能 gather,只能保持串行,减少轮次或改成“提交任务 -> 返回 run_id -> 异步查询结果”。
  • 下游有配额或容量上限:并发越大越容易触发 429、连接池排队或服务雪崩,所以必须用压测确定 Semaphore,不是越大越快。

最后要做的是超时预算向内递减。例如客户端 65 秒、网关 60 秒、FastAPI handler 55 秒、FSRAG 整批 50 秒、单条 query 20 秒。这样内层能先超时并返回可识别的错误;如果所有层都设成 60 秒,往往是客户端先断开,服务端还在做一堆已经没人要的工作。

这组案例补上了一个很重要的区别:

异步解决“等待时线程能不能去做别的事”,并发解决“多个等待能不能同时开始”,超时解决“最多愿意等多久”。三者不是一回事。

BackgroundTasks:响应返回之后再干活

有些活不需要让用户等:发邮件、写日志、上报埋点。FastAPI 内置了 BackgroundTasks 处理这类事:

python
from fastapi import BackgroundTasks, FastAPI

app = FastAPI()

def bg_work(name: str):
    time.sleep(0.3)
    print(f"[后台任务] {name} 完成", flush=True)

@app.post("/notify")
def notify(background_tasks: BackgroundTasks):
    background_tasks.add_task(bg_work, "任务A")
    background_tasks.add_task(bg_work, "任务B")
    return {"status": "已受理"}

实测:客户端 0.035s 就拿到了响应,而两个任务合计 0.6s 的活在响应之后才陆续完成——日志里任务A、任务B 先后打印,顺序执行

几个要点,都和上一篇的知识对得上:

  • 响应先于任务。Starlette 的实现是把响应 body 发给客户端之后,才开始逐个执行后台任务。所以用户感知不到后台任务的耗时。
  • 任务之间是顺序的,不是并发的。这点和 asyncio.create_task 不一样——create_task 登记完立刻并发调度,BackgroundTasks 是响应后按添加顺序一个一个跑。
  • 同步任务自动进线程池。和路径函数同一套规则:def 任务扔线程池,async def 任务直接 await。所以在 async 后台任务里调阻塞库,一样会堵事件循环。
  • 上一篇的坑全部适用:任务异常不会传给任何人(客户端早就拿到响应了),想感知就得自己 try/except 记日志;进程退出时没跑完的任务会被取消。

什么时候 BackgroundTasks 不够用?活太重(跑几分钟)、要求可靠投递(进程重启任务不能丢)、要定时调度——这时候该上正经的任务队列:Celery、Dramatiq、或者 asyncio 原生的 ARQ。BackgroundTasks 的定位是"顺手的小事",不是分布式任务系统。

部署:进程、线程、协程各就各位

单进程开发够了,生产环境通常需要多个进程。当前最直接的两种起法是:

bash
# FastAPI CLI,底层仍然是 Uvicorn
fastapi run app.py --workers 4

# 或者直接使用 Uvicorn
uvicorn app:app --workers 4

如果现有基础设施依赖 Gunicorn 的进程管理能力,可以安装独立的 uvicorn-worker 包:

bash
gunicorn app:app -w 4 -k uvicorn_worker.UvicornWorker

旧写法 -k uvicorn.workers.UvicornWorker 所在的 uvicorn.workers 模块已经被官方标记为弃用,别再从旧教程里复制。具体变化可以看 Uvicorn 部署文档

不管用哪个进程管理器,结构都一样:每个 worker 是一个独立进程,进程里跑一个独立的事件循环。用 --workers 2 起服务,然后打一个带全局计数器的接口验证:

python
import os

counter = 0

@app.get("/count")
async def count():
    global counter
    counter += 1
    return {"pid": os.getpid(), "counter": counter}

并发打 20 个请求,结果:

text
pid 6297: 处理了 2 个请求,  计数器=[1, 2]
pid 6298: 处理了 18 个请求, 计数器=[7, 8, ..., 24]

两个进程的计数器各自独立累加——上一篇讲的进程内存隔离,在这里立刻变成一个非常实际的问题:多进程部署下,进程内的全局变量每个 worker 各有一份。想用全局字典做缓存,每个 worker 缓存一份(内存 ×N,还可能不一致);想存登录态、Agent 会话、计数器、分布式锁这类需要全局一致的状态,必须放到 Redis 或数据库里。

worker 数量怎么定?Gunicorn 文档里那个著名的 2 × CPU 核数 + 1 公式是给同步 worker的经验值,不该原样套到异步服务上。Uvicorn worker 在一个进程里就能并发等待大量 IO;而每增加一个 worker,又会多复制一份模型、缓存、线程池和数据库连接池。因此不存在通用答案,更稳妥的做法是:

  1. 从 1 个或接近 CPU 核数的较小值开始。
  2. 用真实流量模型压测,同时观察 CPU、内存、事件循环延迟、P95/P99 和下游限流。
  3. 逐个增加 worker;吞吐不再上涨或尾延迟开始恶化时就停。
  4. 记得核算总资源,例如每进程 20 个数据库连接、4 个 worker 就可能建立 80 个连接。

在 Kubernetes 这类编排环境里,通常还会选择“一个容器一个 Uvicorn 进程”,把副本数交给编排器管理;单机部署再考虑 --workers。FastAPI 的多进程部署文档也专门区分了这两类场景。

至此,三种并发手段在完整部署里的位置全部清楚了:

text
                nginx(反向代理)
                    │
          ┌─────────┴─────────┐
          ▼                   ▼
    worker 进程 1        worker 进程 2        <- 进程:吃满多核 CPU
    ┌──────────────┐     ┌──────────────┐
    │  事件循环      │     │  事件循环      │  <- 协程:单线程挂成千上万连接
    │  协程 协程 协程 │     │  协程 协程 协程 │
    │  线程池 x40   │     │  线程池 x40   │  <- 线程:跑 def 函数和同步库
    └──────────────┘     └──────────────┘
  • 进程负责利用多核:GIL 之下一个进程只能用一个核跑字节码,多开 worker 才能把机器吃满
  • 协程负责扛住连接:等待 IO 的时间全部让出来服务别的请求,单线程挂住成千上万并发
  • 线程负责兼容同步世界:def 路径函数、只有同步版的第三方库,在线程池里等着 GIL 释放后并行 IO

为什么 Agent 服务特别爱用 FastAPI

先把说法收窄一点:不是所有 Agent 都需要 FastAPI,FastAPI 也不是 Agent 框架。Agent 的核心是模型调用、工具编排、状态和记忆;FastAPI 解决的是怎么把这些能力稳定地暴露成 HTTP API。它更像 Agent 外面那层服务外壳。

但这层外壳确实经常是 FastAPI,因为一次典型的 Agent 运行天然长这样:

text
用户请求
  -> 等 LLM 首轮推理
  -> 等向量库召回
  -> 等一个或多个 HTTP / 数据库工具
  -> 再等 LLM 汇总
  -> 持续把 token / tool event 推给前端

普通 CRUD 接口可能只等一次数据库,Agent 一次运行却会串起多轮外部 IO。真正执行 Python 的时间很短,大部分时间都花在“等”。这恰好是 ASGI + 协程最擅长的负载:一个 Agent 在等模型时,事件循环可以继续推进其他 Agent,而不需要为每个会话常驻一个线程。

具体说,FastAPI 和 Agent 的需求有六个对位点:

Agent 服务需要什么FastAPI 提供什么工程上的价值
同时等待 LLM、向量库和工具ASGI、async def、异步依赖一个进程可以挂住大量等待中的运行
边生成边展示 token 和步骤StreamingResponse、SSE、WebSocket不必等完整答案生成后才响应
严格定义输入、输出和工具参数Pydantic + OpenAPI参数校验、JSON Schema、交互文档一次得到
复用模型客户端、连接池和图实例依赖注入 + lifespan启动时创建,所有请求复用,退出时统一释放
直接调用 AI/数据生态Python 运行时SDK、编排框架、数据处理代码不用跨语言封装
暂时接入只有同步版的 SDKdef 线程池、asyncio.to_thread可以渐进迁移,不必先把所有依赖异步化

所以大家选择 FastAPI,通常不是因为某个单项性能数字,而是因为它让 Python AI 代码 -> 类型化 API -> 流式交互 -> ASGI 部署 这条路径很短。十几行就能把原型暴露出去,后面又有足够的机制补上鉴权、限流、监控和多进程部署。

一个能拿去改的 Agent 接口层

下面这段不是 Agent 算法,而是一层接近生产的 API 骨架。假设真正的 Agent 运行时提供普通和 SSE 流式两个接口;无论后面接的是编排框架、自研循环还是独立模型服务,这层并发护栏都一样:

python
import asyncio
import json
import os
from collections.abc import AsyncIterator
from contextlib import asynccontextmanager

import httpx
from fastapi import FastAPI, HTTPException, Request
from fastapi.responses import StreamingResponse
from pydantic import BaseModel, Field

MAX_AGENT_RUNS = 20
TOTAL_TIMEOUT_SECONDS = 60


class AgentInput(BaseModel):
    session_id: str = Field(min_length=1, max_length=128)
    message: str = Field(min_length=1, max_length=20_000)


@asynccontextmanager
async def lifespan(app: FastAPI):
    # 一个进程只创建一个客户端,复用连接池;不要每个请求 new 一个。
    async with httpx.AsyncClient(
        base_url=os.environ["AGENT_RUNTIME_URL"],
        timeout=httpx.Timeout(55.0, connect=5.0),
        limits=httpx.Limits(
            max_connections=100,
            max_keepalive_connections=20,
        ),
    ) as client:
        app.state.agent_http = client
        # Semaphore 是进程内限制:4 个 worker 的总上限是 4 × 20。
        app.state.agent_slots = asyncio.Semaphore(MAX_AGENT_RUNS)
        yield


app = FastAPI(lifespan=lifespan)


@app.post("/agents/run")
async def run_agent(body: AgentInput, request: Request):
    client: httpx.AsyncClient = request.app.state.agent_http
    slots: asyncio.Semaphore = request.app.state.agent_slots

    try:
        # 总超时同时覆盖排队时间和真正执行时间。
        async with asyncio.timeout(TOTAL_TIMEOUT_SECONDS):
            async with slots:
                response = await client.post(
                    "/v1/runs",
                    json=body.model_dump(),
                )
                response.raise_for_status()
                return response.json()
    except TimeoutError as exc:
        raise HTTPException(504, "agent run timed out") from exc
    except httpx.HTTPError as exc:
        raise HTTPException(502, "agent runtime unavailable") from exc


def sse_error(code: str) -> bytes:
    data = json.dumps({"type": "error", "code": code})
    return f"event: error\ndata: {data}\n\n".encode()


@app.post("/agents/run/stream")
async def stream_agent(body: AgentInput, request: Request):
    client: httpx.AsyncClient = request.app.state.agent_http
    slots: asyncio.Semaphore = request.app.state.agent_slots

    async def events() -> AsyncIterator[bytes]:
        try:
            async with asyncio.timeout(TOTAL_TIMEOUT_SECONDS):
                async with slots:
                    # 上游返回已经按 SSE 格式编码的数据,这里直接透传。
                    async with client.stream(
                        "POST",
                        "/v1/runs/stream",
                        json=body.model_dump(),
                    ) as upstream:
                        upstream.raise_for_status()
                        async for chunk in upstream.aiter_bytes():
                            if await request.is_disconnected():
                                break
                            yield chunk
        except TimeoutError:
            # 流开始后不能再把 HTTP 状态码改成 504,只能发错误事件。
            yield sse_error("timeout")
        except httpx.HTTPError:
            yield sse_error("upstream_unavailable")

    return StreamingResponse(
        events(),
        media_type="text/event-stream",
        headers={
            "Cache-Control": "no-cache",
            "X-Accel-Buffering": "no",
        },
    )

这段代码比一个“能跑”的 demo 多做了四件真正影响稳定性的事:

  1. 连接复用:在 lifespan 里创建一次 AsyncClient,避免每次请求重新握手、建连接。
  2. 并发有界Semaphore 不让瞬时流量无限压向模型服务。异步提高的是等待效率,不会提高供应商配额,也不会凭空增加 GPU 显存。
  3. 全链路超时:60 秒既包含等待并发槽位,也包含上游运行,避免请求永远挂着。对会产生副作用的 Agent 运行不要盲目自动重试,应使用幂等键或显式恢复机制。
  4. 正确处理流StreamingResponse 逐块透传 SSE;客户端断开后停止读取上游。流已经开始时不能再修改 HTTP 状态,只能在协议内发送错误事件。

这才是“异步 Agent API”的完整含义。只把函数签名从 def 改成 async def,却没有异步 SDK、并发上限和超时,离能上线还很远。

Agent 场景最容易踩的六个坑

直接后果更合适的做法
每个请求创建一个 LLM/HTTP 客户端连接无法复用,延迟和文件描述符一起上涨lifespan 创建共享客户端,退出时关闭
认为 async 等于无限并发先撞供应商 429、连接池或 GPU OOM每进程 Semaphore,跨进程再用网关或 Redis 限流
把会话和记忆放全局字典多 worker 数据不一致,重启全丢Redis/数据库保存状态,以 session_id 定位
用 BackgroundTasks 跑几分钟的 Agent重启丢任务,无法重试,也无法查询进度持久化任务队列,先返回 run_id,再查状态或订阅事件
在事件循环里跑本地推理、PDF 解析CPU/GPU 调用阻塞整个 worker独立推理服务;CPU 重活交给任务队列或进程池
让模型任意选择 URL、Shell 或数据库操作SSRF、命令执行、越权和数据破坏工具白名单、参数校验、最小权限、沙箱和人工审批

流式交互也要按需求选协议:只有“服务端不断推 token/步骤”时,SSE 更简单,协议还定义了事件 ID 和重连线索;但像上面这样的 POST 流,客户端仍要自己保存最后事件 ID,并由服务端实现断点续传。如果还要在运行中双向发送打断、确认、语音帧或 Human-in-the-loop 指令,再上 WebSocket。不要因为 Agent 听起来复杂,就默认选择更复杂的协议。

FastAPI 不会替 Agent 解决什么

FastAPI 不会让本地模型推理更快,不会替你管理 GPU,也不会自动保证任务可靠、会话一致和工具安全。生产里常见的拆法反而是:

text
FastAPI API 层
  ├─ 鉴权、校验、限流、SSE / WebSocket
  ├─ Redis / PostgreSQL:会话、记忆、运行状态
  ├─ 任务队列:长任务、重试、恢复
  └─ Agent / 模型运行时:独立进程或独立服务

本地大模型推理更适合交给专门的推理服务;数分钟的深度研究 Agent 更适合“提交任务 -> 返回 run_id -> 查询或订阅结果”;纯内部高吞吐 RPC 也可能更适合 gRPC。理解这些边界,才是真正理解了为什么选 FastAPI:它擅长 Agent 的服务接口,不等于它包办 Agent 的全部运行时。

小结

FastAPI 的并发设计可以浓缩成四句话:

  1. 一个请求就是一个协程,事件循环靠 await 点在等待 IO 的间隙服务其他请求——这是高并发的来源。
  2. async defdef 是两条路:前者直上事件循环,后者进线程池。在 async def 里写阻塞调用,堵的是全局唯一的事件循环,全站连坐——这是事故的头号来源。
  3. 生产部署是三层配合:多 worker 进程利用多核,进程内协程扛 IO 并发,线程池兜底同步代码。进程间内存隔离,全局状态必须外置到 Redis/DB。
  4. FastAPI 常被用作 Agent 的服务外壳,是因为 Agent 天生有大量链式 IO、流式输出和类型化工具参数;不是因为 FastAPI 能加速模型推理。上线时还要补齐连接复用、并发上限、超时、持久化任务和工具安全。

把上一篇的机理和这一篇的落地对照着看,FastAPI 文档里那些“建议”背后的原因就全都说得通了:为什么推荐异步库、为什么同步代码要用 def、为什么要多个 worker,以及为什么 Agent 项目尤其容易从它起步。框架替你做的每一个默认行为,底下都是进程、线程、协程这三样基本功;而真正的工程能力,体现在你是否知道默认行为的边界。


系列导航

上一篇:进程、线程、协程,再到 Python 的 async/await:并发编程一次讲透

分享
← 返回博客列表
🎁 有邀请福利哦,点击查看
🎁