20. WebAPI 层与 SSE 流式交互
本文档详细介绍知识库系统的 Web API 层设计与实现,包括文件导入和知识查询两套 API,以及 SSE 流式交互机制。
学习理念:Web 层是知识库系统的"门面"——它把后端的 LangGraph 工作流暴露为 REST API,支持文件上传导入和流式问答查询。核心设计:事件队列 + SSE 实现流式进度推送,FastAPI + Pydantic 实现类型安全的 API 接口。
海外对标:Web 层的"导入(POST /upload) + 查询(SSE GET /query)"双 API 设计对标 LangChain 的 LangServe(REST + Streaming)和 Vercel AI SDK 的
streamText。事件队列 + SSE 的异步模式对标 AWS SQS + WebSocket 的微服务架构。
本节 AI 替代率:~85% | 人工干预率:~15%
| 角色 | 能力范围 |
|---|---|
| 🤖 AI 擅长 | FastAPI 路由定义、Pydantic Schema、事件队列代码、SSE 端点、文件上传逻辑、异常处理 |
| 👤 人类需理解 | 事件队列的工作原理(队列 vs 直接 SSE)、流式 vs 非流式的 API 设计选择 |
阅读指引
| 颜色 | 章节 | AI 替代率 | 人工干预 | 说明 |
|---|---|---|---|---|
| 🟡 | §1 任务目标 | ~95% | ~5% | 学习目标明确 |
| 🟢 | §2 核心概念 | ~90% | ~10% | FastAPI / SSE / 事件队列 |
| 🟡 | §3 整体流程 | ~90% | ~10% | 两套 API 的运作流程 |
| 🟠 | §4 分步实现 | ~80% | ~20% | 事件队列 + SSE 是核心 |
| 🔴 | §4.5 主代码 | ~75% | ~25% | Web API 层完整实现 |
| 🟢 | §5 测试运行 | ~95% | ~5% | curl 命令验证 |
| 🟡 | §6 总结 | ~90% | ~10% | 设计要点回顾 |
技术栈健康度标签体系
| 技术 | 健康度 | 建议 |
|---|---|---|
| FastAPI | 🔥 巅峰 | Python Web 框架事实标准。自动 OpenAPI 文档 + Pydantic 校验。 |
| SSE | 🟢 稳定 | HTTP 单向推送,StreamingResponse + EventSourcePolyfill。 |
| Pydantic | 🔥 巅峰 | 数据校验标准库。BaseModel Schema + Field 验证。 |
| 事件队列 | 🟢 稳定 | queue.Queue + 后台线程,管理 SSE 连接和任务状态。 |
体系说明:🟢🟡🟠🔴 标识学习优先级 / AI 替代率;🔥🟢⏳⚠️💀 标识技术栈健康度。
中英文对照表
| English | 中文 | 本质 |
|---|---|---|
| SSE | 服务器推送事件 | HTTP 单向推送技术 |
| Event Queue | 事件队列 | 管理与 SSE 连接对应的异步事件通道 |
| StreamingResponse | 流式响应 | FastAPI 返回流式数据的类型 |
| REST API | REST 接口 | HTTP 标准化的 Web 服务接口 |
| Pydantic Schema | 数据模式 | Python 类型安全的请求/响应定义 |
💡 程序员比喻
- Web 层 就像 API Gateway——把后端的 LangGraph 工作流封装为 REST + SSE 端点。
- 事件队列 就像 Kafka 的 topic——每个 SSE 连接对应一个 topic,队列满了就推送给客户端。
- SSE 就像
kubectl get pods -w——建立连接后持续 watch,有新事件就 push。
1. 任务目标
本节课将实现知识库系统的 Web API 层,包括文件导入和知识查询两套 API。
该层是整个系统与前端交互的入口,负责:
- 接收 HTTP 请求并路由到对应处理逻辑
- 通过 后台任务 异步执行长耗时操作
- 使用 SSE(Server-Sent Events) 实现实时进度推送和流式答案输出
- 管理会话历史和任务状态
学完本节你将掌握:
- FastAPI 框架的核心概念(路由、依赖注入、后台任务)
- SSE 服务端推送技术的原理与实现
- 事件队列(Event Queue)在流式交互中的应用
- 导入流程与查询流程的 API 设计模式
- 前后端实时通信的完整架构
2. 核心概念扫盲
2.1 系统架构总览

2.2 FastAPI 核心概念
2.2.1 路由定义
from fastapi import FastAPI
app = FastAPI()
@app.post("/upload")
async def upload_files(files: List[UploadFile]):
"""POST 请求处理"""
pass
@app.get("/status/{task_id}")
async def get_task(task_id: str):
"""GET 请求处理,路径参数"""
pass2.2.2 后台任务
from fastapi import BackgroundTasks
@app.post("/upload")
async def upload_files(
background_tasks: BackgroundTasks,
files: List[UploadFile]
):
# 立即返回响应
task_id = str(uuid.uuid4())
# 将耗时操作放入后台执行
background_tasks.add_task(run_import_task, task_id, files)
return {"task_id": task_id, "message": "Processing started"}后台任务的优势:
- 快速响应客户端请求
- 异步执行耗时操作
- 适合文件处理、数据导入等场景
2.2.3 Pydantic 数据模型
from pydantic import BaseModel, Field
class QueryRequest(BaseModel):
"""查询请求模型"""
query: str = Field(..., description="查询内容")
session_id: Optional[str] = Field(None, description="会话ID")
is_stream: bool = Field(False, description="是否流式返回")2.3 SSE
2.3.1 什么是 SSE?
SSE 是一种服务器向客户端单向推送事件的技术:
┌──────────┐ HTTP 长连接 ┌──────────┐
│ │ ──────────────────────────>│ │
│ Client │ │ Server │
│ │ <── event: ready ──────────│ │
│ │ <── event: delta ──────────│ │
│ │ <── event: delta ──────────│ │
│ │ <── event: final ──────────│ │
└──────────┘ └──────────┘2.3.2 SSE 消息格式
event: delta
data: {"delta": "根据"}
event: delta
data: {"delta": "参考内容"}
event: final
data: {"answer": "完整答案...", "status": "completed"}2.2.3 SSE vs WebSocket
| 特性 | SSE | WebSocket |
|---|---|---|
| 通信方向 | 单向(服务器→客户端) | 双向 |
| 协议 | HTTP | 独立协议 |
| 自动重连 | 内置支持 | 需手动实现 |
| 浏览器支持 | 原生支持 | 原生支持 |
| 适用场景 | 实时推送、流式输出 | 双向通信、聊天 |
2.4 事件队列
事件队列是 SSE 实现的核心数据结构:
import queue
from typing import Dict
# 全局会话队列存储
# Key: session_id, Value: queue.Queue
_session_stream: Dict[str, queue.Queue] = {}
def create_sse_queue(session_id: str) -> queue.Queue:
"""创建并注册一个新的 SSE 队列"""
q = queue.Queue()
_session_stream[session_id] = q
return q
def push_to_session(session_id: str, event: str, data: dict):
"""向指定会话推送事件"""
stream_queue = _session_stream.get(session_id)
if stream_queue:
stream_queue.put({"event": event, "data": data})工作原理:

2.5 任务状态管理
系统使用内存字典管理任务状态:
# 任务状态常量
TASK_STATUS_PENDING = "pending" # 等待处理
TASK_STATUS_PROCESSING = "processing" # 处理中
TASK_STATUS_COMPLETED = "completed" # 已完成
TASK_STATUS_FAILED = "failed" # 失败
# 内存态任务追踪
_tasks_running_list: Dict[str, List[str]] = {} # 正在运行的节点
_tasks_done_list: Dict[str, List[str]] = {} # 已完成的节点
_tasks_status: Dict[str, str] = {} # 任务状态
_tasks_result: Dict[str, Dict[str, str]] = {} # 任务结果3. WebAPI 业务处理流程(总)
3.1 导入流程交互图

3.2 查询流程交互图(流式模式)

3.3 查询流程交互图(非流式模式)

4. WebAPI 业务处理流程(分)
4.1 目标
实现两套 FastAPI 应用,分别处理文件导入和知识查询请求,支持同步/异步两种模式,并通过 SSE 实现实时进度推送。
4.2 需求分析
| 需求项 | 说明 |
|---|---|
| 文件上传 | 支持多文件上传,保存到本地和 MinIO |
| 后台处理 | 耗时任务放入后台异步执行 |
| 进度追踪 | 支持查询任务状态和已完成节点 |
| 流式输出 | 支持 SSE 实时推送生成内容 |
| 会话管理 | 支持查询和清除历史记录 |
| 跨域支持 | 允许前端跨域访问 |
4.3 实现流程
4.3.1 实现流程图

4.3.2 具体实现步骤
Step 1: 创建 FastAPI 应用
创建应用工厂函数,配置中间件和路由:
from fastapi import FastAPI
from fastapi.middleware.cors import CORSMiddleware
from fastapi.staticfiles import StaticFiles
def create_app() -> FastAPI:
"""创建 FastAPI 应用"""
app = FastAPI(
title="Query Service",
description="知识库查询服务"
)
# 允许跨域(开发环境)
app.add_middleware(
CORSMiddleware,
allow_origins=["*"], # 允许所有来源
allow_credentials=True, # 允许携带凭证
allow_methods=["*"], # 允许所有方法
allow_headers=["*"], # 允许所有头部
)
# 静态文件挂载(前端页面)
page_dir = get_front_page_dir()
if os.path.exists(page_dir):
app.mount("/front", StaticFiles(directory=page_dir), name="front")
# 注册路由
_register_routes(app)
return appCORS 配置说明:
allow_origins: 允许访问的前端域名allow_credentials: 是否允许携带 Cookieallow_methods: 允许的 HTTP 方法allow_headers: 允许的请求头
Step 2: 定义请求/响应模型
使用 Pydantic 定义数据模型:
# schemas/query_schema.py
from typing import Optional
from pydantic import BaseModel, Field
class QueryRequest(BaseModel):
"""查询请求模型"""
query: str = Field(..., description="查询内容")
session_id: Optional[str] = Field(None, description="会话ID")
is_stream: bool = Field(False, description="是否流式返回")# schemas/task_schema.py
from typing import List
from pydantic import BaseModel, Field
class TaskStatusResponse(BaseModel):
"""任务状态响应"""
status: str = Field(..., description="任务状态")
done_list: List[str] = Field(..., description="已完成节点列表")
running_list: List[str] = Field(..., description="正在运行节点列表")Pydantic 的优势:
- 自动数据验证
- 自动生成 OpenAPI 文档
- 类型提示支持
Step 3: 实现导入 API 路由
处理文件上传和状态查询:
# api/import_router.py
@app.post("/upload", response_model=UploadResponse)
async def upload_files(
background_tasks: BackgroundTasks,
files: List[UploadFile] = File(...)
):
"""上传文件"""
# 1. 获取文件导入 service
service = get_file_import_service()
# 2. 处理导入文件(保存到本地和 MinIO)
task_ids, date_dir = service.process_files(files)
# 3. 添加后台任务处理导入流程
for task_id, file in zip(task_ids, files):
file_dir = os.path.join(date_dir, task_id)
import_file_path = os.path.join(file_dir, file.filename)
background_tasks.add_task(
_run_import_task,
task_id,
file_dir,
import_file_path
)
# 4. 立即返回任务 ID(不等待处理完成)
return UploadResponse(
message="Files uploaded successfully",
task_ids=task_ids
)
@app.get("/status/{task_id}", response_model=TaskStatusResponse)
async def get_task(task_id: str):
"""获取任务状态"""
task_info = TaskService.get_task_info(task_id)
return TaskStatusResponse(**task_info)后台任务处理函数:
def _run_import_task(task_id: str, file_dir: str, import_file_path: str):
"""后台任务:执行导入流程"""
setup_logging()
service = get_file_import_service()
service.run_import_task(task_id, file_dir, import_file_path)Step 4: 实现查询 API 路由
处理查询请求,支持流式和非流式两种模式:
# api/query_router.py
@app.post("/query")
async def query(request: QueryRequest, background_tasks: BackgroundTasks):
"""处理查询请求"""
# 1. 解析参数
user_query = request.query
session_id = request.session_id or str(uuid.uuid4())
is_stream = request.is_stream
# 2. 更新任务状态
update_task_status(session_id, TASK_STATUS_PROCESSING, is_stream)
# 3. 根据模式选择处理方式
if is_stream:
# 流式模式:后台任务 + SSE
background_tasks.add_task(
_run_query_graph, session_id, user_query, is_stream
)
await asyncio.sleep(0.1) # 等待队列创建
return {"message": "Query submitted", "session_id": session_id}
else:
# 非流式模式:同步执行
_run_query_graph(session_id, user_query, is_stream)
answer = get_task_result(session_id, "answer", "")
return {
"message": "处理完成",
"session_id": session_id,
"answer": answer
}SSE 流式输出端点:
@app.get("/stream/{session_id}")
async def stream(session_id: str, request: Request):
"""SSE 实时返回结果"""
return StreamingResponse(
sse_generator(session_id, request),
media_type="text/event-stream"
)历史记录管理:
@app.get("/history/{session_id}")
async def history(session_id: str, limit: int = 50):
"""查询当前会话历史记录"""
records = get_recent_messages(session_id, limit=limit)
items = []
for r in records:
items.append({
"_id": str(r.get("_id")) if r.get("_id") else "",
"session_id": r.get("session_id", ""),
"role": r.get("role", ""),
"text": r.get("text", ""),
"rewritten_query": r.get("rewritten_query", ""),
"item_names": r.get("item_names", []),
"ts": r.get("ts")
})
return {"session_id": session_id, "items": items}
@app.delete("/history/{session_id}")
async def clear_chat_history(session_id: str):
"""清除会话历史记录"""
count = clear_history(session_id)
return {"message": "History cleared", "deleted_count": count}Step 5: 实现 SSE 事件队列
SSE 工具模块的核心实现:
# tools/sse_utils.py
import json
import queue
import asyncio
from typing import Dict, Any, AsyncGenerator
from fastapi import Request
class SSEEvent:
"""SSE 事件类型常量"""
READY = "ready" # 连接建立
PROGRESS = "progress" # 任务节点进度
DELTA = "delta" # LLM 流式输出增量
FINAL = "final" # 最终完整答案
ERROR = "error" # 错误信息
CLOSE = "__close__" # 关闭连接信号
# 全局 SSE 会话队列存储
_session_stream: Dict[str, queue.Queue] = {}
def create_sse_queue(session_id: str) -> queue.Queue:
"""创建并注册一个新的 SSE 队列"""
q = queue.Queue()
_session_stream[session_id] = q
return q
def get_sse_queue(session_id: str):
"""获取指定 session 的队列"""
return _session_stream.get(session_id)
def remove_sse_queue(session_id: str):
"""移除指定 session 的队列"""
_session_stream.pop(session_id, None)
def _sse_pack(event: str, data: Dict[str, Any]) -> str:
"""打包 SSE 消息格式"""
payload = json.dumps(data, ensure_ascii=False)
return f"event: {event}\ndata: {payload}\n\n"
def push_to_session(session_id: str, event: str, data: Dict[str, Any]):
"""向指定会话推送事件"""
stream_queue = get_sse_queue(session_id)
if stream_queue:
stream_queue.put({"event": event, "data": data})SSE 生成器(异步):
async def sse_generator(
session_id: str, request: Request
) -> AsyncGenerator[str, None]:
"""SSE 生成器,用于 FastAPI 的 StreamingResponse"""
stream_queue = get_sse_queue(session_id)
if stream_queue is None:
return
loop = asyncio.get_running_loop()
try:
# 发送连接建立信号
yield _sse_pack("ready", {})
while True:
# 检测客户端断开
if await request.is_disconnected():
break
try:
# 使用 run_in_executor 避免阻塞 async 事件循环
msg = await loop.run_in_executor(
None, stream_queue.get, True, 1.0
)
except queue.Empty:
continue
event = msg.get("event")
data = msg.get("data")
# 特殊关闭事件
if event == "__close__":
break
yield _sse_pack(event, data)
except (asyncio.CancelledError, ConnectionResetError, BrokenPipeError):
# 生成器被取消/对端断开:静默退出
return
finally:
# 清理资源
remove_sse_queue(session_id)Step 6: 实现任务状态管理
任务追踪工具模块:
# tools/task_utils.py
from typing import Dict, List
from knowledge.tools.sse_utils import push_to_session
# 内存态任务追踪
_tasks_running_list: Dict[str, List[str]] = {}
_tasks_done_list: Dict[str, List[str]] = {}
_tasks_status: Dict[str, str] = {}
_tasks_result: Dict[str, Dict[str, str]] = {}
# 节点名 -> 中文名映射
_NODE_NAME_TO_CN: Dict[str, str] = {
"upload_file": "上传文件",
"entry": "检查文件",
"pdf_to_md": "PDF转Markdown",
"document_split": "文档切分",
"item_name_confirm": "确认问题产品",
"answer_output": "生成答案",
# ... 更多节点映射
}
def add_running_task(task_id: str, node_name: str, push_queue: bool = False):
"""添加正在运行的节点任务"""
_ensure_task(task_id)
running = _tasks_running_list[task_id]
if node_name not in running:
running.append(node_name)
if push_queue:
task_push_queue(task_id)
def add_done_task(task_id: str, node_name: str, push_queue: bool = False):
"""添加已完成的节点任务"""
_ensure_task(task_id)
# 从 running 中移除
running = _tasks_running_list[task_id]
_tasks_running_list[task_id] = [n for n in running if n != node_name]
# 追加到 done
done = _tasks_done_list[task_id]
if node_name not in done:
done.append(node_name)
if push_queue:
task_push_queue(task_id)
def update_task_status(task_id: str, status_name: str, push_queue: bool = False):
"""更新任务状态"""
_tasks_status[task_id] = status_name
if push_queue:
task_push_queue(task_id)
def task_push_queue(task_id: str):
"""推送任务进度到 SSE"""
push_to_session(task_id, "progress", {
"status": get_task_status(task_id),
"done_list": get_done_task_list(task_id),
"running_list": get_running_task_list(task_id),
})Step 7: 实现后台任务处理
导入流程后台任务:
def _run_import_task(task_id: str, file_dir: str, import_file_path: str):
"""后台任务:执行导入流程"""
service = get_file_import_service()
service.run_import_task(task_id, file_dir, import_file_path)查询流程后台任务:
def _run_query_graph(session_id: str, user_query: str, is_stream: bool):
"""后台任务:执行查询流程图"""
# 流式模式:创建 SSE 队列
if is_stream:
create_sse_queue(session_id)
# 构建初始状态
default_state = {
"original_query": user_query,
"session_id": session_id,
"is_stream": is_stream
}
# 执行查询流程图
query_app.invoke(default_state)
# 更新任务状态
update_task_status(session_id, TASK_STATUS_COMPLETED, is_stream)Step 8: 实现服务层
文件导入服务:
# services/file_import_service.py
class FileImportService:
"""文件导入服务类"""
def __init__(self, file_dir: str):
self.file_dir = file_dir
os.makedirs(file_dir, exist_ok=True)
def process_files(self, files: List[UploadFile]) -> Tuple[List[str], str]:
"""处理文件上传"""
date_dir = self.get_date_dir()
task_ids = []
for file in files:
# 1. 生成任务 ID
task_id = str(uuid.uuid4())
task_ids.append(task_id)
# 2. 创建任务状态
TaskService.create_task_node(task_id, "upload_file")
# 3. 保存文件
file_dir = os.path.join(date_dir, task_id)
import_file_path = self.save_upload_file(file, file_dir)
# 4. 上传到 MinIO
self.upload_to_minio(file, import_file_path)
# 5. 标记上传完成
TaskService.complete_task_node(task_id, "upload_file")
return task_ids, date_dir
def run_import_task(self, task_id: str, file_dir: str, import_file_path: str):
"""运行导入任务"""
try:
# 1. 更新任务状态
TaskService.update_task_status(task_id, "processing")
# 2. 创建初始状态
initial_state = create_default_state(
task_id=task_id,
file_dir=file_dir,
import_file_path=import_file_path
)
# 3. 运行图并追踪节点完成
for event in kb_import_app.stream(initial_state):
for key, value in event.items():
TaskService.complete_task_node(task_id, key)
# 4. 更新任务状态
TaskService.update_task_status(task_id, "completed")
except Exception as e:
TaskService.update_task_status(task_id, "failed")任务状态服务:
# services/task_service.py
class TaskService:
"""任务服务类"""
@staticmethod
def create_task_node(task_id: str, initial_node: str = "upload_file"):
"""创建新任务"""
add_running_task(task_id, initial_node)
@staticmethod
def complete_task_node(task_id: str, node_name: str):
"""标记节点完成"""
add_done_task(task_id, node_name)
@staticmethod
def get_task_info(task_id: str) -> Dict[str, any]:
"""获取任务完整信息"""
return {
"status": get_task_status(task_id),
"done_list": get_done_task_list(task_id),
"running_list": get_running_task_list(task_id),
}4.4 代码实现
查询 API 路由完整代码
"""知识库查询 API 路由"""
import os
import uuid
import asyncio
from typing import Optional
from fastapi import FastAPI, BackgroundTasks, HTTPException, Request
from fastapi.responses import FileResponse, StreamingResponse
from fastapi.middleware.cors import CORSMiddleware
from fastapi.staticfiles import StaticFiles
from knowledge.tools.task_utils import (
update_task_status, get_task_result, clear_task,
TASK_STATUS_PROCESSING, TASK_STATUS_COMPLETED
)
from knowledge.tools.sse_utils import (
create_sse_queue, push_to_session, sse_generator, SSEEvent
)
from knowledge.tools.mongo_history_utils import get_recent_messages, clear_history
from knowledge.processor.query_process.main_graph import query_app
from knowledge.schemas.query_schema import QueryRequest
def create_app() -> FastAPI:
"""创建 FastAPI 应用"""
app = FastAPI(title="Query Service", description="知识库查询服务")
app.add_middleware(
CORSMiddleware,
allow_origins=["*"],
allow_credentials=True,
allow_methods=["*"],
allow_headers=["*"],
)
page_dir = get_front_page_dir()
if os.path.exists(page_dir):
app.mount("/front", StaticFiles(directory=page_dir), name="front")
_register_routes(app)
return app
def _register_routes(app: FastAPI):
"""注册路由"""
@app.get("/chat.html")
async def chat_page():
return FileResponse(os.path.join(get_front_page_dir(), "chat.html"))
@app.post("/query")
async def query(request: QueryRequest, background_tasks: BackgroundTasks):
user_query = request.query
session_id = request.session_id or str(uuid.uuid4())
is_stream = request.is_stream
update_task_status(session_id, TASK_STATUS_PROCESSING, is_stream)
if is_stream:
background_tasks.add_task(_run_query_graph, session_id, user_query, is_stream)
await asyncio.sleep(0.1)
return {"message": "Query submitted", "session_id": session_id}
else:
_run_query_graph(session_id, user_query, is_stream)
answer = get_task_result(session_id, "answer", "")
return {"session_id": session_id, "answer": answer}
@app.get("/stream/{session_id}")
async def stream(session_id: str, request: Request):
return StreamingResponse(
sse_generator(session_id, request),
media_type="text/event-stream"
)
@app.get("/history/{session_id}")
async def history(session_id: str, limit: int = 50):
records = get_recent_messages(session_id, limit=limit)
items = [
{
"role": r.get("role", ""),
"text": r.get("text", ""),
"ts": r.get("ts")
}
for r in records
]
return {"session_id": session_id, "items": items}
@app.delete("/history/{session_id}")
async def clear_chat_history(session_id: str):
count = clear_history(session_id)
return {"deleted_count": count}
def _run_query_graph(session_id: str, user_query: str, is_stream: bool):
"""后台任务:执行查询流程图"""
if is_stream:
create_sse_queue(session_id)
default_state = {
"original_query": user_query,
"session_id": session_id,
"is_stream": is_stream
}
query_app.invoke(default_state)
update_task_status(session_id, TASK_STATUS_COMPLETED, is_stream)
app = create_app()
if __name__ == "__main__":
import uvicorn
uvicorn.run(app, host="0.0.0.0", port=8001)5. 测试运行
5.1 运行 API 服务测试
# 启动导入服务(端口 8000)
python -m knowledge.api.import_router
# 启动查询服务(端口 8001)
python -m knowledge.api.query_router5.2 测试代码
测试导入 API
import requests
# 上传文件
files = [
('files', ('test.pdf', open('test.pdf', 'rb'), 'application/pdf'))
]
response = requests.post('http://localhost:8000/upload', files=files)
print(response.json())
# {"message": "Files uploaded successfully", "task_ids": ["xxx-xxx-xxx"]}
# 查询状态
task_id = response.json()['task_ids'][0]
status = requests.get(f'http://localhost:8000/status/{task_id}')
print(status.json())
# {"status": "processing", "done_list": ["上传文件"], "running_list": ["检查文件"]}测试查询 API(非流式)
import requests
response = requests.post('http://localhost:8001/query', json={
"query": "万用表怎么测电压?",
"session_id": "test_001",
"is_stream": False
})
print(response.json())
# {"session_id": "test_001", "answer": "根据参考内容..."}测试查询 API(流式)
import requests
# 1. 提交查询
response = requests.post('http://localhost:8001/query', json={
"query": "万用表怎么测电压?",
"session_id": "test_002",
"is_stream": True
})
session_id = response.json()['session_id']
# 2. 连接 SSE 流
import sseclient
url = f'http://localhost:8001/stream/{session_id}'
response = requests.get(url, stream=True)
client = sseclient.SSEClient(response)
for event in client.events():
print(f"Event: {event.event}, Data: {event.data}")5.3 预期输出
导入流程状态轮询输出
第 1 次轮询:
{"status": "processing", "done_list": ["上传文件"], "running_list": ["检查文件"]}
第 2 次轮询:
{"status": "processing", "done_list": ["上传文件", "检查文件", "PDF转Markdown"], "running_list": ["文档切分"]}
第 3 次轮询:
{"status": "completed", "done_list": ["上传文件", "检查文件", "PDF转Markdown", "文档切分", "向量生成", "导入向量数据库", "处理完成"], "running_list": []}查询流程 SSE 输出
Event: ready, Data: {}
Event: delta, Data: {"delta": "根据"}
Event: delta, Data: {"delta": "参考内容"}
Event: delta, Data: {"delta": ","}
Event: delta, Data: {"delta": "万用表"}
Event: delta, Data: {"delta": "测量"}
Event: delta, Data: {"delta": "电压"}
Event: delta, Data: {"delta": "的步骤"}
Event: delta, Data: {"delta": "如下:"}
...
Event: final, Data: {"answer": "根据参考内容,万用表测量电压的步骤如下...", "status": "completed"}5.4 处理前后对比
导入流程数据流
# 请求
POST /upload
Content-Type: multipart/form-data
files: [test.pdf]
# 响应(立即返回)
{
"message": "Files uploaded successfully",
"task_ids": ["a1b2c3d4-e5f6-7890-abcd-ef1234567890"]
}
# 后台处理完成后查询状态
GET /status/a1b2c3d4-e5f6-7890-abcd-ef1234567890
{
"status": "completed",
"done_list": [
"上传文件",
"检查文件",
"PDF转Markdown",
"Markdown图片处理",
"文档切分",
"主体名称识别",
"向量生成",
"导入向量数据库",
"导入知识图谱",
"处理完成"
],
"running_list": []
}查询流程数据流
# 请求(流式)
POST /query
{
"query": "万用表怎么测电压?",
"session_id": "session_001",
"is_stream": true
}
# 响应(立即返回)
{
"message": "Query submitted",
"session_id": "session_001"
}
# SSE 流
GET /stream/session_001
event: ready
data: {}
event: delta
data: {"delta": "根据参考内容,"}
event: delta
data: {"delta": "万用表测量电压"}
...
event: final
data: {"answer": "完整答案...", "status": "completed"}6. 总结
6.1 功能概览

6.2 设计要点
要点 1:后台任务异步处理
@app.post("/upload")
async def upload_files(background_tasks: BackgroundTasks, files: List[UploadFile]):
# 快速返回响应
task_ids = process_files(files)
# 耗时操作放入后台
for task_id in task_ids:
background_tasks.add_task(_run_import_task, task_id, ...)
return {"task_ids": task_ids} # 立即返回优势: 用户无需等待处理完成,提升体验。
要点 2:SSE 事件队列模式
# 生产者(后台任务)
def _run_query_graph(session_id, ...):
create_sse_queue(session_id)
# ... 处理过程中
push_to_session(session_id, "delta", {"delta": "内容"})
# 消费者(SSE 生成器)
async def sse_generator(session_id, request):
queue = get_sse_queue(session_id)
while True:
msg = await loop.run_in_executor(None, queue.get, True, 1.0)
yield _sse_pack(msg["event"], msg["data"])优势: 解耦生产者和消费者,支持实时推送。
要点 3:流式/非流式双模式
if is_stream:
# 流式:后台任务 + SSE
background_tasks.add_task(_run_query_graph, ...)
return {"session_id": session_id} # 立即返回
else:
# 非流式:同步执行
_run_query_graph(...)
return {"answer": get_task_result(...)} # 等待完成优势: 满足不同场景需求(实时交互 vs API 集成)。
要点 4:任务状态追踪
# 节点完成时更新状态
for event in kb_import_app.stream(initial_state):
for key, value in event.items():
TaskService.complete_task_node(task_id, key)
# 前端轮询获取进度
@app.get("/status/{task_id}")
async def get_task(task_id: str):
return TaskService.get_task_info(task_id)优势: 前端可实时展示处理进度。
要点 5:资源自动清理
async def sse_generator(...):
try:
while True:
# 检测客户端断开
if await request.is_disconnected():
break
# ... 推送事件
finally:
# 无论如何都清理资源
remove_sse_queue(session_id)优势: 防止内存泄漏,确保资源正确释放。
企业痛点映射
| 痛点 | 传统方案 | Web API 方案 | 效率提升 |
|---|---|---|---|
| 前端无法看到后端处理进度 | 轮询或直接等待 | SSE 事件队列实时推送(ready/progress/delta/final) | 用户体验提升 ~80% |
| 文件导入与查询 API 耦合 | 两套逻辑混在同一个端点 | 导入(GET/POST upload) + 查询(SSE POST query) 分离 | 职责清晰,可独立扩展 |
| 流式 vs 非流式难以兼容 | 两套不同的返回值处理 | SSE 统一 + is_stream 标志 | 前端统一用 EventSource |
Remote & Agent 应用场景价值
Remote 场景价值:FastAPI 自动生成 OpenAPI 文档,远程团队可直接在 Swagger UI 中测试 API。SSE 流式输出支持远程网络环境下的实时数据传输。
Agent 落地场景:Web 层的 POST /query 端点可直接作为 Agent 的 Tool API——Agent 发送查询请求,接收 SSE 流式返回的结果。
EventQueue机制确保多个 Agent 并发查询时互不干扰。
Git Commit 对应
本节 Web API 层对应提交记录(参考值):
<待补充>cd shopkeeper_brain
git log --oneline --all -- knowledge/web/