Skip to content

知识库查询 —— 答案生成节点

本文档详细介绍知识库查询流程中的答案生成节点(answer_output)的设计与实现。该节点是整个查询流程的终点,负责将检索到的上下文、历史对话、知识图谱关系组装成提示词,调用 LLM 生成最终答案,并支持流式输出和历史记录持久化。

学习理念:答案生成是查询流程的"最后一公里"——整合所有检索结果,通过字符预算控制(12000 chars)智能分配上下文配额。核心设计:流式/非流式统一接口 + 已有答案短路跳过 LLM + MongoDB 历史持久化。

海外对标:Context Budget 控制对标 LangChain 的 ConversationBufferWindowMemory,流式 SSE 推送对标 Vercel AI SDK 的 streamText。字符预算分配策略(文档>历史>图谱)与 Google Vertex AI Context 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 EventSSE 事件服务端推送事件类型
Prompt Assembly提示词组装多区块拼接为 LLM Prompt
History Persistence历史持久化写入 MongoDB 供后续使用

💡 程序员比喻

  • AnswerOutputNode 就像 git commit -m——前面的检索(git add)把变化放到暂存区,答案生成(git commit)打包成一条完整消息。
  • 字符预算控制 就像 max_line_length=80 的 lint——文档(核心代码)> 历史(注释)> 图谱(文档字符串)。
  • 流式 vs 非流式 就像 kubectl logs -f vs kubectl 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 有输入长度限制,需要控制提示词总长度:

python
# 总字符预算
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 事件类型

python
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: 检查已有答案

某些情况下(如商品确认提示),答案已在前置节点生成:

python
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)

    # ... 后续步骤

已有答案的推送逻辑:

python
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: 获取查询文本

获取用于提示词的查询文本:

python
question = state.get("rewritten_query") or state.get("original_query", "")
item_names = state.get("item_names") or []

优先级说明:

  • rewritten_query:经过商品名确认节点重写的完整查询
  • original_query:用户原始输入(降级选项)
Step 3: 格式化检索文档

将重排序后的文档格式化为提示词文本:

python
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: 格式化历史对话

将历史对话格式化为提示词文本:

python
@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: 格式化图谱三元组

将知识图谱三元组格式化为文本描述:

python
@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 提示词:

python
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 生成答案

根据模式选择流式或非流式生成:

python
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)

流式生成实现:

python
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 的传递。

python
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:

python
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 记录结构:

json
{
    "session_id": "sess_001",
    "role": "assistant",
    "text": "根据参考内容,RS-12数字万用表测量电压的步骤如下...",
    "item_names": ["RS-12数字万用表"],
    "timestamp": "2024-01-15T10:30:00Z"
}
Step 9: 发送结束事件

流式模式下发送最终完成事件:

python
# 流式模式发送结束事件
if is_stream:
    push_to_session(
        session_id,
        SSEEvent.FINAL,
        {"answer": state.get("answer", ""), "status": "completed"}
    )

return state

SSE 事件格式:

event: final
data: {"answer": "根据参考内容,RS-12数字万用表测量电压...", "status": "completed"}

4.4 代码实现

完整的节点实现代码:

python
"""答案输出节点

组装提示词、调用 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 运行答案生成节点测试

bash
# 确保配置了必要的环境变量
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_output

5.2 测试代码

python
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 事件

状态变化示例:

python
# 处理前
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:字符预算控制

python
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:流式与非流式统一处理

python
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:已有答案的处理

python
# 商品确认节点可能已生成提示答案
if state.get("answer"):
    self._push_existing_answer(state)  # 直接推送,跳过 LLM
else:
    prompt = self._build_prompt(state)
    self._generate(state, prompt)      # 正常生成流程

场景示例:

  • 用户问题中有多个可能的商品名
  • 商品确认节点生成澄清提示:"请问您是指 A 还是 B?"
  • 答案生成节点直接推送此提示,无需调用 LLM

要点 4:降级与容错

python
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 对应

本节答案生成节点对应提交记录(参考值):

<待补充>
bash
cd shopkeeper_brain
git log --oneline --all -- knowledge/processor/query_process/nodes/answer_output.py

OPC 超级个体实战指南