知识库查询 —— 答案生成节点
本文档详细介绍知识库查询流程中的答案生成节点(answer_output)的设计与实现。该节点是整个查询流程的终点,负责将检索到的上下文、历史对话、知识图谱关系组装成提示词,调用 LLM 生成最终答案,并支持流式输出和历史记录持久化。
学习理念:答案生成是查询流程的"最后一公里"——整合所有检索结果,通过字符预算控制(12000 chars)智能分配上下文配额。核心设计:流式/非流式统一接口 + 已有答案短路跳过 LLM + MongoDB 历史持久化。
海外对标:Context Budget 控制对标 LangChain 的
ConversationBufferWindowMemory,流式 SSE 推送对标 Vercel AI SDK 的streamText。字符预算分配策略(文档>历史>图谱)与 Google Vertex AIContext Cache优先级一致。
本节 AI 替代率:~80% | 人工干预率:~20%
| 角色 | 能力范围 |
|---|---|
| 🤖 AI 擅长 | 格式化文档/历史/图谱代码、字符预算控制、流式/非流式 LLM 调用、SSE 推送、MongoDB 写入、测试代码 |
| 👤 人类需理解 | 字符预算分配策略(文档>历史>图谱的优先级)、流式 vs 非流式选择 |
阅读指引
| 颜色 | 章节 | AI 替代率 | 人工干预 | 说明 |
|---|---|---|---|---|
| 🟡 | §1 任务目标 | ~95% | ~5% | 学习目标明确 |
| 🟢 | §2 核心概念 | ~90% | ~10% | Prompt工程 / 预算控制 / SSE |
| 🟡 | §3 整体流程 | ~90% | ~10% | 节点位置 + 输入输出 |
| 🟠 | §4 分步实现 | ~75% | ~25% | 9 步流程 |
| 🔴 | §4.4 主代码 | ~70% | ~30% | AnswerOutputNode + _format_docs |
| 🟢 | §5 测试运行 | ~95% | ~5% | 看预期输出 |
| 🟡 | §6 总结 | ~90% | ~10% | 要点回顾 |
技术栈健康度标签体系
| 技术 | 健康度 | 建议 |
|---|---|---|
| LangChain ChatOpenAI | 🔥 巅峰 | LLM 调用标准封装。llm.invoke / llm.stream 统一接口。 |
| SSE | 🟢 稳定 | HTTP 单向推送,push_to_session + SSEEvent 枚举。 |
| MongoDB | 🟢 稳定 | 历史持久化。save_chat_message 模式固定。 |
| 字符预算 | 🟢 稳定 | 12000 chars 配额,文档>历史>图谱分配策略。 |
体系说明:🟢🟡🟠🔴 标识学习优先级 / AI 替代率;🔥🟢⏳⚠️💀 标识技术栈健康度。
中英文对照表
| English | 中文 | 本质 |
|---|---|---|
| Context Budget | 上下文预算 | LLM 输入长度的配额管理 |
| Stream vs Batch | 流式 vs 非流式 | llm.stream() vs llm.invoke() |
| SSE Event | SSE 事件 | 服务端推送事件类型 |
| Prompt Assembly | 提示词组装 | 多区块拼接为 LLM Prompt |
| History Persistence | 历史持久化 | 写入 MongoDB 供后续使用 |
💡 程序员比喻
- AnswerOutputNode 就像
git commit -m——前面的检索(git add)把变化放到暂存区,答案生成(git commit)打包成一条完整消息。- 字符预算控制 就像
max_line_length=80的 lint——文档(核心代码)> 历史(注释)> 图谱(文档字符串)。- 流式 vs 非流式 就像
kubectl logs -fvskubectl logs。- 已有答案短路 就像
if cache.get(key): return——商品确认已有答案则跳过 LLM。
1. 任务目标
本节课将实现知识库查询流程中的 答案生成节点(answer_output)。
该节点是整个查询流程的终点,负责将检索到的上下文、历史对话、知识图谱关系组装成提示词,调用 LLM 生成最终答案,并支持 流式输出 和 历史记录持久化。
学完本节你将掌握:
- LLM 提示词工程的最佳实践
- 字符预算控制(Context Budget)技术
- SSE(Server-Sent Events)流式输出实现
- LangChain ChatOpenAI 的流式/非流式调用
- MongoDB 历史记录持久化设计
2. 核心概念扫盲
2.1 为什么答案生成是最复杂的节点?
答案生成节点需要整合整个查询流程的所有结果:

2.2 提示词工程要点
提示词的质量直接决定答案的质量。本系统的提示词包含以下几个区块:

2.3 字符预算控制
LLM 有输入长度限制,需要控制提示词总长度:
# 总字符预算
budget = config.max_context_chars # 默认 12000
# 依次分配给各区块
context_str, budget = self._format_docs(docs, budget) # 文档优先
history_str, budget = self._format_history(history, budget) # 历史次之
graph_str, budget = self._format_triples(triples, budget) # 图谱最后预算分配策略:
- 检索文档优先级最高(核心参考内容)
- 历史对话次之(理解上下文)
- 图谱关系最后(补充信息)
2.4 流式输出 vs 非流式输出

2.5 SSE 事件类型
class SSEEvent:
READY = "ready" # 连接建立
PROGRESS = "progress" # 任务节点进度
DELTA = "delta" # LLM 流式输出增量
FINAL = "final" # 最终完整答案
ERROR = "error" # 错误信息
CLOSE = "__close__" # 关闭连接信号事件流示例:
event: ready
data: {}
event: delta
data: {"delta": "根据"}
event: delta
data: {"delta": "参考内容"}
event: delta
data: {"delta": ",万用表"}
... (更多增量)
event: final
data: {"answer": "根据参考内容,万用表测量电压的步骤...", "status": "completed"}3. 答案生成业务处理流程(总)
3.1 节点在流程中的位置

3.2 节点输入输出

4. 答案生成业务处理流程(分)
4.1 目标
实现一个答案生成节点,将检索结果组装成提示词,调用 LLM 生成答案,支持流式输出并持久化历史记录。
4.2 需求分析
| 需求项 | 说明 |
|---|---|
| 多源整合 | 整合文档、历史、图谱、商品名等多种上下文 |
| 预算控制 | 控制提示词总长度,避免超出 LLM 限制 |
| 流式支持 | 支持边生成边推送,提升用户体验 |
| 降级处理 | LLM 调用失败时返回友好错误信息 |
| 历史持久化 | 将答案写入 MongoDB 供后续对话使用 |
| 已有答案 | 支持跳过生成直接输出已有答案 |
4.3 实现流程
4.3.1 实现流程图

4.3.2 具体实现步骤
Step 1: 检查已有答案
某些情况下(如商品确认提示),答案已在前置节点生成:
def process(self, state: QueryGraphState) -> QueryGraphState:
session_id = state.get("session_id")
is_stream = state.get("is_stream")
# 已有答案(如商品确认提示)→ 直接推送
if state.get("answer"):
self._push_existing_answer(state)
else:
# 构建提示词 → 生成答案
prompt = self._build_prompt(state)
state["prompt"] = prompt
self._generate(state, prompt)
# ... 后续步骤已有答案的推送逻辑:
def _push_existing_answer(self, state: QueryGraphState):
"""将已有答案推送到流或任务结果。"""
answer = state["answer"]
if state.get("is_stream"):
# 流式模式:通过 SSE 推送
push_to_session(state["session_id"], SSEEvent.DELTA, {"delta": answer})
else:
# 非流式模式:写入任务结果
set_task_result(state["session_id"], "answer", answer)Step 2: 获取查询文本
获取用于提示词的查询文本:
question = state.get("rewritten_query") or state.get("original_query", "")
item_names = state.get("item_names") or []优先级说明:
rewritten_query:经过商品名确认节点重写的完整查询original_query:用户原始输入(降级选项)
Step 3: 格式化检索文档
将重排序后的文档格式化为提示词文本:
def _format_docs(self, docs: List[Dict], budget: int) -> tuple:
"""格式化重排序文档,带字符预算控制。"""
lines = []
used = 0
for i, doc in enumerate(docs, 1):
text = (doc.get("text") or "").strip()
if not text:
continue
# 构建元信息行
meta = [f"[{i}]"]
# 添加来源、ID、URL、标题
for key, fmt in [
("source", "[{}]"),
("chunk_id", "[chunk_id={}]"),
("url", "[url={}]"),
("title", "[title={}]"),
]:
val = str(doc.get(key) or "").strip()
if val:
meta.append(fmt.format(val))
# 添加相关性得分
score = doc.get("score")
if score is not None:
meta.append(f"[score={float(score):.4f}]")
# 组合元信息和内容
doc_str = " ".join(meta) + "\n" + text
# 检查预算
if used + len(doc_str) > budget:
break
lines.append(doc_str)
used += len(doc_str) + 2 # +2 for \n\n
return "\n\n".join(lines), budget - used格式化结果示例:
[1] [local] [chunk_id=chunk_001] [title=万用表使用手册] [score=0.9234]
数字万用表测量电压的步骤:首先将旋钮转到V档位,然后将黑色表笔插入COM孔...
[2] [web] [url=https://example.com/guide] [title=电压测量指南] [score=0.8756]
使用万用表测量直流电压时,需要注意选择合适的量程...Step 4: 格式化历史对话
将历史对话格式化为提示词文本:
@staticmethod
def _format_history(history: List[Dict], budget: int) -> tuple:
"""格式化历史对话。"""
lines = []
used = 0
for msg in history:
for role, key in [("用户", "user"), ("助手", "assistant")]:
text = msg.get(key)
if not text:
continue
line = f"{role}: {text}"
used += len(line) + 1
# 检查预算
if used > budget:
return "\n".join(lines), budget - used
lines.append(line)
return "\n".join(lines), budget - used格式化结果示例:
用户: 万用表怎么测电压?
助手: 测量电压时,首先需要选择合适的档位...
用户: 如果量程选错了会怎样?
助手: 如果量程选择过小,可能会损坏万用表...Step 5: 格式化图谱三元组
将知识图谱三元组格式化为文本描述:
@staticmethod
def _format_triples(triples: List, budget: int) -> tuple:
"""格式化图谱三元组。"""
lines = []
used = 0
for tr in triples:
line = (str(tr) if tr is not None else "").strip()
if not line:
continue
# 检查预算
if used + len(line) > budget:
break
lines.append(line)
used += len(line) + 1
return "\n".join(lines), budget - used格式化结果示例:
万用表 -[包含]-> 表笔
万用表 -[包含]-> 旋钮
表笔 -[用于]-> 测量电压
表笔 -[用于]-> 测量电流Step 6: 组装完整提示词
使用模板组装完整的 LLM 提示词:
def _build_prompt(self, state: QueryGraphState) -> str:
"""根据检索结果、历史对话、图谱关系组装 LLM 提示词。"""
config = get_config()
budget = config.max_context_chars # 总字符预算
question = state.get("rewritten_query") or state.get("original_query", "")
item_names = state.get("item_names") or []
# 依次分配预算
context_str, budget = self._format_docs(
state.get("reranked_docs") or [], budget
)
history_str, budget = self._format_history(
state.get("history") or [], budget
)
graph_str, budget = self._format_triples(
state.get("kg_triples") or [], budget
)
# 填充模板
return ANSWER_PROMPT.format(
context=context_str or "无参考内容",
history=history_str or "无历史对话",
item_names=", ".join(item_names) if item_names else "无指定商品",
graph_relation_description=graph_str or "无图谱关系",
question=question,
)完整提示词示例:
你是一个智能助手,请根据参考内容回答用户的问题。
要求:
1 尽量基于【参考内容】和【图谱关系描述】 作答,不要编造不存在的事实。
2 如果用户的问题需要通过图片来辅助说明...
【参考内容】
[1] [local] [chunk_id=chunk_001] [score=0.9234]
数字万用表测量电压的步骤...
[2] [web] [url=https://example.com] [score=0.8756]
使用万用表测量直流电压时...
【历史对话】
用户: 万用表是什么?
助手: 万用表是一种多功能电子测量仪器...
【相关商品/实体】
RS-12数字万用表
【图谱关系描述】
万用表 -[包含]-> 表笔
表笔 -[用于]-> 测量电压
【用户问题】
RS-12数字万用表怎么测电压?
请回答:Step 7: 调用 LLM 生成答案
根据模式选择流式或非流式生成:
def _generate(self, state: QueryGraphState, prompt: str):
"""调用 LLM 生成答案(流式/非流式)。"""
self.log_step("generate", "生成答案")
llm = get_llm_client()
session_id = state.get("session_id")
if state.get("is_stream"):
state["answer"] = self._stream_generate(llm, prompt, session_id)
else:
state["answer"] = self._invoke_generate(llm, prompt, session_id)流式生成实现:
def _stream_generate(self, llm, prompt: str, session_id: str) -> str:
"""流式生成,逐 chunk 推送。"""
result = ""
try:
for chunk in llm.stream(prompt):
# 提取增量内容
delta = getattr(chunk, "content", "") or ""
if delta:
result += delta
# 实时推送给前端
push_to_session(session_id, "delta", {"delta": delta})
except Exception as e:
self.logger.error(f"流式生成出错: {e}")
return result非流式生成实现:
🟡 【P1 看注释就行】 字符预算控制——
budget = config.max_context_chars。策略:文档优先 > 历史次之 > 图谱最后。注意剩余 budget 的传递。
def _invoke_generate(self, llm, prompt: str, session_id: str) -> str:
"""非流式生成。"""
try:
response = llm.invoke(prompt)
answer = response.content
# 写入任务结果供轮询获取
set_task_result(session_id, "answer", answer)
return answer
except Exception as e:
self.logger.error(f"生成回答出错: {e}")
return "抱歉,生成回答时出现错误。"Step 8: 写入历史记录
将答案持久化到 MongoDB:
def _write_history(self, state: QueryGraphState):
"""将助手回答写入 Mongo 历史记录。"""
answer = (state.get("answer") or "").strip()
if not answer:
return
try:
save_chat_message(
session_id=state.get("session_id", "default"),
role="assistant",
text=answer,
rewritten_query="",
item_names=state.get("item_names") or [],
)
except Exception as e:
self.logger.warning(f"写入历史记录失败: {e}")MongoDB 记录结构:
{
"session_id": "sess_001",
"role": "assistant",
"text": "根据参考内容,RS-12数字万用表测量电压的步骤如下...",
"item_names": ["RS-12数字万用表"],
"timestamp": "2024-01-15T10:30:00Z"
}Step 9: 发送结束事件
流式模式下发送最终完成事件:
# 流式模式发送结束事件
if is_stream:
push_to_session(
session_id,
SSEEvent.FINAL,
{"answer": state.get("answer", ""), "status": "completed"}
)
return stateSSE 事件格式:
event: final
data: {"answer": "根据参考内容,RS-12数字万用表测量电压...", "status": "completed"}4.4 代码实现
完整的节点实现代码:
"""答案输出节点
组装提示词、调用 LLM 生成答案,支持流式/非流式输出,写入历史记录。
"""
from typing import List, Dict, Any
from knowledge.processor.query_process.base import BaseNode
from knowledge.processor.query_process.state import QueryGraphState
from knowledge.processor.query_process.prompt import ANSWER_PROMPT
from knowledge.processor.query_process.config import get_config
from knowledge.tools.sse_utils import push_to_session, SSEEvent
from knowledge.tools.task_utils import set_task_result
from knowledge.tools.llm_utils import get_llm_client
from knowledge.tools.mongo_history_utils import save_chat_message
class AnswerOutputNode(BaseNode):
"""答案输出节点。
流程: 检查已有答案 → 构建提示词 → LLM 生成 → 写入历史 → 发送结束事件
"""
name = "answer_output"
# ================================================================== #
# 主流程 #
# ================================================================== #
def process(self, state: QueryGraphState) -> QueryGraphState:
session_id = state.get("session_id")
is_stream = state.get("is_stream")
# Step 1: 已有答案(如商品确认提示)→ 直接推送
if state.get("answer"):
self._push_existing_answer(state)
else:
# Step 2-7: 构建提示词 → 生成答案
prompt = self._build_prompt(state)
state["prompt"] = prompt
self._generate(state, prompt)
# Step 8: 写入历史
if state.get("answer"):
self._write_history(state)
# Step 9: 流式模式发送结束事件
if is_stream:
push_to_session(
session_id,
SSEEvent.FINAL,
{"answer": state.get("answer", ""), "status": "completed"}
)
return state
# ================================================================== #
# 已有答案推送 #
# ================================================================== #
def _push_existing_answer(self, state: QueryGraphState):
"""将已有答案推送到流或任务结果。"""
answer = state["answer"]
if state.get("is_stream"):
push_to_session(state["session_id"], SSEEvent.DELTA, {"delta": answer})
else:
set_task_result(state["session_id"], "answer", answer)
# ================================================================== #
# 提示词构建 #
# ================================================================== #
def _build_prompt(self, state: QueryGraphState) -> str:
"""根据检索结果、历史对话、图谱关系组装 LLM 提示词。"""
config = get_config()
budget = config.max_context_chars
question = state.get("rewritten_query") or state.get("original_query", "")
item_names = state.get("item_names") or []
# Step 3: 格式化检索文档
context_str, budget = self._format_docs(
state.get("reranked_docs") or [], budget
)
# Step 4: 格式化历史对话
history_str, budget = self._format_history(
state.get("history") or [], budget
)
# Step 5: 格式化图谱关系
graph_str, budget = self._format_triples(
state.get("kg_triples") or [], budget
)
# Step 6: 组装完整提示词
return ANSWER_PROMPT.format(
context=context_str or "无参考内容",
history=history_str or "无历史对话",
item_names=", ".join(item_names) if item_names else "无指定商品",
graph_relation_description=graph_str or "无图谱关系",
question=question,
)
def _format_docs(self, docs: List[Dict], budget: int) -> tuple:
"""格式化重排序文档,带字符预算控制。"""
lines = []
used = 0
for i, doc in enumerate(docs, 1):
text = (doc.get("text") or "").strip()
if not text:
continue
meta = [f"[{i}]"]
for key, fmt in [
("source", "[{}]"), ("chunk_id", "[chunk_id={}]"),
("url", "[url={}]"), ("title", "[title={}]"),
]:
val = str(doc.get(key) or "").strip()
if val:
meta.append(fmt.format(val))
score = doc.get("score")
if score is not None:
meta.append(f"[score={float(score):.4f}]")
doc_str = " ".join(meta) + "\n" + text
if used + len(doc_str) > budget:
break
lines.append(doc_str)
used += len(doc_str) + 2
return "\n\n".join(lines), budget - used
@staticmethod
def _format_history(history: List[Dict], budget: int) -> tuple:
"""格式化历史对话。"""
lines = []
used = 0
for msg in history:
for role, key in [("用户", "user"), ("助手", "assistant")]:
text = msg.get(key)
if not text:
continue
line = f"{role}: {text}"
used += len(line) + 1
if used > budget:
return "\n".join(lines), budget - used
lines.append(line)
return "\n".join(lines), budget - used
@staticmethod
def _format_triples(triples: List, budget: int) -> tuple:
"""格式化图谱三元组。"""
lines = []
used = 0
for tr in triples:
line = (str(tr) if tr is not None else "").strip()
if not line:
continue
if used + len(line) > budget:
break
lines.append(line)
used += len(line) + 1
return "\n".join(lines), budget - used
# ================================================================== #
# LLM 生成 #
# ================================================================== #
def _generate(self, state: QueryGraphState, prompt: str):
"""调用 LLM 生成答案(流式/非流式)。"""
self.log_step("generate", "生成答案")
llm = get_llm_client()
session_id = state.get("session_id")
if state.get("is_stream"):
state["answer"] = self._stream_generate(llm, prompt, session_id)
else:
state["answer"] = self._invoke_generate(llm, prompt, session_id)
def _stream_generate(self, llm, prompt: str, session_id: str) -> str:
"""流式生成,逐 chunk 推送。"""
result = ""
try:
for chunk in llm.stream(prompt):
delta = getattr(chunk, "content", "") or ""
if delta:
result += delta
push_to_session(session_id, "delta", {"delta": delta})
except Exception as e:
self.logger.error(f"流式生成出错: {e}")
return result
def _invoke_generate(self, llm, prompt: str, session_id: str) -> str:
"""非流式生成。"""
try:
response = llm.invoke(prompt)
answer = response.content
set_task_result(session_id, "answer", answer)
return answer
except Exception as e:
self.logger.error(f"生成回答出错: {e}")
return "抱歉,生成回答时出现错误。"
# ================================================================== #
# 历史记录 #
# ================================================================== #
def _write_history(self, state: QueryGraphState):
"""将助手回答写入 Mongo 历史记录。"""
answer = (state.get("answer") or "").strip()
if not answer:
return
try:
save_chat_message(
session_id=state.get("session_id", "default"),
role="assistant",
text=answer,
rewritten_query="",
item_names=state.get("item_names") or [],
)
except Exception as e:
self.logger.warning(f"写入历史记录失败: {e}")
# 兼容入口
_node_instance = AnswerOutputNode()
def node_answer_output(state: QueryGraphState) -> QueryGraphState:
"""兼容原有调用方式的入口函数。"""
return _node_instance(state)5. 测试运行
5.1 运行答案生成节点测试
# 确保配置了必要的环境变量
export OPENAI_API_KEY="your-api-key"
export OPENAI_API_BASE="https://your-llm-endpoint"
export LLM_MODEL="qwen3-32b"
# 运行测试(需要编写测试脚本)
python -m knowledge.processor.query_process.nodes.answer_output5.2 测试代码
if __name__ == "__main__":
from dotenv import load_dotenv
import json
# 加载环境变量
load_dotenv()
# 初始化日志
from knowledge.processor.query_process.base import setup_logging
setup_logging()
print("=" * 60)
print("开始测试: 答案生成节点 (AnswerOutputNode)")
print("=" * 60)
# 构造模拟状态
mock_state = {
"session_id": "test_session_001",
"is_stream": False, # 非流式测试
"original_query": "万用表怎么测电压?",
"rewritten_query": "RS-12数字万用表如何测量电压?",
"item_names": ["RS-12数字万用表"],
"reranked_docs": [
{
"text": "数字万用表测量电压步骤:1. 将旋钮转到V档位;2. 黑表笔插COM孔,红表笔插V孔;3. 将表笔并联到被测点两端。",
"source": "local",
"chunk_id": "chunk_001",
"title": "万用表使用手册",
"score": 0.9234
},
{
"text": "测量直流电压时需注意正负极性,红表笔接正极,黑表笔接负极。",
"source": "web",
"url": "https://example.com/guide",
"title": "电压测量指南",
"score": 0.8756
}
],
"history": [
{"user": "万用表是什么?", "assistant": "万用表是一种多功能电子测量仪器..."}
],
"kg_triples": [
"万用表 -[包含]-> 表笔",
"万用表 -[包含]-> 旋钮",
"表笔 -[用于]-> 测量电压"
]
}
print("【输入状态】:")
print(f" query: {mock_state['rewritten_query']}")
print(f" item_names: {mock_state['item_names']}")
print(f" reranked_docs: {len(mock_state['reranked_docs'])} 篇")
print(f" kg_triples: {len(mock_state['kg_triples'])} 条")
print("-" * 60)
# 执行答案生成
result = node_answer_output(mock_state)
# 打印结果
print("\n【生成结果】:")
print("-" * 60)
print(result.get("answer", "无答案"))
print("-" * 60)
# 打印提示词(调试用)
if result.get("prompt"):
print("\n【构建的提示词】:")
print(result["prompt"][:500] + "...")
print("\n测试完成")5.3 预期输出
============================================================
开始测试: 答案生成节点 (AnswerOutputNode)
============================================================
【输入状态】:
query: RS-12数字万用表如何测量电压?
item_names: ['RS-12数字万用表']
reranked_docs: 2 篇
kg_triples: 3 条
------------------------------------------------------------
[answer_output] [generate] 生成答案
【生成结果】:
------------------------------------------------------------
根据参考内容,RS-12数字万用表测量电压的步骤如下:
1. **选择档位**:将旋钮转到V档位(直流电压选DCV,交流电压选ACV)
2. **连接表笔**:
- 黑色表笔插入COM孔
- 红色表笔插入V孔
3. **进行测量**:将两只表笔并联到被测点两端
4. **注意事项**:
- 测量直流电压时注意正负极性
- 红表笔接正极,黑表笔接负极
- 选择合适的量程,避免超量程损坏仪表
------------------------------------------------------------
【构建的提示词】:
你是一个智能助手,请根据参考内容回答用户的问题。
要求:
1 尽量基于【参考内容】和【图谱关系描述】 作答...
【参考内容】
[1] [local] [chunk_id=chunk_001] [title=万用表使用手册] [score=0.9234]
数字万用表测量电压步骤...
测试完成5.4 处理前后对比
| 对比项 | 处理前 | 处理后 |
|---|---|---|
answer | 不存在或为空 | LLM 生成的完整答案 |
prompt | 不存在 | 组装好的完整提示词 |
| MongoDB | 无记录 | 新增助手回答记录 |
| SSE(流式) | 无事件 | DELTA + FINAL 事件 |
状态变化示例:
# 处理前
state = {
"session_id": "sess_001",
"rewritten_query": "RS-12数字万用表如何测量电压?",
"reranked_docs": [...],
"kg_triples": [...],
# answer 不存在
# prompt 不存在
}
# 处��后
state = {
"session_id": "sess_001",
"rewritten_query": "RS-12数字万用表如何测量电压?",
"reranked_docs": [...],
"kg_triples": [...],
"answer": "根据参考内容,RS-12数字万用表测量电压的步骤如下...",
"prompt": "你是一个智能助手,请根据参考内容回答用户的问题..."
}6. 总结
6.1 节点功能概览

6.2 节点设计要点
要点 1:字符预算控制
budget = config.max_context_chars # 12000
# 依次分配,优先级递减
context_str, budget = self._format_docs(docs, budget)
history_str, budget = self._format_history(history, budget)
graph_str, budget = self._format_triples(triples, budget)设计思路:
- 检索文档是核心参考,优先级最高
- 历史对话提供上下文理解,优先级次之
- 图谱关系是补充信息,优先级最低
- 预算用尽后自动截断,保证不超限
要点 2:流式与非流式统一处理
if state.get("is_stream"):
state["answer"] = self._stream_generate(llm, prompt, session_id)
else:
state["answer"] = self._invoke_generate(llm, prompt, session_id)两种模式的差异:
| 模式 | API | 输出方式 | 适用场景 |
|---|---|---|---|
| 流式 | llm.stream() | SSE 逐 chunk 推送 | Web 实时交互 |
| 非流式 | llm.invoke() | 任务结果一次返回 | 后台任务/API |
要点 3:已有答案的处理
# 商品确认节点可能已生成提示答案
if state.get("answer"):
self._push_existing_answer(state) # 直接推送,跳过 LLM
else:
prompt = self._build_prompt(state)
self._generate(state, prompt) # 正常生成流程场景示例:
- 用户问题中有多个可能的商品名
- 商品确认节点生成澄清提示:"请问您是指 A 还是 B?"
- 答案生成节点直接推送此提示,无需调用 LLM
要点 4:降级与容错
def _invoke_generate(self, llm, prompt, session_id):
try:
response = llm.invoke(prompt)
return response.content
except Exception as e:
self.logger.error(f"生成回答出错: {e}")
return "抱歉,生成回答时出现错误。" # 友好错误信息企业痛点映射
| 痛点 | 传统方案 | AI Agent 方案 | 效率提升 |
|---|---|---|---|
| LLM 输入长度超限 | 信息丢失 | 字符预算控制 12000 chars 智能分配 | 不超限 ~100% |
| 流式/非流式要两套代码 | 重复实现 | if is_stream: stream else: invoke 统一 | 开发成本降 ~50% |
| 商品确认后仍需 LLM | 浪费 token | 已有答案短路直接推送 | Token 浪费 ~0% |
Remote & Agent 应用场景价值
Remote 场景价值:LLM API 远程调用,SSE 通过 Redis Pub/Sub 跨服务推送。MongoDB 远程持久化,团队共享会话状态。
Agent 落地场景:AnswerOutputNode 可封装为"答案生成 Agent"——接收多源文档 → 预算分配 → LLM 生成 → 写入历史。
_stream_generate适合实时 Agent,_invoke_generate适合后台 Agent。
Git Commit 对应
本节答案生成节点对应提交记录(参考值):
<待补充>cd shopkeeper_brain
git log --oneline --all -- knowledge/processor/query_process/nodes/answer_output.py