一、前言
项目在智能数据问答场景中,用户往往需要实时感知查询请求的处理进度(如「召回字段」「过滤指标」「生成 SQL」等步骤),而非被动等待最终结果。本文将基于 FastAPI 框架,结合 SSE(Server-Sent Events)协议、LangGraph 工作流、依赖注入等技术,实现一个支持流式响应执行进度 + 最终结果的智能数据查询接口,完整覆盖接口设计、依赖管理、生命周期控制、日志追踪等核心环节。
本专栏最后一篇,本章旨在实现一个查询接口,用于接收用户查询,并实时响应工作流执行进度和查询结果。接口使用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/
二、技术栈与核心知识点
- Web 框架:FastAPI(高性能、原生支持异步、自动生成 OpenAPI 文档)
- 流式响应:StreamingResponse + SSE 协议(服务端单向推送数据)
- 工作流引擎:LangGraph(处理复杂的多节点数据查询流程)
- 工程化实践:依赖注入、生命周期事件、ContextVar 请求追踪、结构化日志
- 数据存储:MySQL(元数据 / 数据仓库)、Elasticsearch(字段取值检索)、Qdrant(向量检索)
三、项目结构设计
遵循「分层解耦、职责单一」的设计原则,项目结构如下:
data-agent/
├─ main.py # FastAPI入口脚本(注册路由、生命周期)
└─ app/
├─ api/ # 接口层:定义路由、请求体、依赖项
│ ├─ routers/
│ │ └─ query_router.py # 查询接口路由定义
│ ├─ schemas/
│ │ └─ query_schema.py # 请求体结构(Pydantic模型)
│ └─ dependencies.py # 依赖注入(仓库、服务实例化)
│
├─ services/ # 服务层:核心业务逻辑封装
│ └─ query_service.py # 调用LangGraph工作流,处理流式响应
│
└─ core/ # 核心配置层:生命周期、上下文、日志
├─ lifespan.py # FastAPI生命周期(客户端初始化/关闭)
├─ context.py # ContextVar存储请求ID
└─ log.py # 结构化日志(带请求ID追踪)
四、核心功能实现
4.1 入口脚本(main.py)
作为 FastAPI 应用的入口,负责注册生命周期函数、中间件、路由:
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
# 创建FastAPI应用,绑定生命周期函数(启动/关闭时初始化/销毁客户端)
app = FastAPI(lifespan=lifespan)
# 中间件:为每个请求生成唯一RequestID,用于日志追踪
@app.middleware("http")
async def add_process_time_header(request: Request, call_next):
# 生成UUID作为请求ID,存入ContextVar(异步安全)
request_id_ctx_var.set(str(uuid.uuid4()))
response = await call_next(request)
return response
# 注册查询接口路由
app.include_router(query_router)
if __name__ == "__main__":
import uvicorn
uvicorn.run("main:app", host="0.0.0.0", port=8000, reload=True)
4.2 接口层实现
4.2.1 请求体定义(query_schema.py)
使用 Pydantic 定义请求体结构,自动校验入参:
from pydantic import BaseModel
class QuerySchema(BaseModel):
"""查询接口请求体"""
query: str # 用户的自然语言查询(如「2026年Q1 GMV是多少」)
4.2.2 依赖注入(dependencies.py)
解耦服务与底层仓库,实现「依赖倒置」,便于测试和扩展:
from fastapi import Depends
from langchain_huggingface import HuggingFaceEndpointEmbeddings
from sqlalchemy.ext.asyncio import AsyncSession
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
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
from app.services.query_service import QueryService
# 1. 数据库会话依赖(异步)
async def get_meta_session():
async with meta_mysql_client_manager.session_factory() as session:
yield session
async def get_dw_session():
async with dw_mysql_client_manager.session_factory() as session:
yield session
# 2. 各类仓库依赖(封装数据访问逻辑)
async def get_embedding_client():
return embedding_client_manager.client
async def get_column_qdrant_repository():
return ColumnQdrantRepository(qdrant_client_manager.client)
async def get_value_es_repository():
return ValueESRepository(es_client_manager.client)
async def get_metric_qdrant_repository():
return MetricQdrantRepository(qdrant_client_manager.client)
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)
# 3. 核心服务依赖(组装所有仓库,供接口调用)
async def get_query_service(
embedding_client: HuggingFaceEndpointEmbeddings = Depends(get_embedding_client),
column_qdrant_repository: ColumnQdrantRepository = Depends(get_column_qdrant_repository),
value_es_repository: ValueESRepository = Depends(get_value_es_repository),
metric_qdrant_repository: MetricQdrantRepository = Depends(get_metric_qdrant_repository),
meta_mysql_repository: MetaMySQLRepository = Depends(get_meta_mysql_repository),
dw_mysql_repository: DWMySQLRepository = Depends(get_dw_mysql_repository)
) -> QueryService:
return QueryService(
embedding_client=embedding_client,
column_qdrant_repository=column_qdrant_repository,
value_es_repository=value_es_repository,
metric_qdrant_repository=metric_qdrant_repository,
meta_mysql_repository=meta_mysql_repository,
dw_mysql_repository=dw_mysql_repository
)
4.2.3 接口路由(query_router.py)
定义 POST 接口,返回流式响应(SSE 格式):
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", summary="智能数据查询接口", description="接收自然语言查询,流式返回处理进度和结果")
async def query(
query: QuerySchema, # 自动校验请求体
query_service: QueryService = Depends(get_query_service) # 注入查询服务
):
# 返回StreamingResponse,指定SSE媒体类型
return StreamingResponse(
query_service.query(query.query),
media_type="text/event-stream" # SSE协议标准媒体类型
)
4.3 服务层实现(query_service.py)
核心业务逻辑层,调用 LangGraph 工作流,将工作流的流式输出封装为 SSE 格式:
import json
from langchain_huggingface import HuggingFaceEndpointEmbeddings
from app.agent.context import DataAgentContext
from app.agent.graph import graph # 预定义的LangGraph工作流
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:
"""查询服务:封装LangGraph工作流调用逻辑"""
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):
"""
执行查询逻辑,生成流式响应
:param query: 用户自然语言查询
:return: 生成器,逐行返回SSE格式数据
"""
# 1. 构建LangGraph上下文(传递仓库/客户端实例)
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
)
# 2. 初始化LangGraph状态(传入用户查询)
state = DataAgentState(query=query)
try:
# 3. 流式执行LangGraph工作流(stream_mode="custom"自定义输出)
async for chunk in graph.astream(input=state, context=context, stream_mode="custom"):
# 4. 封装为SSE格式(data: {JSON数据}\\n\\n)
yield f"data: {json.dumps(chunk, ensure_ascii=False, default=str)}\\n\\n"
except Exception as e:
# 5. 异常处理:返回错误类型的SSE数据
yield f"data: {json.dumps({'type': 'error', 'message': str(e)}, ensure_ascii=False, default=str)}\\n\\n"
4.4 核心配置层
4.4.1 生命周期管理(lifespan.py)
统一初始化 / 销毁外部存储客户端(MySQL、ES、Qdrant、Embedding),避免重复创建连接:
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):
"""
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()
4.4.2 请求上下文(context.py)
使用 ContextVar 存储请求 ID,异步场景下安全传递:
from contextvars import ContextVar
# 定义上下文变量,默认值为"1",实际由中间件设置为UUID
request_id_ctx_var = ContextVar("request_id", default="1")
4.4.3 结构化日志(log.py)
日志中携带 RequestID,实现分布式场景下的请求追踪:
import sys
from pathlib import Path
from loguru import logger
from app.conf.app_config import app_config
from app.core.context import request_id_ctx_var
# 日志格式:时间 + 级别 + RequestID + 位置 + 消息
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):
"""为日志记录注入RequestID"""
record["extra"]["request_id"] = request_id_ctx_var.get()
# 清空默认日志配置
logger.remove()
# 注入RequestID到日志
logger = logger.patch(inject_request_id)
# 控制台日志(开发环境)
if app_config.logging.console.enable:
logger.add(sink=sys.stdout, level=app_config.logging.console.level, format=log_format)
# 文件日志(生产环境)
if app_config.logging.file.enable:
log_path = Path(app_config.logging.file.path)
log_path.mkdir(parents=True, exist_ok=True)
logger.add(
sink=log_path / "app.log",
level=app_config.logging.file.level,
format=log_format,
rotation=app_config.logging.file.rotation, # 日志轮转(如按大小/时间)
retention=app_config.logging.file.retention, # 日志保留时长
encoding="utf-8"
)
五、接口测试
5.1 启动服务
# 进入main.py所在目录
cd data-agent
# 启动FastAPI开发服务器(自动重载)
fastapi dev main.py
5.2 接口调用
使用 Apifox/Postman 调用POST http://127.0.0.1:8000/api/query,请求体:
{
"query": "2026年Q1 GMV环比增长多少?"
}
5.3 响应示例(SSE 格式)
服务端会流式返回工作流执行进度和最终结果:
data: {"type":"progress","step":"召回字段信息","status":"running"}
data: {"type":"progress","step":"召回字段信息","status":"success"}
data: {"type":"progress","step":"过滤指标","status":"running"}
data: {"type":"progress","step":"过滤指标","status":"success"}
data: {"type":"result","content":"2026年Q1 GMV为1200万元,环比增长15.2%"}
六、前后端联调
6.1 启动前端项目
前端代码可从课程资料获取,另外启动项目需要node运行环境,node的安装包也可从课程资料获取。
在准备好node环境后,可在前端项目的根目录执行如下命令启动项目:
安装项目所需依赖
npm install
启动项目
npm run dev
6.2 访问前端页面
根据命令行输出反问指定页面即可,具体效果如下:

七、关键设计亮点
7.1 流式响应 + SSE:提升用户体验
通过StreamingResponse结合 SSE 协议,将 LangGraph 工作流的每个节点进度实时推送给前端,避免用户长时间等待,符合「智能问答」场景的交互需求。
7.2 依赖注入:解耦与可扩展
所有仓库、服务均通过依赖注入创建,便于替换实现(如测试时用 Mock 仓库),同时符合「开闭原则」。
7.3 全链路请求追踪
通过ContextVar+ 结构化日志,每个请求生成唯一 ID 并贯穿全流程,日志中可精准定位单个请求的所有操作,便于问题排查。
7.4 生命周期管理:资源复用
统一初始化 / 销毁外部客户端,避免每次请求创建新连接,提升性能并减少资源泄露风险。
八、总结
本专栏完整实现了一个「高性能、可扩展、易维护」的智能数据查询智能体项目,核心价值在于:
该架构可直接复用至智能 BI、数据问答、指标分析等场景,只需替换 LangGraph 工作流的具体节点逻辑即可适配不同业务需求。

![pip install安装markitdown时出现ERROR: ‘packages/markitdown[all]‘ is not a valid editable requirement解决方案-171主机测评](https://www.171host.com/wp-content/uploads/2026/08/20260823004656-6a8a430045592-220x150.png)
![pip install安装markitdown时出现ERROR: ‘packages/markitdown[all]‘ is not a valid editable requirement解决方案-171主机测评](https://www.171host.com/wp-content/uploads/2026/08/20260823004339-6a8a423bd3e97-220x150.png)

