Skip to content

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/


二、代码组织规划

bash
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(中间件)。

python
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 消费。

python
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

python
from pydantic import BaseModel

class QuerySchema(BaseModel):
    query: str

3.4 QueryService(核心业务逻辑)

🔥 【P0 必须理解】 QueryService.query() 是一个异步生成器,遍历 graph.astream() 的每个 chunk,将其序列化为 SSE 格式 "data: {json}\n\n" 逐行 yield。FastAPI 的 StreamingResponse 消费这个生成器实现流式输出。

python
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 系统自动解析。

python
# 每一层返回一个依赖对象
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 异步上下文管理器注册。

python
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 上下文变量:

python
from contextvars import ContextVar
request_id_ctx_var = ContextVar("request_id", default="1")

middleware — 每次请求开始时生成新 ID:

python
@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

log.py — loguru 的 format 中加入 {extra[request_id]},并通过 patch() 动态注入:

python
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环境后,可在前端项目的根目录执行如下命令启动项目:

安装项目所需依赖

bash
npm install

启动项目

bash
npm run dev

img


七、企业痛点-方案映射

痛点传统方案AI Agent 方案效率提升
查询操作等待时间长(10s+)无进度提示,用户不知是否挂掉SSE 流式推送每个节点进度用户体验提升 80%
多用户并发日志混乱单一日志文件难以追踪request_id 全链路追踪排错速度提升 5x
客户端连接资源泄漏忘记关闭连接生命周期统一管理 init/close资源泄漏减少 90%

八、本阶段文件索引

优先级文件路径
🔥 P0main.pymain.py
🔥 P0query_router.pyapp/api/routers/query_router.py
🔥 P0query_service.pyapp/services/query_service.py
🔥 P0dependencies.pyapp/api/dependencies.py
🟡 P1lifespan.pyapp/core/lifespan.py
🟡 P1log.pyapp/core/log.py
🟢 P2context.pyapp/core/context.py
🟢 P2query_schema.pyapp/api/schemas/query_schema.py

OPC 超级个体实战指南