1. 请求
FastAPI 建立在 ASGI 之上,可以在一个服务进程中并发处理多个请求。对 async def 路径函数来说,每个正在处理的请求可以理解为一条由事件循环推进的异步调用链。
1import asyncio2from fastapi import FastAPI34app = FastAPI()56@app.get("/health")7async def health():8await asyncio.sleep(0.1)9return {"status": "ok"}
请求运行到 await asyncio.sleep(0.1) 时,当前请求暂停,事件循环可以推进同一个 worker 中的其他请求。等定时器完成后,这个请求再获得后续执行机会。
真实接口等待的通常不是 sleep(),而是数据库、HTTP 服务、模型供应商或文件流。它们共同的特点是:等待外部 IO 时不必一直占用 CPU。
一个 worker 中可能同时存在:
- 正在事件循环上运行的请求 Task;
- 等待网络或数据库 Future 的请求;
- 被 FastAPI 放到线程池中的同步路径函数和依赖;
- 应用生命周期中创建的连接池与后台 Task。
事件循环并不会让所有 Python 代码同时执行。一个 async def 路径函数从上一个 await 恢复后,会在事件循环线程上继续执行;如果其中包含耗时同步代码,同一 worker 的其他异步请求也会被拖延。
2. 接口
FastAPI 同时支持 async def 和普通 def 路径函数,但两者的执行位置不同。
异步接口
使用原生异步库时,路径函数应该声明为 async def:
1from fastapi import FastAPI23app = FastAPI()45@app.get("/users/{user_id}")6async def get_user(user_id: int):7user = await async_user_repository.get(user_id)8return user
await 只在操作尚未完成时暂停当前请求,并不会阻塞事件循环线程。
同步接口
只能使用同步阻塞库时,可以把整个路径函数声明为普通 def:
1import time2from fastapi import FastAPI34app = FastAPI()56@app.get("/reports/{report_id}")7def get_report(report_id: int):8time.sleep(1)9return {"id": report_id, "status": "done"}
FastAPI 不会在事件循环线程中直接调用这个普通函数,而是把它交给外部线程池,再异步等待结果。因此,一个同步接口阻塞时,事件循环仍能处理其他请求,但线程池容量有限,超出容量的同步工作会排队。
两者可以这样选择:
| 路径函数 | 内部调用 | 执行方式 |
|---|---|---|
async def | 原生异步 API | 在事件循环中执行,并在 await 处暂停 |
async def | 直接调用同步阻塞函数 | 阻塞事件循环,错误用法 |
def | 同步阻塞 API | FastAPI 在线程池中调用 |
def | 想要调用异步 API | 不能直接使用 await,应重新设计为异步路径 |
最危险的写法,是在异步接口中直接调用同步阻塞函数:
01import time02from fastapi import FastAPI0304app = FastAPI()0506def load_report():07time.sleep(2)08return {"status": "done"}0910@app.get("/report")11async def get_report():12return load_report() # 阻塞事件循环
FastAPI 只会自动处理它负责调用的路径函数和依赖。路径函数内部直接调用的普通工具函数,不会因为运行在 FastAPI 中就自动进入线程池。
暂时无法替换同步工具时,可以显式使用 asyncio.to_thread():
01import asyncio02import time03from fastapi import FastAPI0405app = FastAPI()0607def load_report():08time.sleep(2)09return {"status": "done"}1011@app.get("/report")12async def get_report():13return await asyncio.to_thread(load_report)
不要在 FastAPI 的异步路径函数中调用 asyncio.run()。FastAPI 已经运行事件循环,应用代码只需要 await 协程;再次启动顶层事件循环会抛出错误。
3. 资源
异步服务通常需要长期复用 HTTP 客户端、数据库连接池和模型客户端。每个请求都重新创建客户端,会失去连接复用,还会增加握手和资源分配成本。
FastAPI 推荐使用 lifespan 管理应用级资源。yield 之前负责初始化,yield 之后负责关闭:
01import asyncio02from contextlib import asynccontextmanager0304import httpx05from fastapi import FastAPI, Request0607@asynccontextmanager08async def lifespan(app: FastAPI):09limits = httpx.Limits(10max_connections=100,11max_keepalive_connections=20,12)13timeout = httpx.Timeout(10, connect=3)1415async with httpx.AsyncClient(16limits=limits,17timeout=timeout,18) as client:19app.state.http_client = client20app.state.model_limit = asyncio.Semaphore(8)21yield2223app = FastAPI(lifespan=lifespan)2425@app.get("/users/{user_id}")26async def get_user(user_id: int, request: Request):27client = request.app.state.http_client28response = await client.get(29f"https://example.com/users/{user_id}"30)31response.raise_for_status()32return response.json()
应用开始接收请求前,FastAPI 会进入 lifespan;应用关闭时,它会退出 async with 并关闭 HTTP 客户端。这样既能复用连接池,也能确保资源得到清理。
需要注意几个边界:
app.state中的对象属于当前 worker,不会跨进程共享;- 四个 worker 会分别创建四个客户端和四个连接池;
max_connections=100是每个客户端实例的上限,不是整个服务的全局上限;- 连接池容量和业务 Semaphore 是两层不同限制。
不应该只把客户端放在模块顶部就不再管理。能够被显式关闭的网络、数据库和模型资源,应拥有与应用生命周期对应的创建和清理位置。
4. 依赖
FastAPI 依赖也可以是异步函数。需要为每个请求获取和释放资源时,可以使用带 yield 的异步依赖:
01from typing import Annotated0203from fastapi import Depends, FastAPI04from sqlalchemy.ext.asyncio import AsyncSession0506from db import User, session_factory0708app = FastAPI()0910async def get_session():11async with session_factory() as session:12yield session1314Session = Annotated[AsyncSession, Depends(get_session)]1516@app.get("/users/{user_id}")17async def get_user(user_id: int, session: Session):18return await session.get(User, user_id)
进入依赖时创建或取出 Session,yield 把它交给路径函数,请求依赖生命周期结束后再退出上下文并清理。连接池属于应用级资源,Session 通常属于请求级资源,这两个生命周期不应混在一起。
普通 def 依赖会像同步路径函数一样在线程池中执行。异步依赖直接运行在异步调用链中,因此其中也不能直接执行耗时同步 IO。
依赖可以负责资源获取、认证和请求范围上下文,但不要让一个依赖偷偷创建无法追踪的长期后台 Task。资源所有权越清楚,超时、取消和关闭时越容易正确清理。
5. LangChain
LangChain Runnable 通常同时提供同步和异步接口:
1async def ask(chain, question):2return await chain.ainvoke({"question": question})34async def stream_answer(chain, question):5async for chunk in chain.astream({"question": question}):6yield chunk
ainvoke() 返回一个可等待结果,astream() 返回异步迭代器。它们可以直接接入 FastAPI 的异步路径函数。
但是,方法名带有 a 只表示提供异步调用方式,不保证每个集成都使用原生异步网络客户端。LangChain Runnable 的默认异步实现可以在线程池中调用同步版本;具体模型、Retriever 或 Tool 也可以覆盖它,提供真正的原生异步实现。
因此需要区分两层:
| 层级 | 要确认的问题 |
|---|---|
| LangChain API | 是否提供 ainvoke()、astream() 等异步入口 |
| 底层集成 | 是原生异步 IO,还是在线程池中兼容同步方法 |
对于常用模型调用,直接使用集成提供的 ainvoke(),不要再额外套一层 to_thread()。如果性能和线程池容量很重要,应查看对应集成的实现与文档,并通过压测确认行为。
普通响应
01from fastapi import FastAPI, Request02from pydantic import BaseModel0304app = FastAPI()0506class ChatInput(BaseModel):07question: str0809@app.post("/chat")10async def chat(payload: ChatInput, request: Request):11chain = request.app.state.chain12answer = await chain.ainvoke(13{"question": payload.question}14)15return {"answer": answer}
这里假设 chain 的最后一步已经使用输出解析器返回字符串。如果返回的是 AIMessage,应根据消息结构读取 content,不要假设所有 Runnable 都返回同一种类型。
流式响应
01from fastapi import FastAPI, Request02from fastapi.responses import StreamingResponse0304app = FastAPI()0506@app.get("/chat/stream")07async def stream_chat(question: str, request: Request):08chain = request.app.state.chain0910async def generate():11async for chunk in chain.astream(12{"question": question}13):14yield chunk1516return StreamingResponse(17generate(),18media_type="text/plain; charset=utf-8",19)
这个例子同样假设 chain 经过字符串输出解析器,每个 chunk 都是字符串。如果模型返回消息块或内容块,需要先转换成协议要求的文本或 SSE 数据。
astream() 的默认实现可能只是等待 ainvoke() 后一次产出结果。只有链条中的具体组件实现了流式能力,客户端才会真正逐块收到模型输出。接口可以异步迭代,不等于底层一定是逐 Token 流式传输。
6. 并发
一个请求内部也可以并发执行互不依赖的操作。例如,用户画像和长期记忆可以同时查询,等两者完成后再调用模型:
01import asyncio0203async def build_answer(04question,05profile_store,06memory_store,07model,08):09async with asyncio.TaskGroup() as group:10profile_task = group.create_task(11profile_store.ainvoke(question)12)13memory_task = group.create_task(14memory_store.ainvoke(question)15)1617context = {18"question": question,19"profile": profile_task.result(),20"memory": memory_task.result(),21}2223return await model.ainvoke(context)
TaskGroup 保证两个子任务都在当前请求的结构中完成。一个查询失败后,它会取消并等待另一个查询,再把异常交给调用方。
模型生成依赖画像和记忆结果,所以必须放在它们之后。创建 Task 只能让相互独立的等待重叠,不能消除真实的数据依赖。
并发上限
如果每个请求都创建多个上游调用,服务级并发很容易膨胀。可以在 lifespan 中创建 Semaphore,并让所有请求共享当前 worker 的并发额度:
1from fastapi import Request23async def invoke_model(payload, request: Request):4limit = request.app.state.model_limit5model = request.app.state.model67async with limit:8return await model.ainvoke(payload)
Semaphore 限制的是当前进程同时进入模型调用的数量,不是每分钟请求次数。供应商的 RPM、TPM 等频率限制还需要单独的限流和重试策略。
四个 worker 每个都创建 Semaphore(8) 时,整个部署最多可能同时进入 32 个模型调用。需要全局限制时,应在网关、共享限流器或任务调度层统一协调。
7. 超时
对外部服务的等待必须有时间边界。Python 3.11 之后,可以用 asyncio.timeout() 限制一段异步调用链:
01import asyncio0203from fastapi import FastAPI, HTTPException, Request04from pydantic import BaseModel0506app = FastAPI()0708class ChatInput(BaseModel):09question: str1011@app.post("/chat")12async def chat(payload: ChatInput, request: Request):13chain = request.app.state.chain1415try:16async with asyncio.timeout(30):17answer = await chain.ainvoke(18{"question": payload.question}19)20except TimeoutError as error:21raise HTTPException(22status_code=504,23detail="模型响应超时",24) from error2526return {"answer": answer}
超时发生时,asyncio.timeout() 会通过取消当前异步操作来停止等待,并在代码块外转换为 TimeoutError。底层异步库需要正确响应取消并释放连接。
不要使用宽泛的 except BaseException 吞掉 asyncio.CancelledError。FastAPI 服务关闭、父任务取消或结构化并发清理都依赖取消传播;如果确实捕获取消,应完成必要清理后继续抛出。
还要区分应用超时与客户端超时:
- HTTP 客户端的连接、读取和写入超时负责限制单次网络阶段;
asyncio.timeout()可以限制整条业务流程;- 反向代理和浏览器还可能有自己的请求超时。
多层超时应该由内向外逐步留出清理时间,而不是全部设置成同一个数字。
如果 LangChain 的异步接口实际在线程中执行同步方法,取消等待并不能强制终止已经运行的线程函数。底层同步客户端仍然必须配置自身超时。
8. 后台
FastAPI 的 BackgroundTasks 可以在响应发送后执行一个小型进程内工作:
01from fastapi import BackgroundTasks, FastAPI, status0203app = FastAPI()0405def write_log(message):06with open("events.log", "a", encoding="utf-8") as file:07file.write(f"{message}\n")0809@app.post("/events", status_code=status.HTTP_202_ACCEPTED)10async def create_event(11message: str,12background_tasks: BackgroundTasks,13):14background_tasks.add_task(write_log, message)15return {"accepted": True}
后台函数既可以是 def,也可以是 async def。同步后台函数会使用线程池,异步后台函数则运行在异步环境中;异步函数内部仍然不能执行阻塞代码。
BackgroundTasks 的「后台」只表示响应已经可以先返回,不表示工作进入了可靠任务系统。它仍然属于当前 FastAPI 进程:进程崩溃、部署重启或 worker 被终止时,任务可能丢失。
| 工作 | 合适方式 |
|---|---|
| 轻量日志、非关键通知 | BackgroundTasks |
| 必须成功、需要重试的任务 | 外部任务队列 |
| CPU 密集型或长时间任务 | 独立 worker 或计算服务 |
| 需要跨服务器统一调度 | 持久化队列与调度器 |
也不要在路径函数中随手调用 asyncio.create_task() 后丢掉引用。这样的任务缺少清晰所有者,异常可能无人处理,应用关闭时也难以统一等待和取消。
9. 生命周期
确实需要每个服务进程维护一个周期任务时,可以把 Task 放进 lifespan,并在关闭阶段取消和等待:
01import asyncio02from contextlib import asynccontextmanager0304from fastapi import FastAPI0506async def refresh_cache():07try:08while True:09print("刷新当前进程缓存")10await asyncio.sleep(60)11finally:12print("停止缓存刷新")1314@asynccontextmanager15async def lifespan(app: FastAPI):16task = asyncio.create_task(17refresh_cache(),18name="refresh-cache",19)2021try:22yield23finally:24task.cancel()2526try:27await task28except asyncio.CancelledError:29pass3031app = FastAPI(lifespan=lifespan)
这种 Task 适合维护每个进程自己的缓存或心跳。四个 worker 会分别执行一次 lifespan,也就会创建四个刷新 Task。全系统只能运行一份的定时任务,不应该放在这里。
lifespan 也是加载模型、创建连接池和启动监控资源的合适位置。初始化失败时应用不应开始接收请求;关闭阶段则应停止新工作,并在服务允许的优雅关闭时间内释放资源。
官方目前推荐使用 lifespan 统一管理启动和关闭逻辑。旧的 @app.on_event("startup") 与 @app.on_event("shutdown") 已不再推荐;设置 lifespan 后,不应再假设这两套机制会同时执行。
10. 线程池
FastAPI 基于 Starlette 和 AnyIO。同步路径函数、同步依赖、同步后台任务以及部分文件处理都会使用 AnyIO 的线程能力。
Starlette 当前默认线程限制器提供 40 个 token,也就是默认最多允许相应数量的同步工作同时占用线程额度。这个容量会被 FastAPI 同步依赖和 Starlette 内部的一些功能共同使用。
这不表示每个 FastAPI 服务永远都应该配置 40 个线程,也不表示把它调大就一定更快。同步数据库连接池只有 10 个连接时,把线程额度扩到 100 可能只会制造更多排队和内存占用。
还要注意,asyncio.to_thread() 使用 asyncio 默认 Executor,而 FastAPI/Starlette 的同步路径执行基于 AnyIO。不要把它们简单视为同一套容量配置;监控时应分别观察同步路径、手动线程任务和下游连接池。
线程池只能避免同步 IO 阻塞事件循环。普通 def 路径中的纯 Python CPU 计算仍然受默认 CPython GIL 影响,还会占用线程和 CPU。此类工作应使用进程池或独立 worker,而不是不断增加线程。
11. 状态
并发请求会访问共享资源,但不是所有共享对象都应该用同一种方式处理。
| 状态类型 | 合适位置 |
|---|---|
| 只读配置 | 应用配置对象 |
| HTTP 客户端、连接池、模型实例 | lifespan 创建的当前 worker 资源 |
| 当前请求的 Session、用户信息 | 依赖或请求上下文 |
| 当前进程缓存 | app.state,接受多 worker 不一致 |
| 库存、余额、订单 | 数据库或外部一致性存储 |
进程内可变状态仍然可能产生协程竞态:
01import asyncio0203counter = 004lock = asyncio.Lock()0506async def increase():07global counter0809async with lock:10current = counter11await asyncio.sleep(0)12counter = current + 113return counter
这个 Lock 只能保护当前事件循环中的状态。多个 worker 是多个独立进程,每个进程都有自己的 counter 和 lock。
业务一致性不能依靠进程内全局变量。库存扣减、余额变更等操作应该使用数据库事务、条件更新、唯一约束或其他经过验证的跨进程协调机制。
12. Worker
开发环境通常只有一个服务进程。生产环境可以使用多个 worker 利用多核 CPU,并增加进程级隔离:
1uv run fastapi run --workers 4 main.py
--workers 4 会启动四个服务进程。每个 worker 通常都有自己的:
- Python 解释器和事件循环;
- 线程池与同步原语;
- lifespan 资源和数据库连接池;
- 模型对象、缓存与全局变量;
- 后台 Task 和 Semaphore。
因此,worker 数量会同时放大吞吐能力和资源占用。一个模型实例占用 1 GB 内存时,四个 worker 可能加载四份;每个连接池上限为 20 时,四个 worker 可能建立 80 个连接。
worker 数量也不等于请求并发数。一个异步 worker 可以在 IO 等待期间并发推进许多请求;增加 worker 主要提供多进程并行、故障隔离和更多事件循环容量。
在 Kubernetes 等编排环境中,常见做法是每个容器运行一个服务进程,再通过多个副本扩展。是否使用容器内多 worker,要结合平台、CPU 配额、内存和连接池容量决定,不能照搬固定数字。
13. 容量
一个异步接口的真实容量由多层限制共同决定:
| 层级 | 常见限制 |
|---|---|
| HTTP 入口 | worker、连接数、反向代理并发 |
| 事件循环 | 阻塞代码、就绪 Task 数量、循环延迟 |
| 同步调用 | 线程池 token 与排队数量 |
| HTTP 客户端 | 最大连接数、Keep-Alive 数量 |
| 数据库 | 连接池大小、事务耗时 |
| 模型供应商 | 并发、RPM、TPM 和配额 |
| 应用逻辑 | Semaphore、Queue、超时和重试 |
假设 100 个请求同时进入,每个请求并发调用 3 个模型,理论上会产生 300 个上游调用。事件循环能够创建这些 Task,不表示模型服务、连接池和本机内存能够承受。
设计容量时通常需要:
- 为外部调用设置连接与读取超时;
- 使用连接池复用网络连接;
- 用 Semaphore 限制关键上游并发;
- 对大批量工作使用 Queue 建立反压;
- 只重试明确可恢复的错误,并加入退避;
- 按所有 worker 的总量核算连接和配额。
并发上限、连接池和 worker 数量应该共同压测。单独把其中一个数字调大,通常只是把排队位置移动到下一层。
14. 观测
异步服务出现变慢时,只看平均接口耗时往往不够。至少应该观察:
- 请求总耗时和分位数;
- 事件循环延迟与阻塞告警;
- 当前请求数和排队数量;
- 线程池使用量与等待时间;
- HTTP、数据库连接池的占用和等待;
- 模型调用首 Token 时间、总时间和错误率;
- 超时、取消、重试与限流次数;
- 后台 Task 异常和应用关闭耗时。
如果事件循环延迟很高,但 CPU 并没有满,通常要排查同步 IO 或锁等待;如果 CPU 持续满载,增加异步 Task 一般不会改善问题;如果线程池排队,则要确认同步依赖、文件处理和同步 SDK 是否共享了有限容量。
压测时应使用接近生产的 worker 数量、连接池、模型限额和数据规模。只在单请求下验证成功,无法说明高并发时不会耗尽资源。