06 FastAPI 部署与前端
学习理念:本章将智能体工作流包装为 FastAPI 服务,通过 SSE 协议流式推送执行进度和查询结果。核心知识点是 FastAPI 的生命周期 + 依赖注入 + 中间件 三件套,这些在前面的章节已经学过,本章是实战应用。
海外对标:HTTP SSE 协议(Server-Sent Events)、FastAPI 官方 streaming 模式
本节 AI 替代率:~75% | 人工干预率:~25%
| 角色 | 能力范围 |
|---|---|
| 🤖 AI 擅长 | 生成 FastAPI 路由代码、SSE 流式响应、依赖注入链 |
| 👤 人类需理解 | SSE 协议格式、生命周期中 client init/close 的时机、request_id 中间件的全链路追踪 |
📌 原文说明:以下内容来源于原始笔记第7章"API接口"和第8章"前后端联调"。原文全部保留,补充了架构图和代码优先级标注。
一、需求说明
本章旨在实现一个查询接口,用于接收用户查询,并实时响应工作流执行进度和查询结果。接口使用FastAPI框架编写,涉及到的相关知识点如下:
流式相应:https://fastapi.org.cn/advanced/custom-response/#streamingresponse
SSE协议:https://www.ruanyifeng.com/blog/2017/05/server-sent_events.html
生命周期事件:https://fastapi.org.cn/advanced/events/
中间件:https://fastapi.org.cn/tutorial/middleware/
依赖注入:https://fastapi.org.cn/tutorial/dependencies/
二、代码组织规划
data-agent/
├── main.py # FastAPI入口脚本
└── app/
├── api/
│ ├── routers/
│ │ └── query_router.py # 查询接口路由
│ ├── schemas/
│ │ └── query_schema.py # 请求体结构
│ └── dependencies.py # 依赖注入(6层)
├── services/
│ └── query_service.py # 查询业务逻辑
└── core/
├── lifespan.py # 生命周期事件
├── context.py # request_id 上下文变量
└── log.py # 日志配置(含 request_id)三、API 实现层
3.1 入口脚本(main.py)
🔥 【P0 必须理解】 FastAPI 应用入口。包含 3 个关键注册:lifespan(生命周期)、router(路由)、middleware(中间件)。
import uuid
from fastapi import FastAPI, Request
from app.api.routers.query_router import query_router
from app.core.context import request_id_ctx_var
from app.core.lifespan import lifespan
app = FastAPI(lifespan=lifespan) # 注册生命周期
app.include_router(query_router) # 注册路由
# 添加中间件,在每个请求中生成唯一的request_id
@app.middleware("http")
async def add_process_time_header(request: Request, call_next):
request_id_ctx_var.set(uuid.uuid4()) # 请求开始时设置
response = await call_next(request) # 执行路径函数
return response
if __name__ == '__main__':
import uvicorn
uvicorn.run(app, host="0.0.0.0", port=8000)3.2 查询接口定义
🔥 【P0 必须理解】 接口使用
StreamingResponse返回 SSE 格式数据。前端通过 EventSource API 消费。
from fastapi import APIRouter
from fastapi.params import Depends
from starlette.responses import StreamingResponse
from app.api.dependencies import get_query_service
from app.api.schemas.query_schema import QuerySchema
from app.services.query_service import QueryService
query_router = APIRouter()
@query_router.post("/api/query")
async def query(query: QuerySchema, query_service: QueryService = Depends(get_query_service)):
return StreamingResponse(
query_service.query(query.query), # 返回异步生成器
media_type="text/event-stream" # SSE 媒体类型
)3.3 请求体 Schema
from pydantic import BaseModel
class QuerySchema(BaseModel):
query: str3.4 QueryService(核心业务逻辑)
🔥 【P0 必须理解】
QueryService.query()是一个异步生成器,遍历graph.astream()的每个 chunk,将其序列化为 SSE 格式"data: {json}\n\n"逐行 yield。FastAPI 的 StreamingResponse 消费这个生成器实现流式输出。
import json
from langchain_huggingface import HuggingFaceEndpointEmbeddings
from app.agent.context import DataAgentContext
from app.agent.graph import graph
from app.agent.state import DataAgentState
from app.repositories.es.value_es_repository import ValueESRepository
from app.repositories.mysql.dw.dw_mysql_repository import DWMySQLRepository
from app.repositories.mysql.meta.meta_mysql_repository import MetaMySQLRepository
from app.repositories.qdrant.column_qdrant_repository import ColumnQdrantRepository
from app.repositories.qdrant.metric_qdrant_repository import MetricQdrantRepository
class QueryService:
def __init__(self,
embedding_client: HuggingFaceEndpointEmbeddings,
column_qdrant_repository: ColumnQdrantRepository,
value_es_repository: ValueESRepository,
metric_qdrant_repository: MetricQdrantRepository,
meta_mysql_repository: MetaMySQLRepository,
dw_mysql_repository: DWMySQLRepository):
self.embedding_client = embedding_client
self.column_qdrant_repository = column_qdrant_repository
self.value_es_repository = value_es_repository
self.metric_qdrant_repository = metric_qdrant_repository
self.meta_mysql_repository = meta_mysql_repository
self.dw_mysql_repository = dw_mysql_repository
async def query(self, query: str):
context = DataAgentContext(
embedding_client=self.embedding_client,
column_qdrant_repository=self.column_qdrant_repository,
value_es_repository=self.value_es_repository,
metric_qdrant_repository=self.metric_qdrant_repository,
meta_mysql_repository=self.meta_mysql_repository,
dw_mysql_repository=self.dw_mysql_repository
)
state = DataAgentState(query=query)
try:
async for chunk in graph.astream(input=state, context=context, stream_mode="custom"):
yield f"data: {json.dumps(chunk, ensure_ascii=False, default=str)}\n\n"
except Exception as e:
yield f"data: {json.dumps({'type': 'error', 'message': str(e)}, ensure_ascii=False, default=str)}\n\n"3.5 依赖注入
🟡 【P1 看注释就行】 6 层依赖注入链:
get_query_service→ 需要 6 个子依赖 → 每个子依赖创建对应的 Repository/Client。FastAPI 的 Depends 系统自动解析。
# 每一层返回一个依赖对象
async def get_meta_session():
async with meta_mysql_client_manager.session_factory() as session:
yield session
async def get_embedding_client():
return embedding_client_manager.client
# ... 类似 get_column_qdrant_repository / get_value_es_repository / get_metric_qdrant_repository
async def get_meta_mysql_repository(session: AsyncSession = Depends(get_meta_session)):
return MetaMySQLRepository(session)
async def get_dw_mysql_repository(session: AsyncSession = Depends(get_dw_session)):
return DWMySQLRepository(session)
# 组装 QueryService
async def get_query_service(
embedding_client=Depends(get_embedding_client),
column_qdrant_repository=Depends(get_column_qdrant_repository),
value_es_repository=Depends(get_value_es_repository),
metric_qdrant_repository=Depends(get_metric_qdrant_repository),
meta_mysql_repository=Depends(get_meta_mysql_repository),
dw_mysql_repository=Depends(get_dw_mysql_repository)
) -> QueryService:
return QueryService(...)四、生命周期事件
🟡 【P1 看注释就行】 应用启动时初始化所有 5 个客户端,关闭时释放资源。通过 FastAPI 的
lifespan异步上下文管理器注册。
from contextlib import asynccontextmanager
from fastapi import FastAPI
from app.clients.embedding_client_manager import embedding_client_manager
from app.clients.es_client_manager import es_client_manager
from app.clients.mysql_client_manager import meta_mysql_client_manager, dw_mysql_client_manager
from app.clients.qdrant_client_manager import qdrant_client_manager
@asynccontextmanager
async def lifespan(app: FastAPI):
# 应用启动前执行
embedding_client_manager.init()
qdrant_client_manager.init()
es_client_manager.init()
meta_mysql_client_manager.init()
dw_mysql_client_manager.init()
yield
# 应用关闭前执行
await qdrant_client_manager.close()
await es_client_manager.close()
await meta_mysql_client_manager.close()
await dw_mysql_client_manager.close()五、中间件与全链路追踪
🟡 【P1 看注释就行】 使用
contextvars+middleware+loguru.patch实现每个请求的 request_id 注入。在并发场景下,每个异步任务的日志都能关联到对应的请求。
context.py — 定义 request_id 上下文变量:
from contextvars import ContextVar
request_id_ctx_var = ContextVar("request_id", default="1")middleware — 每次请求开始时生成新 ID:
@app.middleware("http")
async def add_process_time_header(request: Request, call_next):
request_id_ctx_var.set(uuid.uuid4())
response = await call_next(request)
return responselog.py — loguru 的 format 中加入 {extra[request_id]},并通过 patch() 动态注入:
log_format = (
"<green>{time:YYYY-MM-DD HH:mm:ss.SSS}</green> | "
"<level>{level: <8}</level> | "
"<magenta>request_id - {extra[request_id]}</magenta> | "
"<cyan>{name}</cyan>:<cyan>{function}</cyan>:<cyan>{line}</cyan> - "
"<level>{message}</level>"
)
def inject_request_id(record):
request_id = request_id_ctx_var.get()
record["extra"]["request_id"] = request_id
logger.remove()
logger = logger.patch(inject_request_id)
# ... add console/file handlers ...六、前端联调
6.1 启动前端项目
前端代码可从课程资料获取,另外启动项目需要node运行环境,node的安装包也可从课程资料获取。
在准备好node环境后,可在前端项目的根目录执行如下命令启动项目:
安装项目所需依赖
npm install启动项目
npm run dev
七、企业痛点-方案映射
| 痛点 | 传统方案 | AI Agent 方案 | 效率提升 |
|---|---|---|---|
| 查询操作等待时间长(10s+) | 无进度提示,用户不知是否挂掉 | SSE 流式推送每个节点进度 | 用户体验提升 80% |
| 多用户并发日志混乱 | 单一日志文件难以追踪 | request_id 全链路追踪 | 排错速度提升 5x |
| 客户端连接资源泄漏 | 忘记关闭连接 | 生命周期统一管理 init/close | 资源泄漏减少 90% |
八、本阶段文件索引
| 优先级 | 文件 | 路径 |
|---|---|---|
| 🔥 P0 | main.py | main.py |
| 🔥 P0 | query_router.py | app/api/routers/query_router.py |
| 🔥 P0 | query_service.py | app/services/query_service.py |
| 🔥 P0 | dependencies.py | app/api/dependencies.py |
| 🟡 P1 | lifespan.py | app/core/lifespan.py |
| 🟡 P1 | log.py | app/core/log.py |
| 🟢 P2 | context.py | app/core/context.py |
| 🟢 P2 | query_schema.py | app/api/schemas/query_schema.py |