RAG Pipeline 主流程与 Prompt 生成
贯通 RAG Pipeline 的八个执行阶段,明确 FAQ 快速路径、上下文筛选与截断、Prompt Profile 选择、变量注入及答案引用增强。
第一部分:技术背景 — 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.py、context.py |
第五部分 5.1~5.4、第六部分 6.1~6.3 |
| Stage 6 | stream_llm_answer()、引用增强、生成后核验 |
steps.py、citations.py、confidence.py |
第七部分 7.1~7.5 |
| Stage 7 | history.add_turn()、finish_success()、finish_error() |
rag.py、runtime.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=true、run_doc=false、rerank=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 不同或最终答案不同的影响。
为保证正确性,只有以下条件同时满足才允许复用:
rewritten_query仍等于原问题,说明没有追问改写。query_variants的第一个元素仍是原问题。- 快速探测与完整计划的有效
source_filter相同。 - 快速探测的候选容量不少于完整计划的
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()
这里有两个容易混淆的时间点:
build_answer_prompt_profile()在prepare_retrieval()中完成。它根据意图和问题类别先确定回答策略。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 选择模板。模板选择是确定性的规则判断,原因有三点:
- 检索策略和回答口径必须一致,不能检索阶段认为是 FAQ,生成阶段又随机改成综合分析模板。
- 少一次 LLM 判断,减少延迟、成本和不稳定性。
- 费用、合规等高风险问题必须稳定使用保守模板,不能把模板选择交给模型自由判断。
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()将这两个字符串转换成SystemMessage和HumanMessage,开始流式生成。第 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_answer在steps.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_confidence由combine_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_match、history_rewrite_used、insufficient_context。
注意:当前answer_confidence是可解释的工程信号,不是经过概率校准的“答案正确率”。它适合用于前端提示、Trace 排查、Bad Case 分层和评测报告观察;默认不替代 Recall@K、MRR、关键词覆盖率和人工复核,也不把历史聊天当成真值来源。历史只作为“是否依赖追问改写”的稳定性信号参与扣分。
后续如果要把它升级成更严格的质量门禁,需要先积累评测集、人工标注和线上反馈,再做分数校准。否则直接用一个固定置信度阈值阻断发布,容易把“有证据但表达保守”和“无证据乱答”混在一起。
第八部分:Stage 7 统一收口
Stage 7 不负责再检索或生成内容,而是保证所有分支都以一致的方式结束:保存最终答案、补齐诊断信息、写入 Trace,并向前端发出end或error事件。它是整条在线链路的出口。
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 事件协议的设计原则
- 类型安全:每个事件都有
type字段,前端用switch分派处理,不靠字段存在与否判断 - 诊断信息附带:
end事件携带完整的 retrieval 诊断信息,前端可以用 JS 渲染到页面上,帮助用户理解"系统为什么这样回答" - 错误不崩溃:异常转为
error事件,不抛到 WebSocket 路由。用户看到错误提示后可以继续下一轮提问 - 历史写入在最后:
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.py的TTLMemoryCache,由qa_core/cache/manager.py使用 |
cache_epoch这类低频 namespace 元数据 |
避免每次检索都查询 MySQL,默认短 TTL,版本激活时会主动清空 |
| L2 Redis 业务缓存 | qa_core/cache/stores.py的RedisJsonCache |
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.py、qa_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 验证后再做优化”,而不是为了缓存而缓存。
如果将来实现,边界必须更严格:
- 只在引用补强和生成后核验都通过后写入缓存,缓存内容不只是一段文字,还应包含
answer、sources和最终answer_confidence。 - 使用独立的答案缓存 Key,不能复用检索缓存 Key。Key 至少要包含
scenario、租户和数据集、可见范围与角色、知识库版本与cache_epoch、source_filter、Prompt Profile 和版本、LLM 模型与参数,以及必要的会话语义指纹。 - 命中缓存后仍要先校验当前用户权限,不能因为缓存命中就跳过权限边界。
一句话记住:缓存负责复用稳定的计算和证据;最终回答仍然根据当前请求、当前权限和当前证据生成并核验。
所以当前实现选择缓存更底层、更稳定的对象:query embedding 和 FAQ/Doc 检索候选。检索结果仍然必须绑定kb_version、cache_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()会在同一个事务中完成三件事:
- 更新
kb_active_versions.active_kb_version。 - 写入
kb_version_activations激活或回滚流水。 - 调用
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 的评论区。