16 - 电商问数:查询接口实现与依赖组装
本章课程目标:
把第 15 章的假流式接口替换成真实问数工作流。
使用 FastAPI 依赖注入组装 Repository、Session、Client 和 Service。
使用 lifespan 管理 Qdrant、ES、MySQL、Embedding 等应用级资源。
学习建议: 这一章开始把 API 和问数工作流接起来。读的时候沿着一条请求走:HTTP 进入路由,路由拿到 QueryService,QueryService 调用 LangGraph,执行进度再通过 SSE 返回前端。依赖函数看起来多,但本质是在按层次组装底层对象;别把装配代码误读成业务逻辑。
对应代码分支: 16-api-query-service
第 15 章已经验证了 StreamingResponse + SSE 可以正常工作,但当时返回的是 fake_streamer() 模拟数据。本章要把它换成真实的问数智能体。
最终调用链路如下:
1 2 3 4 5 6 7 前端 / Apifox -> POST /api/query -> query_router.py 接收请求 -> Depends(get_query_service) 获取 QueryService -> QueryService.query(...) -> graph.astream(...) -> StreamingResponse 按 SSE 格式持续返回
完成这一章后,问数智能体就不再只是命令行里的内部脚本,而是一个可以通过 HTTP 调用的后端接口。
1、API 代码目录结构 这一阶段 API 代码主要涉及下面几个文件:
1 2 3 4 5 6 7 8 9 10 11 12 shopkeeper-agent/ ├─ main.py # FastAPI 入口脚本,负责创建应用并注册路由 └─ app/ ├─ api/ │ ├─ routers/ │ │ └─ query_router.py # 查询接口路由,接收请求并返回流式响应 │ ├─ schemas/ │ │ └─ query_schema.py # 查询接口请求体结构 │ ├─ dependencies.py # 查询接口依赖项,负责组装 QueryService │ └─ lifespan.py # FastAPI 生命周期事件,负责初始化和关闭外部客户端 └─ services/ └─ query_service.py # 查询接口核心业务逻辑,负责调用问数工作流
可以先把它们分成五层:
层次
文件
主要职责
入口层
main.py
创建 FastAPI 应用,注册生命周期和路由
HTTP 层
query_router.py
定义 /api/query,接收请求,返回流式响应
业务层
query_service.py
创建 State / Context,调用 LangGraph,并包装 SSE 消息
依赖层
dependencies.py
组装 Service、Repository、Session、Client
生命周期
lifespan.py
应用启动时初始化客户端,应用关闭时释放连接
这几个文件的关系可以记成一句话:
1 2 3 4 5 main 挂路由 router 接请求 service 调工作流 dependencies 组装对象 lifespan 管理应用级资源
2、入口和路由:让 HTTP 请求先进来 这一部分只解决一个问题:外部请求怎样进入我们的 Python 代码。
2.1 main.py:创建应用并挂载路由 项目对应文件路径:shopkeeper-agent/main.py
1 2 3 4 5 6 7 8 9 10 from fastapi import FastAPIfrom app.api.lifespan import lifespanfrom app.api.routers.query_router import query_routerapp = FastAPI(lifespan=lifespan) app.include_router(query_router)
这里最关键的是两行:
1 2 app = FastAPI(lifespan=lifespan) app.include_router(query_router)
lifespan 负责服务启动和关闭时的资源管理,include_router 负责把 query_router.py 中定义的接口挂到 FastAPI 应用上。
如果忘了 app.include_router(query_router),即使 query_router.py 里写了 /api/query,FastAPI 应用也不知道这个接口存在,/docs 页面里自然也看不到它。
2.2 query_router.py:路由层只做 HTTP 相关的事 项目对应文件路径:shopkeeper-agent/app/api/routers/query_router.py
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 from typing import Annotatedfrom fastapi import APIRouter, Dependsfrom starlette.responses import StreamingResponsefrom app.api.dependencies import get_query_servicefrom app.api.schemas.query_schema import QuerySchemafrom app.services.query_service import QueryServicequery_router = APIRouter() @query_router.post("/api/query" ) async def query_handler ( query: QuerySchema, query_service: Annotated[QueryService, Depends(get_query_service )], ): return StreamingResponse( query_service.query(query.query), media_type="text/event-stream" , )
路由层只做三件事:
1 2 3 1. 接收 QuerySchema 请求体 2. 通过 Depends 拿到 QueryService 3. 把 QueryService.query(...) 交给 StreamingResponse
它不直接创建 Qdrant、ES、MySQL Repository,也不直接执行图节点。路由层越薄,后续接口越容易维护。
本节最重要的是这一行:
1 query_service: Annotated[QueryService, Depends(get_query_service)]
它表达的是:当前接口需要一个 QueryService,这个对象由 get_query_service() 提供。路由只声明“我需要什么”,至于怎么创建,交给 dependencies.py。
3、QueryService:把一次请求变成一次图执行 QueryService 是这一章的核心。它把 HTTP 层传入的自然语言问题,转换成一次 LangGraph 工作流执行。
如果不抽出 QueryService,路由函数里就会堆满这些逻辑:
1 2 3 4 5 6 7 解析请求体 创建 DataAgentState 创建 DataAgentContext 准备 Repository 和 Client 调用 graph.astream(...) 把 chunk 包成 SSE 处理异常
这会让路由层既懂 HTTP,又懂工作流,又懂底层依赖,边界很快就乱了。抽出 QueryService 后,分工会清楚很多:
1 2 3 4 5 query_router.py -> 负责 HTTP 层:接请求、返响应 query_service.py -> 负责业务层:调问数智能体、组织流式输出
3.1 从测试脚本迁移到 API 前面测试 LangGraph 工作流时,核心代码大致是:
1 2 3 4 5 6 7 8 9 state = DataAgentState(query="统计华北地区的销售总额" ) context = DataAgentContext(...) async for chunk in graph.astream( input =state, context=context, stream_mode="custom" , ): print (chunk)
迁移到 API 后,变化主要有三处:
测试脚本
API 接口版本
用户问题写死在代码里
使用请求体传入的 query
上下文在测试脚本里手动创建
由 QueryService 接收依赖后创建
print(chunk) 输出到控制台
yield data: ...\n\n 流式写给前端
3.2 QueryService 核心代码 项目对应文件路径:shopkeeper-agent/app/services/query_service.py
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 import jsonfrom langchain_huggingface import HuggingFaceEndpointEmbeddingsfrom app.agent.context import DataAgentContextfrom app.agent.graph import graphfrom app.agent.state import DataAgentStatefrom app.repositories.es.value_es_repository import ValueESRepositoryfrom app.repositories.mysql.dw.dw_mysql_repository import DWMySQLRepositoryfrom app.repositories.mysql.meta.meta_mysql_repository import MetaMySQLRepositoryfrom app.repositories.qdrant.column_qdrant_repository import ColumnQdrantRepositoryfrom app.repositories.qdrant.metric_qdrant_repository import MetricQdrantRepositoryclass QueryService : def __init__ ( self, meta_mysql_repository: MetaMySQLRepository, embedding_client: HuggingFaceEndpointEmbeddings, dw_mysql_repository: DWMySQLRepository, column_qdrant_repository: ColumnQdrantRepository, metric_qdrant_repository: MetricQdrantRepository, value_es_repository: ValueESRepository, ): self .meta_mysql_repository = meta_mysql_repository self .dw_mysql_repository = dw_mysql_repository self .embedding_client = embedding_client self .column_qdrant_repository = column_qdrant_repository self .metric_qdrant_repository = metric_qdrant_repository self .value_es_repository = value_es_repository async def query (self, query: str ): state = DataAgentState(query=query) context = DataAgentContext( column_qdrant_repository=self .column_qdrant_repository, embedding_client=self .embedding_client, metric_qdrant_repository=self .metric_qdrant_repository, value_es_repository=self .value_es_repository, meta_mysql_repository=self .meta_mysql_repository, dw_mysql_repository=self .dw_mysql_repository, ) 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: error = {"type" : "error" , "message" : str (e)} yield f"data: {json.dumps(error, ensure_ascii=False , default=str )} \n\n"
这段代码里最重要的不是 json.dumps,而是两个对象:state 和 context。
3.3 state 和 context 怎么区分 state 保存的是本次任务会不断变化的业务数据,例如:
1 2 3 4 5 6 7 8 query keywords retrieved_column_infos retrieved_metric_infos table_infos metric_infos sql result
context 保存的是节点运行时需要使用的外部能力,例如:
1 2 3 4 5 6 Embedding Client ColumnQdrantRepository MetricQdrantRepository ValueESRepository MetaMySQLRepository DWMySQLRepository
可以这样记:
1 2 state:任务数据,图执行过程中会变 context:工具资源,节点执行时拿来用
不要把 Repository、Client 这类对象塞进 state。它们不是业务中间结果,而是节点执行时需要调用的外部能力,更适合放在 context 里。
3.4 为什么异常也要包装成 SSE 普通接口出错时,可以直接返回 500 状态码。但流式接口不太一样:一旦 StreamingResponse 开始往外写数据,HTTP 响应头通常已经发送出去了,后面就不能再随便改状态码。
所以当前版本先把异常包装成一条 SSE 消息:
1 2 error = {"type" : "error" , "message" : str (e)} yield f"data: {json.dumps(error, ensure_ascii=False , default=str )} \n\n"
前端拿到这条消息后,可以按 type=error 展示错误状态。下一章会继续完善节点级异常处理、前后端联调和 request_id 日志追踪。
4、dependencies.py:把对象创建交给依赖层 现在路由里已经声明了:
1 query_service: Annotated[QueryService, Depends(get_query_service)]
这意味着项目必须提供 get_query_service()。
如果直接在路由里创建 QueryService,路由就会知道太多底层细节:
1 2 3 4 5 Qdrant 客户端怎么取 ES 客户端怎么取 MySQL Session 怎么创建和释放 Repository 怎么实例化 QueryService 需要哪些参数
这些不是 HTTP 层该关心的事。所以本项目把依赖组装统一放到:
1 shopkeeper-agent/app/api/dependencies.py
4.1 先看最终依赖树 本章最终要组装的是 QueryService,依赖树大致如下:
1 2 3 4 5 6 7 8 9 get_query_service -> get_meta_mysql_repository -> get_meta_session -> get_embedding_client -> get_dw_mysql_repository -> get_dw_session -> get_column_qdrant_repository -> get_metric_qdrant_repository -> get_value_es_repository
再往底层,这些依赖会使用生命周期中初始化好的客户端管理器:
1 2 3 4 5 embedding_client_manager.client qdrant_client_manager.client es_client_manager.client meta_mysql_client_manager.session_factory dw_mysql_client_manager.session_factory
所以 dependencies.py 和 lifespan.py 是一组配合关系:
1 2 3 4 5 lifespan.py -> 应用启动时初始化客户端管理器 dependencies.py -> 每次请求中取出客户端或 Session,组装 Repository 和 QueryService
4.2 MySQL Session 是请求级资源 项目中的 MySQL Session 用带 yield 的依赖项管理:
1 2 3 4 5 6 7 8 9 10 11 12 13 async def get_meta_session (): """创建一次请求内使用的元数据库 Session""" async with meta_mysql_client_manager.session_factory() as meta_session: yield meta_session async def get_dw_session (): """创建一次请求内使用的数仓 Session""" async with dw_mysql_client_manager.session_factory() as dw_session: yield dw_session
执行顺序可以理解成:
1 2 3 4 5 请求需要 Session -> 创建 Session -> yield 给 Repository 使用 -> 请求结束 -> 退出 async with,释放 Session
这里不要把 Session 做成全局对象。数据库 Session 通常属于一次请求的工作单元,而客户端管理器、连接池这类才适合放到应用生命周期里。
4.3 Repository 和 Client 怎么组装 有了 Session,就可以创建 MySQL Repository:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 async def get_meta_mysql_repository ( session: Annotated[AsyncSession, Depends(get_meta_session )], ) -> MetaMySQLRepository: """基于请求级 Session 创建元数据仓储""" return MetaMySQLRepository(session) async def get_dw_mysql_repository ( session: Annotated[AsyncSession, Depends(get_dw_session )], ) -> DWMySQLRepository: """基于请求级 Session 创建数仓仓储""" return DWMySQLRepository(session)
Qdrant 和 ES Repository 使用的是应用启动阶段已经初始化好的客户端:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 async def get_column_qdrant_repository () -> ColumnQdrantRepository: """创建字段向量检索仓储""" return ColumnQdrantRepository(qdrant_client_manager.client) async def get_metric_qdrant_repository () -> MetricQdrantRepository: """创建指标向量检索仓储""" return MetricQdrantRepository(qdrant_client_manager.client) async def get_value_es_repository () -> ValueESRepository: """创建字段取值全文检索仓储""" return ValueESRepository(es_client_manager.client)
Embedding 客户端也是同样的思路:
1 2 3 4 async def get_embedding_client () -> HuggingFaceEndpointEmbeddings: """获取应用启动阶段初始化好的 Embedding 客户端""" return embedding_client_manager.client
4.4 最后组装 QueryService 所有依赖最终收拢到 get_query_service():
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 async def get_query_service ( meta_mysql_repository: Annotated[ MetaMySQLRepository, Depends(get_meta_mysql_repository ) ], embedding_client: Annotated[ HuggingFaceEndpointEmbeddings, Depends(get_embedding_client ) ], dw_mysql_repository: Annotated[DWMySQLRepository, Depends(get_dw_mysql_repository )], column_qdrant_repository: Annotated[ ColumnQdrantRepository, Depends(get_column_qdrant_repository ) ], metric_qdrant_repository: Annotated[ MetricQdrantRepository, Depends(get_metric_qdrant_repository ) ], value_es_repository: Annotated[ValueESRepository, Depends(get_value_es_repository )], ) -> QueryService: """组装一次查询所需的业务服务""" return QueryService( meta_mysql_repository=meta_mysql_repository, embedding_client=embedding_client, dw_mysql_repository=dw_mysql_repository, column_qdrant_repository=column_qdrant_repository, metric_qdrant_repository=metric_qdrant_repository, value_es_repository=value_es_repository, )
这就是 FastAPI 子依赖的价值:我们只声明依赖关系,FastAPI 会自动从叶子节点往上解析整棵依赖树。
5、lifespan.py:在应用启动时准备外部客户端 dependencies.py 能拿到 qdrant_client_manager.client、es_client_manager.client,前提是这些 manager 已经初始化。
所以需要在应用启动阶段执行:
1 2 3 4 5 qdrant_client_manager.init() embedding_client_manager.init() es_client_manager.init() meta_mysql_client_manager.init() dw_mysql_client_manager.init()
项目对应文件路径:shopkeeper-agent/app/api/lifespan.py
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 from contextlib import asynccontextmanagerfrom fastapi import FastAPIfrom app.clients.embedding_client_manager import embedding_client_managerfrom app.clients.es_client_manager import es_client_managerfrom app.clients.mysql_client_manager import ( dw_mysql_client_manager, meta_mysql_client_manager, ) from app.clients.qdrant_client_manager import qdrant_client_manager@asynccontextmanager async def lifespan (app: FastAPI ): """管理应用启动和关闭两个阶段的外部资源""" qdrant_client_manager.init() embedding_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()
按 yield 拆开看:
1 2 3 4 5 6 7 8 9 10 11 yield 前: FastAPI 应用启动时执行 初始化 Qdrant、Embedding、ES、MySQL 客户端管理器 yield 处: 应用进入运行状态 开始接收请求 yield 后: 应用关闭前执行 释放外部客户端连接
注意这里没有关闭 embedding_client_manager,因为当前项目代码里只对 Qdrant、ES、MySQL manager 提供了异步关闭逻辑。文档跟着实际代码走,不额外虚构关闭方法。
6、启动后端并测试真实查询接口 在 shopkeeper-agent 项目根目录启动后端:
1 uv run fastapi dev main.py
启动前要确认 MySQL、Qdrant、Elasticsearch 等外部服务可用,Embedding 相关配置也已经准备好。
用 Apifox 测试:
1 2 POST http://127.0.0.1:8000/api/query Content-Type: application/json
请求体示例:
1 2 3 { "query" : "统计华北地区销售额" }
如果一切正常,接口不会一次性返回完整结果,而是持续返回多段 SSE 消息。前端或 Apifox 能看到问数智能体执行过程中的进度输出。
本章小结:
本章把 /api/query 从协议验证推进到了真实业务执行。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 main.py -> 创建 FastAPI 应用 -> 注册 lifespan 和 query_router query_router.py -> 接收 QuerySchema -> 通过 Depends 获取 QueryService -> 用 StreamingResponse 返回 query_service.query(...) QueryService -> 创建 DataAgentState -> 创建 DataAgentContext -> 调用 graph.astream(..., stream_mode="custom") -> 把 chunk 转成 JSON 并包装成 SSE dependencies.py -> 组装 Session、Repository、Client、Service lifespan.py -> 应用启动时初始化客户端 -> 应用关闭时释放客户端
这一章完成后,后端查询接口已经能跑真实问数流程。下一章要继续处理交付级问题:节点失败时如何停止工作流、前端如何稳定消费消息,以及并发请求时如何通过 request_id 追踪日志。