欢迎光临
我们一直在努力

基于 FastAPI+LangGraph 实现流式响应的智能数据查询

一、前言

项目在智能数据问答场景中,用户往往需要实时感知查询请求的处理进度(如「召回字段」「过滤指标」「生成 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 生命周期管理:资源复用

统一初始化 / 销毁外部客户端,避免每次请求创建新连接,提升性能并减少资源泄露风险。

八、总结

本专栏完整实现了一个「高性能、可扩展、易维护」的智能数据查询智能体项目,核心价值在于:

  • 技术选型:FastAPI+LangGraph+SSE 的组合,兼顾高性能与业务复杂度;
  • 工程化:依赖注入、生命周期、结构化日志等实践,符合生产级应用标准;
  • 用户体验:流式响应实时反馈进度,解决传统「同步等待」的交互痛点。
  • 该架构可直接复用至智能 BI、数据问答、指标分析等场景,只需替换 LangGraph 工作流的具体节点逻辑即可适配不同业务需求。

    赞(0)
    未经允许不得转载:171主机测评 » 基于 FastAPI+LangGraph 实现流式响应的智能数据查询
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址