第二阶段:核心 RAG 链路第 10 章

RAG Pipeline 主流程与 Prompt 生成

贯通 RAG Pipeline 的八个执行阶段,明确 FAQ 快速路径、上下文筛选与截断、Prompt Profile 选择、变量注入及答案引用增强。

20 分钟阅读

第一部分:技术背景 — Pipeline 设计模式

1.1 Pipeline vs Chain

Chain(链):固定的步骤序列,A → B → C → D,没有分支。

Pipeline(管道):有分支、有快慢路径的流程。每一步可以提前结束(如 FAQ 命中时跳过文档检索),也可以根据上一步的结果调整下一步的参数。

本项目的 RAG 流程是 Pipeline 而非 Chain:

用户问题
  → 查询路由(direct_answer / faq_exact / retrieval)
  → 检索准备(历史 / 意图 / source / 按需改写 / 计划 / 变体)
  → 按 RetrievalPlan 执行 FAQ 检索(先查安全缓存,高置信标准直出可结束)
  → 按 RetrievalPlan 执行文档检索(先查安全缓存)
  → 上下文构建 → LLM 生成

1.2 Pipeline 的模块化拆分

项目将 Pipeline 拆分为多个职责单一的文件:

qa_core/pipeline/
├── rag.py          # 主流程编排(stream_query, debug_retrieval)
├── runtime.py      # 请求上下文(RAGQueryContext)和事件工具函数
├── steps.py        # 查询路由、检索准备、Prompt 准备
├── retrieval_steps.py  # FAQ / 文档检索执行
├── context.py      # 上下文构建(筛选、去重、格式化)
├── rewrite.py      # 查询改写
├── query_variants.py  # 查询变体生成
├── events.py       # 事件构造(start/status/token/end/error)
├── confidence.py   # 最终答案置信度(answer_confidence)
└── citations.py    # 答案引用增强

第二部分:8 个 Stage 主流程

2.0 先按代码主线阅读,再看横切专题

当前模块的主线以qa_core/pipeline/rag.py::stream_query()为准。阅读时先沿着下面的调用链从上到下理解;缓存、性能和事件协议属于跨阶段能力,放在主线之后回看,不把它们误认为新的 Stage。

stream_query()
  -> create_query_context()                         # Stage 0
  -> decide_route()                                 # Stage 1
  -> prepare_retrieval()                            # Stage 2
  -> _search_and_generate()
      -> search_faq()                               # Stage 3
      -> search_doc()                               # Stage 4
      -> prepare_answer()                           # Stage 5
      -> stream_llm_answer()                        # Stage 6
  -> enforce_answer_citations()                     # Stage 6 后处理
  -> finalize_generated_answer_confidence()         # Stage 6 后处理
  -> history.add_turn() + finish_success()          # Stage 7
执行阶段 主入口函数 主要实现文件 项目文档深挖位置
Stage 0 create_query_context() qa_core/pipeline/runtime.py 2.4 主线代码骨架
Stage 1 decide_route() qa_core/pipeline/steps.py 2.5~2.6.1
Stage 2 prepare_retrieval() qa_core/pipeline/steps.py 第三部分 3.1~3.5
Stage 3 search_faq() qa_core/pipeline/retrieval_steps.py 第四部分 4.1
Stage 4 search_doc() qa_core/pipeline/retrieval_steps.py 第四部分 4.2
Stage 5 prepare_answer() qa_core/pipeline/steps.pycontext.py 第五部分 5.1~5.4、第六部分 6.1~6.3
Stage 6 stream_llm_answer()、引用增强、生成后核验 steps.pycitations.pyconfidence.py 第七部分 7.1~7.5
Stage 7 history.add_turn()finish_success()finish_error() rag.pyruntime.py 第八部分 8.1~8.3

这张表解决的是“先读什么”的问题。后面的专题小节可以按需跳转,但不能改变上面这条实际调用顺序。

Stage 0-7 可视化总览

flowchart TD
    Start(["🚀 stream_query() 开始"]) --> Stage0

    Stage0["🏗️ Stage 0:创建上下文<br/>场景/数据域/会话/trace/KB版本"] --> Stage1

    Stage1["🧭 Stage 1:查询路由<br/>decide_route()"] --> RouteCheck{"RouteDecision.route?"}
    RouteCheck -->|"direct_answer"| End1["📝 返回直接答案<br/>问候/越界/转人工/边界"]
    RouteCheck -->|"faq_exact"| End2Fast["🎯 返回 FAQ 标准答案<br/>intent=FAQ_QUERY"]
    RouteCheck -->|"retrieval"| Stage2

    Stage2["🎯 Stage 2:检索准备<br/>历史/意图/source/改写/计划/变体/Prompt Profile"] --> Stage3

    Stage3["🔍 Stage 3:FAQ 检索<br/>先查缓存 / 未命中查 Milvus"] --> Stage3Check{"FAQ 高置信直出?<br/>精确匹配 或 分数>阈值"}
    Stage3Check -->|"✅"| End3["📋 返回标准答案<br/>hit_type: faq_direct"]
    Stage3Check -->|"❌"| Stage4

    Stage4["📚 Stage 4:文档检索<br/>先查缓存 / 未命中查 Milvus<br/>Dense + Sparse Hybrid / Rerank"] --> Stage5
    Stage5["📊 Stage 5:上下文与 Prompt 组装<br/>筛选去重截断 / 选择 Profile / 组织引用来源"] --> Stage5Check{"召回结果是否不足?"}
    Stage5Check -->|"信息不足"| End5["⚠️ 信息不足提示<br/>引导联系人工"]
    Stage5Check -->|"✅ 有资料"| Stage6

    Stage6["🤖 Stage 6:LLM 流式生成 + 引用增强<br/>逐 token 推送/补充来源标注"] --> Stage7

    Stage7["💾 Stage 7:保存历史<br/>写入 Trace"] --> Final(["✅ end 事件<br/>返回 sources/intent/retrieval/answer_confidence"])

    style Stage1 fill:#FEF3C7,stroke:#D97706,stroke-width:2px
    style Stage2 fill:#EFF6FF,stroke:#2563EB,stroke-width:2px
    style Stage3 fill:#ECFDF5,stroke:#059669,stroke-width:2px
    style Stage4 fill:#FFFBEB,stroke:#D97706,stroke-width:2px
    style Stage5 fill:#FFFBEB,stroke:#D97706,stroke-width:2px
    style Stage6 fill:#FEF2F2,stroke:#DC2626,stroke-width:2px
    style Final fill:#ECFDF5,stroke:#059669,stroke-width:3px

代码执行时序图

这张图对应的是qa_core/pipeline/rag.py::stream_query()的真实执行顺序。它比前面的阶段总览更适合代码调试,因为可以直接按函数名下断点。

sequenceDiagram
    autonumber
    participant API as QAService.stream_query()
    participant RAG as rag.stream_query()
    participant Ctx as create_query_context()
    participant Route as decide_route()
    participant Prep as prepare_retrieval()
    participant Prompt as build_answer_prompt_profile()
    participant FAQ as search_faq()
    participant Doc as search_doc()
    participant Ans as prepare_answer()
    participant LLM as stream_llm_answer()
    participant Hist as history.add_turn()
    participant Trace as finish_success()

    API->>RAG: yield from rag_stream_query(...)
    RAG->>Ctx: 创建请求上下文
    RAG->>Route: decide_route(context)
    alt direct_answer / faq_exact
        Route-->>RAG: RouteDecision + answer
        RAG->>Hist: _finish_with_single_answer()
        RAG->>Trace: finish_success()
    else retrieval
        Route-->>RAG: route=retrieval
        RAG->>Prep: prepare_retrieval(context)
        Prep->>Prompt: 根据意图和问题类别选择 Profile
        Prompt-->>Prep: PromptProfile
        Prep-->>RAG: RetrievalPreparation
        RAG->>FAQ: search_faq(context, prepared)
        alt FAQ 直出
            FAQ-->>RAG: direct_answer
            RAG->>Hist: _finish_with_single_answer()
            RAG->>Trace: finish_success()
        else 继续检索
            RAG->>Doc: search_doc(context, prepared)
            RAG->>Ans: prepare_answer(context, prepared, faq_result, doc_result)
            alt 信息不足
                Ans-->>RAG: insufficient_context
                RAG->>Hist: _finish_with_single_answer()
                RAG->>Trace: finish_success()
            else 进入生成
                Ans-->>RAG: AnswerPreparation
                RAG->>LLM: stream_llm_answer(system_prompt, user_prompt)
                LLM-->>RAG: token chunks
                RAG->>Hist: history.add_turn()
                RAG->>Trace: finish_success()
            end
        end
    end

读图重点: 1.decide_route()先决定是直答、FAQ 精确命中,还是进入完整检索链路。 2. FAQ 直出和信息不足都会提前收口,不进入 LLM 生成。 3. 只有检索准备、FAQ/文档召回都完成后,才会进入stream_llm_answer()

单独查看:打开可缩放时序图。

eline)。

2.4 主线代码骨架(先按此顺序阅读)

# qa_core/pipeline/rag.py
def stream_query(history, query, source_filter, session_id, ...):
    # === Stage 0: 创建运行上下文 ===
    context = create_query_context(...)
    yield build_query_start_event(context)  # start 事件

    try:
        # === Stage 1: 查询路由 ===
        yield build_status_event("正在进行查询路由...", context.session_id)
        route = decide_route(context)
        if route.answer:
            yield from _finish_with_single_answer(context, history, query, route.answer)
            return

        # === Stage 2: 检索准备 + Prompt Profile 选择 ===
        yield build_status_event("正在识别问题意图...", context.session_id)
        prepared = prepare_retrieval(context)

        # === Stage 3-6: FAQ 检索 → 文档检索 → 上下文/Prompt 构建 → LLM 生成 ===
        helper_result = yield from _search_and_generate(context, prepared, query, history)
        if helper_result is None:
            return  # 已在内部收尾(FAQ 直出或信息不足)

        # === Stage 6 continuation: 引用补强与生成后核验 ===
        answer = enforce_answer_citations(context.answer, helper_result.context_docs)
        finalize_generated_answer_confidence(
            context,
            answer=answer,
            context_docs=helper_result.context_docs,
        )

        # === Stage 7: 保存历史 + 写入 Trace + 结束事件 ===
        history.add_turn(context.session_id, query, answer)
        yield finish_success(context, answer=answer)

    except Exception as exc:
        yield finish_error(context, exc)

def _search_and_generate(context, prepared, query, history):
    """检索-生成核心链路:FAQ 检索 → 文档检索 → 上下文构建 → LLM 流式生成。

    提取为独立函数使 stream_query 主干更清晰,便于单步调试和异常定位。
    """
    # Stage 3: FAQ 检索 + 直出判断
    yield build_status_event("正在检索业务 FAQ 知识库...", context.session_id)
    faq_result = search_faq(context, prepared)
    direct_answer = get_faq_direct_answer(context, prepared, faq_result)
    if direct_answer:
        yield from _finish_with_single_answer(context, history, query, direct_answer)
        return None

    # Stage 4: 文档检索
    yield build_status_event("正在匹配相关业务资料...", context.session_id)
    doc_result = search_doc(context, prepared)

    # Stage 5: 上下文构建
    answer_prepared = prepare_answer(context, prepared, faq_result, doc_result)
    context.sources = answer_prepared.sources
    context.hit_type = answer_prepared.hit_type

    if context.hit_type == "insufficient_context":
        answer = build_insufficient_context_answer(context)
        yield from _finish_with_single_answer(context, history, query, answer)
        return None

    # Stage 6: LLM 流式生成
    yield build_status_event("正在生成回答...", context.session_id)
    for chunk in stream_llm_answer(answer_prepared.system_prompt, answer_prepared.user_prompt):
        token = str(getattr(chunk, "content", "") or "")
        if not token:
            continue
        yield build_token_event(token, context.session_id)

    return answer_prepared

2.5 Stage 1:查询路由

这是在线问答进入检索准备之前的低成本路由层。它统一处理三类结果:

route intent 含义
direct_answer GREETING/HUMAN_SERVICE/OUT_OF_SCOPE 问候、转人工、越界、source 边界,直接返回
faq_exact FAQ_QUERY FAQ 标准问题精确命中,直接返回标准答案
retrieval 暂不确定 路由不了,进入检索准备

这里的关键点是:intent 描述用户想做什么,route 描述系统下一步怎么处理。FAQ 精确命中不是新的用户意图,而是route=faq_exact,同时携带intent=FAQ_QUERY

@dataclass
class RouteDecision:
    route: Literal["direct_answer", "faq_exact", "retrieval"]
    answer: str | None = None
    intent: IntentResult | None = None
    reason: str = ""

def decide_route(context):
    # 1. 先校验 source_filter
    context.run_stage("validate_source", ...)

    # 2. 协议/安全类直答:问候、越界、短句转人工
    direct_intent = classify_direct_intent(context.query, context.scenario)
    if direct_intent:
        return RouteDecision("direct_answer", direct_intent.direct_answer, direct_intent)

    # 3. source 边界
    boundary_answer = detect_and_apply_boundary_answer(context)
    if boundary_answer:
        return RouteDecision("direct_answer", boundary_answer, out_of_scope_intent)

    # 4. FAQ 精确命中:route 是 faq_exact,intent 仍是 FAQ_QUERY
    if should_try_faq_fast_path(context.query, context.scenario):
        answer, intent = try_fast_faq_direct_answer(context)
        if answer:
            return RouteDecision("faq_exact", answer, intent)

    # 5. 路由不了,再进入检索准备
    return RouteDecision("retrieval")

这也是你截图里最应该调整的地方:你好转人工彩票怎么买这类问题不应该先进入 FAQ 快速路径,而应该在这一阶段直接收口。

2.6 FAQ 精确命中为什么放在路由层

FAQ 精确命中依赖知识库内容、版本、tenant、source_filter 和标准问题文本,它不是“用户意图类型”。所以更准确的表达是:查询路由层可以产出route=faq_exact,并把intent标记为FAQ_QUERY

# qa_core/pipeline/steps.py
from qa_core.config.rules import get_rule_config

def should_try_faq_fast_path(query, scenario):
    """判断短问题是否值得先做 FAQ 精确匹配探测。

    快速路径只处理"短、完整、像标准问答"或能推断业务分类的问题。
    不是语义答案缓存:会先查带版本和权限边界的检索缓存,未命中时访问当前场景的 FAQ Milvus 集合,
    并带上 kb_version、tenant、dataset、visibility 和 role 过滤。
    返回 True 只代表可以先探测 FAQ 候选,不代表已经可以直出。
    """
    rules = get_rule_config().faq_fast_path
    compact_query = (query or "").strip()
    if (
        not compact_query
        or len(compact_query) > rules.max_chars
        or "\n" in compact_query
    ):
        return False  # 长问题、多行问题不适合快速路径
    return bool(
        rules.hint_matches(compact_query)  # FAQ 句式特征
        or infer_source(compact_query, scenario)  # 明确业务分类
    )

def try_fast_faq_direct_answer(context):
    """路由层的 FAQ 精确试探:只允许精确匹配,不允许相似直出。"""
    faq_store = get_faq_store(context.scenario.faq_collection)
    result = faq_store.search_many(
        [context.query],
        # 原问题候选随后可能供完整 FAQ 链路复用,容量不能小于计划所需。
        k=max(plan.faq_top_k, get_settings().faq_short_query_top_k),
        source_filter=effective_source_filter,
        kb_version=context.active_kb_version,
        data_scope=context.data_scope,
        source_type="faq",
        rerank=False,
    )
    # 只允许精确匹配,分数阈值设为无穷大
    answer, _ = _exact_faq_answer(context.query, result)
    return answer, intent  # 不是精确匹配就返回 (None, FAQ_QUERY intent),继续主流程

FAQ 快路径的触发词和最大长度来自config/rules.toml中的faq_fast_path配置,不写死在代码里。max_chars = 48是本项目的初始保护阈值,不是官方标准。它的作用是把 FAQ 快路径限制在“一句话标准问法”上:短问题先试精确命中;长问题、多行问题、带多个条件的问题交给后面的完整检索链路处理。

诊断信息里的retrieval.plan会按 FAQ 快路径的实际执行方式展示:run_faq=truerun_doc=falsererank=false,并标记match_policy=standard_question_exact。这样 Trace 里看到的计划和真实执行链路一致,不会误以为 Stage 1 也进入了完整文档检索。

这里还有一个边界要注意:FAQ 快路径只做“是否精确命中标准 FAQ”的判断,不选择 Prompt Profile。因为精确命中会直接返回标准答案,不调用 LLM,也不需要构造回答 Prompt。Prompt Profile 的选择发生在后面的prepare_retrieval()/prepare_answer(),也就是完整检索和生成路径里。

_exact_faq_answer()— 精确匹配实现

def _exact_faq_answer(query: str, faq_result: RetrievalResult) -> tuple[str | None, RetrievalResult]:
    """从 FAQ 候选中找与标准问题完全一致的答案。

    快速路径只允许精确标准问答直出,不按相似分数直出。
    找到精确命中后会把该命中排到来源列表第一位,方便页面展示。
    """
    for index, hit in enumerate(faq_result.hits):
        answer = direct_faq_answer(query, hit.document, hit.score, threshold=float("inf"))
        if not answer:
            continue
        if index:
            reordered = [hit, *faq_result.hits[:index], *faq_result.hits[index + 1 :]]
            faq_result = RetrievalResult(
                hits=reordered,
                query=faq_result.query,
                source_type=faq_result.source_type,
                elapsed_ms=faq_result.elapsed_ms,
            )
        return answer, faq_result
    return None, faq_result

为什么在检索准备之前做

  • 减少首 token 延迟。标准 FAQ 的精确命中不需要经过历史加载、检索类意图识别、改写、检索计划等步骤。
  • FAQ 快速路径可以命中 Redis 检索缓存;缓存未命中时仍然访问 Milvus,并且始终带版本和数据隔离过滤。
  • 它在同一个decide_route()中排在 direct_answer 之后,避免问候、转人工、越界问题先触发知识库查询。

为什么只允许精确匹配

  • 还没做意图识别,不知道这是 FAQ_QUERY 还是 KNOWLEDGE_QUERY
  • 如果是知识咨询但 FAQ 相似分数高,可能误答。所以只允许用户问题和 FAQ 标准问题完全一致时才直出。

2.6.1 FAQ 快速探测未命中后的候选复用

先用一个具体例子理解。用户问“新人入职需要完成哪些流程?”。路由层为了判断能否精确 FAQ 直出,已经用原问题查过一次 FAQ;结果没有找到完全相同的标准问题,于是进入完整 RAG。完整链路生成的变体可能是:

["新人入职需要完成哪些流程?", "新人入职需要完成哪些 SOP?", "新员工需要办理哪些入职手续?"]

当前请求上下文会保存fast_faq_result、对应的source_filter和已取回的候选容量。复用条件成立时,完整检索读取原问题候选,只向 FAQ collection 查询两个新增变体:

快速探测:原问题 -> FAQ 原始候选
完整检索:读取原问题候选 + 只查询后两个新增变体
        -> 按 FAQ/chunk 去重 -> 一次统一 CrossEncoder 重排 -> 取 top_k

这不是 Redis 中跨请求复用“整段回答”,也不会跳过后面的意图识别、来源过滤、引用补强或 LLM 生成。它只是在同一次请求里复用已经得到的原始 FAQ 候选,因此不会受到对话历史不同、Prompt 不同或最终答案不同的影响。

为保证正确性,只有以下条件同时满足才允许复用:

  1. rewritten_query仍等于原问题,说明没有追问改写。
  2. query_variants的第一个元素仍是原问题。
  3. 快速探测与完整计划的有效source_filter相同。
  4. 快速探测的候选容量不少于完整计划的faq_top_k

如果是追问、来源过滤发生变化,或者完整计划需要更多候选,系统会正常执行完整search_many(query_variants)。这是性能优化的边界:不能以少查一次为理由复用不等价的数据。

第三部分:Stage 2 检索准备与 Prompt Profile

这一部分对应 Stage 2 的prepare_retrieval()。它先根据意图、问题类别和场景选择 Prompt Profile,但此时还没有把检索结果填入 User Prompt。

3.1 Prompt 在在线链路中的位置

完整检索路径的调用主线如下:

stream_query()
  -> prepare_retrieval()
      -> classify_intent()
      -> rewrite_query_if_needed()
      -> build_retrieval_plan()
      -> generate_query_variants()
      -> build_answer_prompt_profile()
  -> search_faq() / search_doc()
  -> prepare_answer()
      -> select_context_docs()
      -> build_context()
      -> 填充 user_template
  -> stream_llm_answer(system_prompt, user_prompt)
  -> enforce_answer_citations()
  -> finalize_generated_answer_confidence()
  -> finish_success()

这里有两个容易混淆的时间点:

  1. build_answer_prompt_profile()prepare_retrieval()中完成。它根据意图和问题类别先确定回答策略。
  2. prepare_answer()在检索结果返回后,才把历史、改写后的问题和最终上下文填入user_template,形成真正发送给 LLM 的user_prompt

FAQ 精确路由、问候、越界和转人工等确定性分支不会调用 LLM,因此不会进入最终回答 Prompt 的生成流程。FAQ 相似候选没有直接结束时,才会进入完整检索链路并选择 Prompt Profile。

3.2 PromptProfile 的数据结构

项目没有把一长段 Prompt 直接散落在rag.py中,而是使用不可变的PromptProfile统一描述一个回答档位:

# qa_core/prompts/profiles.py
from dataclasses import dataclass

@dataclass(frozen=True)
class PromptProfile:
    name: str
    system_template: str
    user_template: str
    reason: str

    def as_dict(self) -> dict[str, str]:
        # Trace 只记录名称和选择原因,不把完整 Prompt 写入诊断日志
        return {"name": self.name, "reason": self.reason}

四个字段的职责是:

字段 作用
name 当前使用的模板档位名称,例如knowledge_answer
system_template 角色、回答边界、引用规则和风险约束
user_template 历史、问题和检索上下文的拼接格式
reason 解释为什么选中该档位,供 Trace 和调试使用

frozen=True是为了避免请求处理中途修改模板。模板属于经过评审的业务规则,变更应该通过代码或配置发布,而不是在一次请求里被悄悄改写。

3.3 模板注册与选择优先级

当前模板分成三个层次:

层次 注册位置 作用
风险类别模板 CATEGORY_PROMPT_PROFILES 费用、合规、排障、总结等问题使用更严格的回答口径
意图模板 PROMPT_PROFILES FAQ、知识查询、追问使用对应的回答结构
默认模板 DEFAULT_PROMPT_PROFILE 新增意图或未知类别的安全兜底

选择优先级固定为:

问题类别专用模板
        ↓ 未命中
意图专用模板
        ↓ 未命中
默认安全模板

当前主要回答档位如下:

问题或意图 Profile 业务目的
费用、价格、退款、发票 pricing_guard 不估算金额,区分已确认和待人工确认
合同、合规、隐私、审计 compliance_guard 只引用依据,不自行判断合规结论
故障排查、异常处理 troubleshooting_steps 按现象、原因、步骤和升级路径回答
资料总结 source_bound_summary 只总结上下文已有内容
FAQ_QUERY faq_answer 短、准、直接复用 FAQ 口径
KNOWLEDGE_QUERY knowledge_answer 对制度、流程和知识资料进行结构化回答
FOLLOW_UP follow_up 结合历史,但只回答当前追问焦点
未知意图 default_answer 通用安全兜底

风险类别优先于意图。例如“报销超过 5000 需要谁审批”可能被意图识别为FAQ_QUERY,但它同时属于费用类问题,因此最终使用pricing_guard,不能只使用普通 FAQ 模板。

3.4 build_answer_prompt_profile() 的实际实现

# qa_core/prompts/selector.py
def build_answer_prompt_profile(
    intent: str,
    scenario: ScenarioDefinition,
    query: str,
) -> PromptProfile:
    question_category = infer_question_category(query)

    profile = (
        CATEGORY_PROMPT_PROFILES.get(question_category)
        or PROMPT_PROFILES.get(intent, DEFAULT_PROMPT_PROFILE)
    )

    context = _scenario_prompt_context(scenario)
    return PromptProfile(
        name=profile.name,
        system_template=profile.system_template.format(**context),
        user_template=profile.user_template,
        reason=profile.reason,
    )

这里没有再次调用 LLM 选择模板。模板选择是确定性的规则判断,原因有三点:

  1. 检索策略和回答口径必须一致,不能检索阶段认为是 FAQ,生成阶段又随机改成综合分析模板。
  2. 少一次 LLM 判断,减少延迟、成本和不稳定性。
  3. 费用、合规等高风险问题必须稳定使用保守模板,不能把模板选择交给模型自由判断。

3.5 场景变量注入

模板本身不写死具体业务名称,而是在选择完成后注入当前场景配置:

# qa_core/prompts/selector.py
def _scenario_prompt_context(scenario: ScenarioDefinition) -> dict[str, str]:
    return {
        "assistant_name": scenario.assistant_name,
        "business_domain": scenario.business_domain,
        "industry": scenario.industry,
        "support_contact": scenario.support_contact,
        "phone": scenario.support_contact,
    }

因此,同一套knowledge_answer模板可以用于人事、财务、IT 支持等场景,只需要从各自的scenario.toml注入不同的助手名称、业务域和人工支持联系方式。

最终回答 Prompt 分为两部分:

System Prompt:定义助手身份、回答边界、引用要求和风险约束
User Prompt:历史消息 + 改写后的问题 + 编号上下文

静态模板统一放在qa_core/prompts/constants.py。所有回答模板都会强调:只能基于上下文回答、不能编造不存在的业务事实、引用编号必须来自当前上下文、资料不足时必须明确说明无法确认。


第四部分:Stage 3-4 FAQ 与文档检索

这一部分承接 Stage 2 的RetrievalPreparation,说明 FAQ 候选如何决定直出,以及未直出时如何进入文档检索。

4.1 FAQ 标准直出 vs FAQ 精确路由

这是两个容易混淆的概念:

FAQ 精确路由(Stage 1) FAQ 标准直出(Stage 3)
时机 decide_route()中,检索准备之前 检索准备之后
route / intent route=faq_exact,携带intent=FAQ_QUERY route=retrieval,保持原始检索意图:FAQ_QUERY/FOLLOW_UP/KNOWLEDGE_QUERY
FAQ 关系 只试 FAQ collection 的标准问题精确匹配 只要RetrievalPlan.run_faq=True,就会执行 FAQ collection 检索
未精确命中 保存原问题候选;完整链路可按条件补查变体 合并原问题候选与变体候选后,再决定是否标准直出
匹配方式 仅精确匹配 精确匹配 + 相似分数阈值
阈值 ∞(只精确) 动态阈值;追问、低决策分、表格类问题会更保守
适用 短标准问答 FAQ 候选足够可靠;不要求意图必须是FAQ_QUERY
风险 低(只精确) 中(相似分数可能误命中)

注意:Stage 3 的“FAQ 检索”不是“把意图改成FAQ_QUERY”。例如用户追问“那审批呢?”时,如果存在历史上下文,检索类意图仍然是FOLLOW_UP;后续只是因为检索计划里run_faq=True,所以会同时查 FAQ collection。FAQ top 命中足够可靠时可以直接返回标准答案,否则继续查文档并进入生成链路。


4.2 Stage 4:文档检索

FAQ 没有达到标准直出条件时,主线才进入search_doc(context, prepared)。这个函数仍然消费 Stage 2 生成的RetrievalPreparation,因此文档检索不会重新判断意图、来源范围或 Prompt Profile。

# qa_core/pipeline/retrieval_steps.py
doc_result = search_doc(context, prepared)

它的执行顺序可以概括为:先按kb_version、DataScope、source 和检索计划检查安全缓存;缓存未命中时执行 Dense + Sparse Hybrid 检索;计划开启rerank时再进行 CrossEncoder 精排;最后把候选、最高分、耗时和缓存命中情况写回context.retrieval_info。Stage 4 只负责拿回候选,不负责筛选最终 Prompt 上下文,最终筛选发生在 Stage 5 的prepare_answer()


第五部分:Stage 5 上下文构建与 Prompt 组装

这一部分对应 Stage 5 的prepare_answer():先筛选和格式化上下文,再把历史、问题和证据填入 Stage 2 已选定的 Prompt Profile,形成真正发送给 LLM 的 User Prompt。

5.1 select_context_docs() 的筛选策略

# qa_core/pipeline/context.py
def select_context_docs(faq_hits: list, doc_hits: list, plan: RetrievalPlan) -> list[Document]:
    """筛选进入 Prompt 的文档片段(只依赖 plan 对象,不需要 scenario 参数)。

    执行流程:
    1. FAQ 命中:过滤分数 → 取前 2 条 → 转成"常见问题 + 标准答案"格式
    2. 文档命中:过滤分数 → prefer_table 时表格行优先 → 优先用 parent_content
    3. 每条追加受 final_context_top_n / max_context_chars / max_context_doc_chars 三重约束
    """
    selected = []
    seen_keys = set()
    used_chars = 0

    # ── FAQ 部分:过滤 min_context_score → 取前 2 条 → 转成标准问答格式 ──
    for hit in [h for h in faq_hits if h.score >= plan.min_context_score][:2]:
        answer = hit.document.metadata.get("answer")
        question = hit.document.metadata.get("standard_question") or hit.document.page_content
        if answer:
            _append_with_budget(
                Document(page_content=f"常见问题:{question}\n标准答案:{answer}"),
                f"faq:{document_key(hit.document)}",
                selected, seen_keys, used_chars, plan)

    # ── 文档部分:过滤分数 → prefer_table 排序 → 优先用 parent_content ──
    eligible = [h for h in doc_hits if h.score >= plan.min_context_score]
    if plan.prefer_table:
        # 表格行(content_type 以 table 开头)排到普通正文前面
        eligible = sorted(eligible,
            key=lambda h: (0 if is_table_document(h.document) else 1, -h.score))

    for hit in eligible:
        parent_content = hit.document.metadata.get("parent_content")
        key = str(hit.document.metadata.get("parent_id") or document_key(hit.document))
        _append_with_budget(
            Document(page_content=str(parent_content or hit.document.page_content)),
            f"doc:{key}",
            selected, seen_keys, used_chars, plan)
    return selected

5.2 build_context() 的格式化输出

def build_context(docs: list[Document]) -> str:
    """构建最终上下文文本(只依赖 doc.metadata,不依赖 scenario 对象)。"""
    lines = []
    seen: set[str] = set()
    for i, doc in enumerate(docs):
        content = doc.page_content.strip()
        if not content or content in seen:
            continue  # 内容去重
        seen.add(content)
        source = _context_source_label(doc.metadata or {})

        # 格式:[编号] 来源:文件名 或 标准问题名 或 表格 sheet+行号
        header = f"[{i+1}] 来源:{source}"
        lines.append(f"{header}\n{content}")

    return "\n\n".join(lines)

输出示例:

[1] 来源:人事制度 / 入职管理
入职流程包括以下步骤:1. 提交入职材料(身份证复印件、学历证书...)

[2] 来源:人事制度 / 审批权限
部门经理负责审批本部门员工的入职申请,审批时限为 3 个工作日...

[3] 来源:行政管理 / 工位分配
新员工入职后由行政部统一分配工位和办公设备...

5.3 Prompt 组装与下游调用

prepare_retrieval()负责选择 Profile,并把它放入RetrievalPreparation

# qa_core/pipeline/steps.py
prompt_profile = context.run_stage(
    "select_prompt_profile",
    lambda: build_answer_prompt_profile(
        intent.intent,
        context.scenario,
        context.rewritten_query,
    ),
)

检索完成后,prepare_answer()才把真实上下文填入模板:

user_prompt = prepared.prompt_profile.user_template.format(
    history=format_messages(prepared.history_messages),
    question=prepared.rewritten_query,
    context=build_context(context_docs)
    or "无可用上下文。必须明确回答:信息不足,无法确认。",
)

return AnswerPreparation(
    context_docs=context_docs,
    sources=sources,
    hit_type=hit_type,
    system_prompt=prepared.prompt_profile.system_template,
    user_prompt=user_prompt,
)

随后由stream_llm_answer()将这两个字符串转换成SystemMessageHumanMessage,开始流式生成。第 10 章后面的引用增强、生成后核验和答案置信度,都是针对这个 Prompt 生成结果继续执行的。

5.4 Prompt 选择的验证方式

当前模块示例通过两种方式验证 Prompt:

python scripts/demo_rag_pipeline.py "公司入职流程文档在哪里查看" --source hr --debug
python -m unittest discover -s tests

--debug输出中的retrieval_plan.prompt_profile会显示最终使用的模板名称和选择原因;完整流式问答的end事件也会保留 Prompt Profile 诊断信息。测试重点验证风险类别优先于意图、知识查询使用knowledge_answer、历史追问仍然保持改写和查询变体闭环。


第六部分:Stage 5 信息不足与生成前证据置信度

Stage 5 的职责是确认“有没有足够可靠的证据进入 Prompt”。它会先筛选上下文;如果上下文不足就确定性收口,如果上下文足够则记录生成前证据置信度,随后才允许进入 Stage 6。

6.1 什么情况判定为信息不足

prepare_answersteps.py中定义,其内部的信息不足判定委托给_build_answer_context

# qa_core/pipeline/steps.py
def prepare_answer(
    context: RAGQueryContext,
    prepared: RetrievalPreparation,
    faq_result: RetrievalResult,
    doc_result: RetrievalResult,
) -> AnswerPreparation:
    """将 FAQ + 文档检索结果整理为 LLM Prompt、引用来源列表和命中类型。

    信息不足判定委托给 _build_answer_context,prepare_answer 负责
    组装最终的 system_prompt 和 user_prompt。
    """
    context_docs, sources, hit_type, top_score = context.run_stage(
        "build_answer_context",
        lambda: _build_answer_context(prepared, faq_result, doc_result),
    )
    _record_context_stats(context, context_docs, prepared.plan, top_score)
    # 这里只记录生成前证据置信度;LLM 输出后会在 rag.py 中执行生成后核验并合并。
    record_evidence_confidence(
        context,
        hit_type=hit_type,
        retrieval_top_score=top_score,
        context_count=len(context_docs),
        source_count=len(sources),
    )

    user_prompt = prepared.prompt_profile.user_template.format(
        history=format_messages(prepared.history_messages),
        question=prepared.rewritten_query,
        context=build_context(context_docs)
        or "无可用上下文。必须明确回答:信息不足,无法确认。",
    )
    return AnswerPreparation(
        context_docs=context_docs,
        sources=sources,
        hit_type=hit_type,
        system_prompt=prepared.prompt_profile.system_template,
        user_prompt=user_prompt,
    )

_build_answer_context负责实际的上下文筛选和命中类型判定:

def _build_answer_context(prepared, faq_result, doc_result):
    """整理上下文文档、引用来源列表、命中类型和最高分数。

    无上下文通过分数过滤时命中类型标记为 insufficient_context。
    """
    context_docs = select_context_docs(faq_result.hits, doc_result.hits, prepared.plan)
    if prepared.plan.prefer_table:
        sources = doc_result.source_payloads(limit=5) + faq_result.source_payloads(limit=2)
    else:
        sources = faq_result.source_payloads(limit=2) + doc_result.source_payloads(limit=5)
    top_score = max(faq_result.top_score, doc_result.top_score)
    return context_docs, sources, "rag" if context_docs else "insufficient_context", top_score

6.2 信息不足的答案

def build_insufficient_context_answer(context: RAGQueryContext) -> str:
    """无可用上下文时返回确定性"信息不足"回答,避免 LLM 幻觉。"""
    context.retrieval_info["insufficient_context_reason"] = "no_context_after_score_filter"
    return f"信息不足,无法确认。当前知识库没有召回到足够可靠的依据,请联系{context.scenario.support_contact}。"

设计意图:信息不足时,系统明确告知用户(而不是让 LLM 即兴发挥),避免 LLM 在没有可靠资料的情况下生成"幻觉"答案。


6.3 生成前证据置信度

目前答案置信度采用“两阶段评估”设计:

证据置信度
    +
生成结果校验
    ↓
最终答案置信度

calculate_evidence_confidence()会在答案生成前判断当前证据是否足以支撑回答。它主要参考以下信号:

信号 含义
retrieval_top_score FAQ / Doc 排序后第一条候选的检索相关性分
context_count 最终进入 Prompt 的上下文条数
source_count 返回给前端的来源数量,目前只作为诊断信号
intent_rule_score 入口规则候选强弱,用于诊断
intent_decision_score 意图网关最终决策分,用于计算证据置信度
history_rewrite_used 是否依赖历史追问改写
hit_type FAQ 直出、RAG、信息不足或确定性直答

这里的retrieval_top_score不是所有候选的平均分,也不是进入 Prompt 的文档数量。它表示当前最强的一条候选证据与问题的匹配程度。例如:

候选1 score=0.91
候选2 score=0.78
候选3 score=0.62

这时 top-1 检索分就是0.91。当前实现不会把上下文数量、来源数量、答案长度等信号做复杂加权,而是采用“路径优先、分数分档、风险降级”的规则。这样更容易解释,也能避免单纯堆更多低相关上下文把置信度抬高。

检索分先转换为normalized_retrieval_score

score <= 0:0
0 < score <= 1:保持原值
score > 1:1 - 1 / (1 + score)

第三种情况用于平滑压缩 CrossEncoder 可能输出的大于 1 的 logit。例如3.0被压缩为0.75,而不是粗暴截断成1.0。该归一化只用于答案置信度,不改变候选排序。

当前基础分档如下,具体逻辑对应qa_core/pipeline/confidence.py中的_base_evidence_score()

场景 证据置信度
FAQ 精确匹配 0.95
确定性业务路由 0.90
RAG 检索强,且意图稳定 0.85
RAG 检索中等 0.65
RAG 检索较弱 0.45
没有检索到上下文 0.20

除此之外,还会进行风险降级:

  • 使用历史问题改写时,证据分最高限制为0.65
  • 意图决策分低于0.70时,非确定性路由且非 FAQ 精确命中的结果最高限制为0.54
  • 没有最终上下文时,证据分最高限制为0.20

因此,source_count目前不会直接给分数加分;context_count会影响系统是否判定为“没有可用证据”。这组分数是工程诊断分档,不是模型训练得到的概率。


第七部分:Stage 6 流式生成、引用增强与生成后核验

Stage 6 的顺序是固定的:先让 LLM 流式输出,再对完整答案补强引用,最后依据最终文本做生成后核验。生成前已经记录的证据分不会在这里被覆盖,而是与生成核验结果保守合并。

stream_llm_answer()
  -> enforce_answer_citations()
  -> finalize_generated_answer_confidence()

7.1 LLM 流式生成

# qa_core/pipeline/rag.py::_search_and_generate()
with context.stage("llm_generation"):
    for chunk in stream_llm_answer(
        answer_prepared.system_prompt,
        answer_prepared.user_prompt,
    ):
        token = str(getattr(chunk, "content", "") or "")
        if not token:
            continue
        context.answer_parts.append(token)
        context.mark_first_token()
        yield build_token_event(token, context.session_id)

这一步只负责把模型输出逐段交给前端,并把原始回答累计到context.answer_parts。此时模型可能遗漏引用或表格关键字段,因此不能立刻结束请求。

7.2 答案引用增强

LLM 生成的答案可能引用上下文中的信息,但不会自动标注"这个信息来自哪个文档"。引用增强发生在 Stage 6 的流式生成完成之后、生成后置信度核验之前:它先补齐可见来源标注,再把最终回答交给calculate_generation_confidence()检查。

# qa_core/pipeline/citations.py
CITATION_RE = re.compile(r"\[\d+\]")

def has_source_citation(answer: str) -> bool:
    """判断答案中是否已经包含 [数字] 形式的来源编号。"""
    return bool(CITATION_RE.search(answer))

def source_reference_label(doc: Document, index: int) -> str:
    """生成简短来源标签(文件名/FAQ 标准问题;表格资料附加 sheet 和行号)。"""
    from qa_core.document_metadata import format_source_label
    return f"[{index}] {format_source_label(dict(doc.metadata or {}))}"

def enforce_table_row_details(answer: str, context_docs: list[Document]) -> str:
    """确保表格类答案在模型遗漏关键单元格时,确定性追加表格行要点。"""
    details = []
    for index, doc in enumerate(context_docs, start=1):
        if not is_table_document(doc) or not needs_table_row_detail(answer, doc):
            continue
        detail = build_table_row_detail(doc, index)
        if detail:
            details.append(detail)
        if len(details) >= 1:
            break
    if not details:
        return answer
    return f"{answer}\n\n" + "\n".join(details)

def enforce_answer_citations(answer: str, context_docs: list[Document]) -> str:
    """确保 RAG 答案带有可见来源编号:模型已写则保留,未写则末尾补充前 3 个来源。

    额外检查表格类答案:模型遗漏核心单元格值(状态/金额/责任人等)时
    确定性追加表格行要点,避免 LLM 忽略关键数据。
    """
    clean_answer = answer.strip()
    if not clean_answer or not context_docs:
        return clean_answer  # 空答案或空来源时原样返回,不阻断流程

    # 确保表格类答案不丢失关键单元格信息
    clean_answer = enforce_table_row_details(clean_answer, context_docs)
    # 答案已包含 [数字] 来源编号时不重复添加
    if has_source_citation(clean_answer):
        return clean_answer

    # 为前 3 个上下文文档生成来源标签(文件名/FAQ 问题名/表格 sheet 和行号)
    references = ";".join(
        source_reference_label(doc, index)
        for index, doc in enumerate(context_docs[:3], start=1)
    )
    return f"{clean_answer}\n\n参考来源:{references}"

7.3 生成后答案核验

如果最终走了 LLM 生成,rag.py会在enforce_answer_citations()之后调用calculate_generation_confidence()。这一步不再看召回分,而是看最终答案文本本身:

核验信号 含义
citation_coverage 有多少事实单元带有有效行内引用
valid_citation_numbers 答案引用的编号是否落在当前上下文范围内
invalid_citation_numbers 是否编造了不存在的[N]来源编号
context_overlap 答案事实单元与上下文的词面支撑比例
answer_char_count 答案是否为空或异常短

这里有一个重要边界:末尾自动追加的参考来源:[1] ...不算行内引用。它只能说明答案整体附带了来源列表,不能说明每个事实句都有依据。所以生成核验会先把参考来源尾部移除,再检查正文里的[1][2]

当前生成核验的结果分档如下:

生成结果 生成置信度
引用有效、引用覆盖率和上下文匹配都达标 0.85
只有引用覆盖率或上下文匹配其中一项达标 0.65
引用编号非法或依据不足 0.35
答案为空 0.00

其中,引用覆盖率达到0.50、上下文词面支撑达到0.45,才认为对应指标达标。上下文匹配属于轻量级词汇重合判断,不是严格的语义蕴含判断。

生成核验的状态有四种:

状态 含义
verified 引用覆盖和上下文支撑都较好
partial 有一定支撑,但引用覆盖或词面支撑不足
failed 答案为空、引用严重缺失或支撑明显不足
not_applicable 没有调用 LLM,如 FAQ 直出、确定性直答、信息不足兜底

7.4 最终合并策略

最终answer_confidencecombine_answer_confidence()生成。合并策略采用保守口径:

final_score = min(evidence_confidence.score, generation_verification.score)

如果本轮没有调用 LLM,generation_verification.status = not_applicable,最终分保留证据分。

这样设计是为了避免下面这种误判:

检索证据很强:0.95
LLM 答案没有行内引用,或引用了不存在的 [9]:0.35
最终 answer_confidence = 0.35

例如:

证据置信度 = 0.85
生成校验置信度 = 0.65
最终答案置信度 = 0.65

如果本轮没有调用 LLM,例如 FAQ 直接返回、确定性路由或上下文不足,系统会保留证据置信度,并将生成核验标记为:

{
  "generation_verification": {
    "status": "not_applicable",
    "score": null
  }
}

此时最终分数直接使用证据置信度,不伪造生成后的核验结果。也就是说,最终答案置信度通过“证据强度 + 生成结果校验 + 取较低值”来衡量答案可靠性,但它仍然不是机器学习意义上的答案正确概率。

最终分限制在[0, 1],并按以下区间展示:

high   >= 0.82
medium >= 0.55 且 < 0.82
low    < 0.55

7.5 答案置信度写入位置

路径 写入位置 说明
确定性直答 finish_success()兜底写入 只写证据分,生成核验为not_applicable
FAQ 直出 try_fast_faq_direct_answer()/_search_and_generate() 只写证据分,生成核验为not_applicable
信息不足 _finish_with_single_answer() 只写证据分,生成核验为not_applicable
RAG 生成 prepare_answer()+rag.py生成后收口 先写证据分,再生成后核验,最后合并

因此前端和 Trace 要这样读:

  • sources[*].score:单条来源的检索相关性排序分。
  • answer_confidence.score:最终答案的综合置信度。
  • answer_confidence.evidence_confidence.score:生成前证据置信度。
  • answer_confidence.generation_verification:生成后核验状态、分数和 reasons。
  • answer_confidence.reasons:为什么高或低,例如faq_exact_matchhistory_rewrite_usedinsufficient_context

注意:当前answer_confidence是可解释的工程信号,不是经过概率校准的“答案正确率”。它适合用于前端提示、Trace 排查、Bad Case 分层和评测报告观察;默认不替代 Recall@K、MRR、关键词覆盖率和人工复核,也不把历史聊天当成真值来源。历史只作为“是否依赖追问改写”的稳定性信号参与扣分。

后续如果要把它升级成更严格的质量门禁,需要先积累评测集、人工标注和线上反馈,再做分数校准。否则直接用一个固定置信度阈值阻断发布,容易把“有证据但表达保守”和“无证据乱答”混在一起。

第八部分:Stage 7 统一收口

Stage 7 不负责再检索或生成内容,而是保证所有分支都以一致的方式结束:保存最终答案、补齐诊断信息、写入 Trace,并向前端发出enderror事件。它是整条在线链路的出口。

8.1 成功路径:保存最终答案、发出 end、写入 Trace

完整 RAG 路径在 Stage 6 的引用增强和置信度核验完成后,才进入 Stage 7:

# qa_core/pipeline/rag.py
with context.stage("save_history"):
    history.add_turn(context.session_id, query, answer)

yield finish_success(context, answer=answer)

这里保存的是已经补齐引用后的最终答案,而不是 LLM 刚生成的原始文本。这样用户下一轮追问时,历史上下文与前端看到的内容一致。

finish_success()位于qa_core/pipeline/runtime.py,按固定顺序完成三件事:

1. 补齐未走 LLM 分支的 answer_confidence,并将生成核验标为 not_applicable
2. finalize_timings() 汇总阶段耗时,构造 WebSocket end 事件
3. record_trace() 写入本轮答案、检索诊断和耗时

因此 Stage 7 的业务意义不是“简单 return”,而是把本次请求的回答、来源、意图、检索诊断、答案置信度和耗时统一交付给前端与 Trace。

8.2 非 LLM 分支也走同一个出口

问候、转人工、越界、FAQ 直出和信息不足不进入 Stage 6 的 LLM 生成,但它们仍需保存历史、写 Trace、发送end事件。项目通过_finish_with_single_answer()复用这一收口:

# qa_core/pipeline/rag.py
context.answer_parts = [answer]
context.mark_first_token()
yield build_token_event(answer, context.session_id)
history.add_turn(context.session_id, query, answer)
yield finish_success(context, answer=answer)

区别只是答案已经确定,不需要逐 token 等待模型,所以一次性发出完整 token。finish_success()会发现本轮没有generation_attempted,将生成核验标为:

{
  "generation_verification": {
    "status": "not_applicable",
    "score": null
  }
}

最终answer_confidence保留前面记录的证据分,不伪造“模型回答已核验”的结果。

8.3 异常路径:finish_error()

任何 Stage 抛出异常,stream_query()都会进入统一异常收口:

# qa_core/pipeline/rag.py
except Exception as exc:
    logger.exception("QA stream failed")
    yield finish_error(context, exc)

finish_error()不会把原始异常直接暴露给用户。它先汇总已有耗时、记录失败 Trace,再返回经过user_facing_error_message()处理的error事件。这样前端可以提示用户重试,服务端则保留完整失败现场用于排查。

口语化总结:Stage 7 就像问答请求的交付台。前面不管走的是模型生成、FAQ 直出,还是信息不足,最后都要在这里把最终答案、诊断和 Trace 统一交出去;如果失败,也在这里把内部错误变成用户可理解的提示。


第九部分:跨阶段性能追踪

9.1 阶段计时

# qa_core/pipeline/runtime.py
class RAGQueryContext:
    """RAG 请求的运行时状态(dataclass,字段名与项目实现一致)。"""
    started: float           # 请求开始时间戳(time.perf_counter())
    stage_timings_ms: dict[str, float] = {}  # 各阶段耗时字典
    first_token_ms: float | None = None      # 首 token 到达时间(毫秒)

    @contextmanager
    def stage(self, name: str):
        """记录某个阶段的耗时。"""
        started = time.perf_counter()
        try:
            yield
        finally:
            self.record_stage(name, started)

    def record_stage(self, stage_name: str, started: float) -> float:
        """将阶段耗时写入 stage_timings_ms 字典。"""
        elapsed_ms = (time.perf_counter() - started) * 1000
        self.stage_timings_ms[stage_name] = round(elapsed_ms, 2)
        return elapsed_ms

    def mark_first_token(self):
        """记录首 token 时间(从请求 started 到首次推送 token 的毫秒数)。"""
        if self.first_token_ms is None:
            self.first_token_ms = round((time.perf_counter() - self.started) * 1000, 2)

追踪信息最终进入end事件:

{
    "type": "end",
    "retrieval": {
        "first_token_ms": 2478.7,
        "stage_timings_ms": [
            {"stage": "intent", "elapsed_ms": 320.5},
            {"stage": "faq_search", "elapsed_ms": 450.2},
            {"stage": "doc_search", "elapsed_ms": 680.1},
            {"stage": "context_build", "elapsed_ms": 15.3},
            {"stage": "llm_generation", "elapsed_ms": 3200.8},
            {"stage": "save_history", "elapsed_ms": 45.2}
        ],
        "slowest_stage": {"stage": "llm_generation", "elapsed_ms": 3200.8}
    }
}

这些数据帮助性能优化:如果文档检索阶段总是很慢,可能需要调整 top_k 或索引参数;如果 LLM 生成阶段很慢,可能需要换更快的模型或调整 max_tokens。


第十部分:跨阶段流式事件协议 — 前后端如何协作

10.1 事件驱动的问答模型

一次 RAG 问答不是"前端发请求 → 等 5 秒 → 收到完整答案"。实际的用户体验是:

前端发送问题 → 看到"正在进行查询路由..." →
看到"正在识别问题意图..." → 看到"正在检索 FAQ..." → 看到"正在匹配业务资料..." → 看到"正在生成回答..." → token 逐字出现 →
看到完整答案 + 来源引用

这就是事件驱动模型。后端通过 WebSocket 持续推送事件,前端根据事件类型更新 UI。

10.2 五种事件类型

sequenceDiagram
    participant Browser as 浏览器
    participant WS as /api/stream (WebSocket)
    participant QASvc as QAService
    participant Pipeline as RAG Pipeline

    Browser->>WS: {"query": "入职流程有哪些步骤", ...}
    WS->>QASvc: stream_query(...)
    QASvc->>Pipeline: 创建生成器

    Pipeline-->>WS: {"type": "start", "session_id": "...", "trace_id": "..."}
    WS-->>Browser: 问答已开始,记录 session_id

    Pipeline-->>WS: {"type": "status", "message": "正在进行查询路由..."}
    WS-->>Browser: 更新状态提示

    Pipeline-->>WS: {"type": "status", "message": "正在识别问题意图..."}
    WS-->>Browser: 更新状态提示

    Pipeline-->>WS: {"type": "status", "message": "正在检索业务 FAQ 知识库..."}
    WS-->>Browser: 更新状态提示

    Pipeline-->>WS: {"type": "status", "message": "正在匹配相关业务资料..."}
    WS-->>Browser: 更新状态提示

    Pipeline-->>WS: {"type": "status", "message": "正在生成回答..."}
    WS-->>Browser: 更新状态提示,准备接收 token

    loop LLM 流式生成
        Pipeline-->>WS: {"type": "token", "token": "入"}
        WS-->>Browser: 追加字符到答案区
        Pipeline-->>WS: {"type": "token", "token": "职"}
        WS-->>Browser: 追加字符到答案区
        Pipeline-->>WS: {"type": "token", "token": "流"}
        WS-->>Browser: 追加字符到答案区
    end

    Pipeline-->>WS: {"type": "end", "sources": [...], "hit_type": "rag", ...}
    WS-->>Browser: 渲染来源引用、展示诊断信息

10.3 每种事件的字段结构

start 事件— 请求已被接收:

{
    "type": "start",
    "session_id": "abc123",
    "trace_id": "xyz789",
    "scenario_id": "enterprise_knowledge",
    "scenario_name": "企业内部知识助手",
    "data_scope": {"tenant_id": "default", "dataset_id": "default"},
    "kb_version": "20260515_a1b2c3d4"
}

status 事件— 阶段性进度通知:

{
    "type": "status",
    "session_id": "abc123",
    "message": "正在检索业务 FAQ 知识库..."
}

前端通常将message显示为一个动态更新的状态栏或加载提示。

token 事件— 流式答案的片段:

{
    "type": "token",
    "session_id": "abc123",
    "token": "入"
}

每个 token 是一个或多个中文字符。前端将这些 token 逐个追加到答案区域,实现打字机效果。

end 事件— 问答完成:

{
    "type": "end",
    "session_id": "abc123",
    "hit_type": "rag",
    "answer": "入职流程包括以下步骤:1. 提交材料 ...",
    "sources": [
        {"file_name": "入职流程.md", "source": "hr", "score": 0.92},
        {"file_name": "FAQ", "standard_question": "入职需要哪些材料", "score": 0.88}
    ],
    "answer_confidence": {
        "score": 0.82,
        "level": "high",
        "label": "高",
        "reasons": ["rag_with_context", "generation_grounded"],
        "evidence_confidence": {
            "score": 0.92,
            "level": "high",
            "label": "高"
        },
        "generation_verification": {
            "score": 0.82,
            "status": "verified",
            "reasons": ["generation_grounded"],
            "signals": {
                "citation_coverage": 1.0,
                "context_overlap": 0.76,
                "valid_citation_numbers": [1, 2],
                "invalid_citation_numbers": []
            }
        },
        "signals": {
            "retrieval_top_score": 0.92,
            "normalized_retrieval_score": 0.92,
            "context_count": 4,
            "source_count": 7,
            "intent_rule_score": 0.84,
            "intent_decision_score": 0.84,
            "history_rewrite_used": false,
            "evidence_confidence_score": 0.92,
            "generation_verification_score": 0.82,
            "generation_verification_status": "verified",
            "citation_coverage": 1.0,
            "context_overlap": 0.76
        }
    },
    "intent": {
        "intent": "KNOWLEDGE_QUERY",
        "rule_score": 0.84,
        "confidence": 0.84,
        "reason": "strong_knowledge_rule"
    },
    "retrieval": {
        "plan": {"faq_top_k": 20, "doc_top_k": 20, "rerank": true},
        "query_variants": ["入职流程", "入职步骤", "入职办理流程"],
        "faq_elapsed_ms": 45.2,
        "doc_elapsed_ms": 120.5,
        "stage_timings_ms": {...},
        "first_token_ms": 350.8,
        "total_elapsed_ms": 4520.3
    },
    "processing_time": 4.52,
    "trace_id": "xyz789"
}

error 事件— 异常恢复:

{
    "type": "error",
    "session_id": "abc123",
    "error": "LLM 服务暂时不可用,请稍后重试。",
    "trace_id": "xyz789"
}

10.4 前端如何消费事件

// static/js/chat.js(简化逻辑)

const ws = new WebSocket(`ws://${location.host}/api/stream`);

ws.onmessage = (event) => {
    const data = JSON.parse(event.data);

    switch (data.type) {
        case "start":
            state.sessionId = data.session_id;
            state.traceId = data.trace_id;
            break;

        case "status":
            updateStatusBar(data.message);  // "正在检索 FAQ..."
            break;

        case "token":
            appendToAnswer(data.content);    // 追加到答案区
            break;

        case "end":
            renderSources(data.sources);     // 渲染来源引用
            renderDiagnostics(data.retrieval); // 展示检索诊断
            updateStatusBar("");             // 清除状态栏
            state.inProgress = false;
            break;

        case "error":
            showError(data.error);           // 显示错误提示
            state.inProgress = false;
            break;
    }
};

10.5 后端如何推进生成器

关键问题:QAService.stream_query()同步生成器(它内部顺序执行意图识别、Milvus 检索、本地 rerank 和 LLM 流式调用),但 WebSocket 路由是异步函数。如果直接调用next(iterator),事件循环会被阻塞。

解决方案:asyncio.to_thread将同步生成器的推进放到独立线程:

# qa_core/api/chat.py
stream = get_qa_service().stream_query(*context.service_args())

while True:
    # 在线程中推进同步生成器,不阻塞事件循环
    has_event, event = await asyncio.to_thread(_next_stream_event, stream)
    if not has_event or event is None:
        break
    await websocket.send_json(event)

    if event.get("type") in {"end", "error"}:
        if event.get("type") == "end":
            # 后台异步刷新历史摘要(不阻塞用户看到结果)
            _schedule_summary_refresh(session_id)
        break
flowchart TD
    subgraph MainThread["主线程(事件循环)"]
        WS["WebSocket 接收消息"]
        Send["发送事件到浏览器"]
        Schedule["调度后台摘要刷新"]
    end

    subgraph WorkerThread["工作线程"]
        Gen["推进同步生成器<br/>next(iterator)"]
        Intent["意图识别"]
        Milvus["Milvus 检索"]
        Rerank["本地重排"]
        LLM["LLM 流式调用"]
    end

    WS -->|"asyncio.to_thread"| Gen
    Gen --> Intent --> Milvus --> Rerank --> LLM
    LLM -->|"yield token"| Gen
    Gen -->|"返回事件"| Send

    style MainThread fill:#EFF6FF,stroke:#3B82F6,stroke-width:2px
    style WorkerThread fill:#ECFDF5,stroke:#059669,stroke-width:2px

为什么需要两个线程?这是一个 Python 异步编程中很经典的"同步生成器 + 异步 WebSocket"阻抗匹配问题。

QAService.stream_query()是一个同步生成器——它内部顺序执行意图识别、Milvus 检索(gRPC 阻塞调用)、本地 Rerank(CPU 密集计算)、LLM 流式调用(HTTP 阻塞读取)。如果把这段逻辑直接放在主线程的事件循环中调用next(iterator),整个事件循环会在每次推进生成器时被阻塞,导致其他 WebSocket 连接、HTTP 请求全部卡住。

解决方案:asyncio.to_thread作为桥梁。主线程通过asyncio.to_thread把同步生成器的推进操作丢给线程池中的工作线程,自己立即返回并继续处理事件循环中的其他任务。工作线程推进完成后,结果通过 Future 传回主线程,主线程再await websocket.send_json(event)发给浏览器。

图中两条线程的分工:

职责 主线程(事件循环) 工作线程
WebSocket 收发 ✅ 接收用户消息、发送事件
意图识别 ✅ 同步调用
Milvus 检索 ✅ gRPC 阻塞调用
本地 Rerank ✅ CPU 密集计算
LLM 流式调用 ✅ HTTP 阻塞读取
后台摘要刷新 ✅ 调度(不阻塞响应)

事件如何跨线程:工作线程每产出一个 token 或状态事件,生成器 yield 一次;主线程的_next_stream_event(stream)捕获这个值并通过asyncio.to_thread的返回值传回;主线程拿到事件后立即send_json给浏览器。这个过程对用户透明——浏览器看到的是连续的status → token... → end事件流。

为什么不在 LLM 流式阶段回到主线程?LLM 的llm.stream()本身返回一个迭代器,每次迭代都是阻塞的 HTTP 读取操作。如果回到主线程逐 token 读取,同样会阻塞事件循环。所以整个生成器——从意图识别到最后一个 token——全部留在工作线程中执行。

10.6 事件协议的设计原则

  1. 类型安全:每个事件都有type字段,前端用switch分派处理,不靠字段存在与否判断
  2. 诊断信息附带end事件携带完整的 retrieval 诊断信息,前端可以用 JS 渲染到页面上,帮助用户理解"系统为什么这样回答"
  3. 错误不崩溃:异常转为error事件,不抛到 WebSocket 路由。用户看到错误提示后可以继续下一轮提问
  4. 历史写入在最后end事件之后才写 MySQL 历史,确保历史记录的是完整答案(含引用增强后的来源)

第十一部分:横切能力——三级缓存

缓存不是 Stage 0-7 中额外的一步,而是嵌入 Stage 3、Stage 4 检索过程的跨阶段能力。第一次阅读建议先跳到 2.4 看完整主线,再回来看这一部分如何影响search_faq()search_doc()

当前缓存不是“答案缓存”。项目不缓存普通 LLM 自由生成答案,只缓存可以被权限、版本和配置边界约束住的安全对象。

更准确地说,当前的“三级缓存”不是三层都在缓存业务结果,而是:

L1:进程内短 TTL epoch 快照
L2:Redis 业务缓存,存 query embedding 和 FAQ/Doc 检索结果
L3:MySQL cache namespace 治理层,存 cache_epoch,用于版本失效

因此 L3 不是“把检索结果再缓存一份到 MySQL”,而是缓存控制面。它只回答一个问题:当前场景、租户、数据集应该使用第几代缓存 key。

三级缓存分别解决三个问题:

层级 实现位置 缓存内容 主要作用
L1 进程内短 TTL 缓存 qa_core/cache/stores.pyTTLMemoryCache,由qa_core/cache/manager.py使用 cache_epoch这类低频 namespace 元数据 避免每次检索都查询 MySQL,默认短 TTL,版本激活时会主动清空
L2 Redis 业务缓存 qa_core/cache/stores.pyRedisJsonCache query embedding、FAQ 检索候选、Doc 检索候选 提升热点问题检索速度,跨 API 请求复用
L3 MySQL namespace 治理 qa_core/cache/namespaces.py和表cache_namespaces scenario_id / tenant_id / dataset_id / cache_epoch 负责缓存失效边界,版本发布、回滚或人工失效时推进 epoch

这三层不是彼此替代的关系:

  • L1 只缓存少量元数据,生命周期在 API 进程内。
  • L2 才保存可复用的检索结果和 query embedding。
  • L3 不保存检索结果,只保存“当前缓存世代”,用于控制哪些 Redis key 还能被命中。
flowchart LR
    Request["用户请求"] --> Manager["CacheManager"]
    Manager --> L1["L1 TTLMemoryCache<br/>读取 cache_epoch 快照"]
    L1 -->|"未命中"| L3["L3 MySQL cache_namespaces<br/>读取或推进 cache_epoch"]
    Manager --> L2["L2 RedisJsonCache<br/>query embedding / FAQ / Doc 候选"]
    L2 -->|"命中"| Hit["直接返回缓存候选"]
    L2 -->|"未命中"| Search["Milvus Hybrid Search"]
    Search --> Write["写入 Redis<br/>等待后续同边界请求命中"]
    L3 --> Manager

当前实现只缓存这些对象:

缓存对象 代码位置 为什么可以缓存
query embedding qa_core/retrieval/models.py返回CachedEmbeddings 只和 query 文本、embedding 模型版本有关,不包含权限数据
FAQ 检索候选 qa_core/pipeline/steps.pyqa_core/pipeline/retrieval_steps.py key 绑定版本、租户、数据集、角色、source 和查询变体
Doc 检索候选 qa_core/pipeline/retrieval_steps.py key 绑定版本、权限域和检索参数,版本切换后 epoch 自动变化

这里要特别注意query embedding的调用入口:RAG Pipeline 本身没有手写embed_query()。项目在创建langchain_milvus.Milvus时把CachedEmbeddings作为embedding_function传进去;当MilvusHybridStore.search()调用similarity_search_with_score()时,langchain-milvus 会在处理dense字段时回调embedding_function.embed_query(query)。所以 query embedding 缓存的真实拦截点在CachedEmbeddings.embed_query(),不是 Pipeline 的某个显式步骤。

这也解释了两种不同的缓存命中:

  • 如果 FAQ/Doc 检索结果缓存命中,Pipeline 会直接复用候选,不进入 Milvus,也不会触发这次embed_query()
  • 如果检索结果缓存未命中,但相同 query 的 embedding 缓存命中,仍然需要访问 Milvus,不过可以少做一次 BGE-M3 在线编码。

当前实现明确不缓存这些对象:

不缓存对象 为什么不缓存
普通 LLM 最终答案 同一个问题会受知识库版本、权限域、Prompt Profile、历史追问和当前上下文影响,直接复用整段答案容易答错或泄露权限数据
文档入库 embedding 入库 chunk 数量大、重复率低,写 Redis 会挤占 query embedding 和检索缓存空间
语义近似答案缓存 “问题语义相近”不等于“企业上下文、权限和版本相同”,当前实现不用近似匹配复用答案

11.1 为什么不缓存 LLM 最终答案

同一个用户问题,不一定应该得到同一个最终答案。例如都问“入职流程是什么?”,最终答案仍可能因为以下条件不同而不同:

  • 知识库版本不同:今天 active 版本是kb_v1,明天激活kb_v2,流程资料可能已经更新,复用旧答案会答错。
  • 权限不同:普通员工只能看公开制度,HR 管理员可能能看到内部操作细则,复用管理员答案给普通员工就是权限泄露。
  • 租户或数据集不同:A 公司和 B 公司都问“入职流程”,业务流程可能完全不同,跨租户复用会污染答案。
  • Prompt Profile 不同:同样上下文,knowledge_answer可能要求总结口径,troubleshooting_steps可能要求步骤化排查,缓存答案会绕过当前模板要求。
  • 历史追问不同:上一轮问的是“实习生入职”,下一轮问“需要哪些材料?”,答案依赖历史改写;这个答案不能给另一个没有相同历史上下文的人。

这并不是说“大模型答案永远不能缓存”,而是说当前实现没有必要立刻增加这个复杂度。标准 FAQ 已经可以精确命中后直接返回,不进入 LLM;另外,查询 embedding 和 FAQ/Doc 检索结果已经能减少重复计算和 Milvus 访问。在还没有线上数据证明“最终生成是主要延迟和成本瓶颈”之前,先做稳定性和可观测性优先。

什么情况下才值得增加最终答案缓存?

先用 Trace 和运行指标回答四个问题:

要观察的信号 想回答的问题
LLM 调用比例和 P95 耗时 检索缓存已经命中时,LLM 是否仍然占大部分延迟和费用?
规范化问题重复率 相同问题、相同会话语义、相同权限域是否反复出现?
潜在缓存命中率 把版本、权限、Prompt 和会话都纳入 Key 后,真正能复用的请求还剩多少?
质量对比 复用缓存后,引用正确率、权限隔离和答案质量是否不下降?

只有当这些数据证明有明显重复热点,且 LLM 生成确实是主要成本和延迟来源时,才考虑扩展。这是“由评测和 Trace 验证后再做优化”,而不是为了缓存而缓存。

如果将来实现,边界必须更严格:

  1. 只在引用补强和生成后核验都通过后写入缓存,缓存内容不只是一段文字,还应包含answersources和最终answer_confidence
  2. 使用独立的答案缓存 Key,不能复用检索缓存 Key。Key 至少要包含scenario、租户和数据集、可见范围与角色、知识库版本与cache_epochsource_filter、Prompt Profile 和版本、LLM 模型与参数,以及必要的会话语义指纹。
  3. 命中缓存后仍要先校验当前用户权限,不能因为缓存命中就跳过权限边界。

一句话记住:缓存负责复用稳定的计算和证据;最终回答仍然根据当前请求、当前权限和当前证据生成并核验。

所以当前实现选择缓存更底层、更稳定的对象:query embedding 和 FAQ/Doc 检索候选。检索结果仍然必须绑定kb_versioncache_epoch、DataScope、source、query variants、top_k、rerank 和模型版本等边界。最终答案则每次基于当前上下文、Prompt 和权限重新生成,并在生成后做引用补强和答案置信度核验。

缓存 key 的业务边界不是“问题文本”,而是:

scenario + tenant_id + dataset_id + visibility + user_roles
+ kb_version + cache_epoch + source_filter + query_variants
+ top_k + rerank + embedding/reranker/chunk_schema 版本

所以同一个问题在不同知识库版本、不同租户、不同角色下不会复用同一份缓存。

11.2 Redis 是精确键缓存,不是模糊查询

Redis 命中依赖结构化参数生成的稳定 hash。只有 query 文本/变体以及版本、权限、source、Top-K、重排开关和模型版本等维度全部一致,才会命中同一个 key;“新人入职怎么办”和“入职流程是什么”不会因为语义相近而互相命中。

因此它的命中率来自重复热点请求、前端快捷问题、会话重试和批量业务查询,不来自 Redis 模糊搜索。即使检索结果缓存未命中,完全相同的 query 文本仍可能命中 embedding 缓存,减少一次 BGE-M3 在线编码。当前实现不做语义答案缓存,因为语义近似阈值、知识版本和权限边界组合后更容易复用错误答案;需要提升长尾命中率时,应先基于 Trace 统计真实重复率,再决定是否引入带评测门禁的语义缓存。

retrieval.cache会进入 end 事件和 Trace:

{
  "enabled": true,
  "hit_count": 1,
  "miss_count": 2,
  "events": [
    {"stage": "faq_retrieval", "hit": false, "source_type": "faq"},
    {"stage": "doc_retrieval", "hit": true, "source_type": "doc"}
  ]
}

这就是缓存闭环:命中能提速,未命中能解释,版本切换能失效,权限边界不会串数据。 状态页/admin的“企业缓存”面板和scripts/quality/cache_acceptance_smoke.py会直接读取这些字段,用来验收首次 miss、二次 hit 和版本失效是否生效。

11.3 缓存如何失效

缓存失效不依赖扫描删除 Redis 全量 key,而是通过cache_epoch生成新 key。

版本发布或回滚时,KnowledgeBaseVersionStore.activate_version()会在同一个事务中完成三件事:

  1. 更新kb_active_versions.active_kb_version
  2. 写入kb_version_activations激活或回滚流水。
  3. 调用bump_cache_epoch_for_scenario_with_conn()推进当前场景的cache_epoch

随后代码会清空 L1:

# qa_core/governance/kb_versions.py
bump_cache_epoch_for_scenario_with_conn(conn, self.scenario.scenario_id)
get_cache_manager().l1_cache.clear()

因此新版本激活后:

  • FAQ/Doc 检索缓存会因为kb_version + cache_epoch改变而重新 miss。
  • Redis 中旧 key 不会立即物理删除,会按 TTL 自然过期。
  • query embedding 缓存可以继续复用,因为它只依赖 query 文本和 embedding 模型版本,不依赖知识库版本。
  • qa_core/retrieval/factory.py中的 Milvus store wrapper 是资源对象缓存,不是业务结果缓存;active 版本变化时会清空这个进程级 wrapper 缓存,避免长生命周期对象持有旧状态。

手工失效也走同一套 epoch 机制:

POST /api/admin/cache/invalidate

11.4 缓存行为检查

本地单测验证 key 边界、epoch 失效和 query embedding 缓存:

python -m pytest tests/test_enterprise_cache.py -q

Docker 环境验证真实 Redis 命中路径:

docker compose --env-file .env.compose exec api python scripts/quality/cache_acceptance_smoke.py --base-url http://127.0.0.1:8000

验收预期是:同一个场景、同一个数据域、同一个问题连续问两次,第一次出现 cache miss,第二次出现 cache hit;如果中间激活了新知识库版本,下一次查询应重新 miss。


讨论

留下你的想法

评论功能尚未配置。启用 Giscus 后,这里会显示基于 GitHub Discussions 的评论区。

图片预览

100%