RAG,中文名为检索增强生成,是当前解决LLM回答出现幻觉,时效性不足的问题和构建企业知识库问答的有效手段之一,它的核心链路有两条:1.文档入库(写)。2.RAG问答(读)。
本文我将以我自己的RAG项目为例分析这两条链路的关键设计。
请求包装
Filter层
针对一个http请求,为了在知识库文档上传请求触发 multipart 解析之前实施并发控制,我在Servlet Filter层定义了一个名为UploadRateLimitFilter的过滤器。该过滤器根据请求方法和 URI 判断请求是否命中知识库文档上传接口,并使用@Order(Ordered.HIGHEST_PRECEDENCE)将其设置为 Spring 管理的过滤器中的最高优先级之一,使限流逻辑尽可能早地执行。当并发上传数超过限制时,可以在请求进入DispatcherServlet并触发 multipart 解析之前直接返回 429,降低大量上传请求生成临时文件、占满服务器磁盘的风险。
UploadRateLimitFilter没有直接实现Filter接口,而是继承了OncePerRequestFilter,并重写了doFilterInternal()。不同与一般的Filter会被多次调用(Controller返回视图、错误页面跳转、异步请求等),OncePerRequestFilter提供了重复执行保护:它会在请求对象中记录“已经过滤”的标记,避免同一请求在过滤器链内部发生重复 dispatch 时再次执行核心过滤逻辑;默认情况下它还会跳过后续的异步和错误 dispatch。
为什么需要做"仅一次"保证? UploadRateLimitFilter继承了OncePerRequestFilter。在 Servlet 模型中,一次 HTTP 请求可能因为转发、异步处理或错误处理产生多次 dispatch,从而多次经过过滤器链。OncePerRequestFilter会通过请求属性标记当前过滤器是否已经执行,避免过滤器链内部的重复进入再次触发doFilterInternal();同时,它默认不参与后续的 ASYNC 和 ERROR dispatch。这对于上传并发控制很重要,因为doFilterInternal()中包含 Redisson permit 的申请和释放逻辑。如果同一个请求重复执行该方法,就可能重复申请 permit,造成并发额度被额外占用,甚至在许可数量较少时出现后一次申请等待前一次释放的情况。因此,这里的“仅一次”主要用于保证一份上传请求只执行一套 permit 申请与释放逻辑。
MVC层
进入DispatcherServlet之后,我通过一个实现了WebMvcConfigurer接口并重写了addInterceptors()的配置类定义了两个拦截器:1.SaInterceptor。2.UserContextInterceptor。二者相同的点为它们都跳过了异步调度请求和OPTIONS预检,不同的点为SaInterceptor用来进行登录校验,UserContextInterceptor用来记下用户身份信息。
AOP 层
在切面链中,共使用了三个自定义注解:
- 1.IdempotentSubmit。
- 2.ChatRateLimit。
- 3.RagTraceNode。
分别对应三个Around方法:
- 1.IdempotentSubmitAspect.idempotentSubmit()。
- 2.ChatRateLimitAspect.limitStreamChat()。
- 3.RagTraceAspect.aroundNode()。
三个切面并不是同时作用于同一个方法然后按照固定优先级依次执行,而是分布在 Controller、Service 和 RAG 核心节点三个不同层次,共同构成流式对话的横切处理链:
IdempotentSubmit → ChatRateLimit → RagTraceNode
IdempotentSubmit用于防止接口被并发重复提交。它被标注在 Controller 方法上,当请求进入目标接口时,切面根据注解中的 SpEL 表达式生成分布式锁 Key,并通过 Redisson 的 tryLock() 尝试获取锁。获取失败时直接抛出异常;获取成功后执行目标方法,并在 finally 中释放锁。
ChatRateLimit被标注在Service实现类的方法上,对流式对话进行全局并发控制。它的切面不会立即执行目标方法,而是先将请求加入 Redis 等待队列。请求获得 Redisson 信号量 permit 后,才会被提交到线程池中执行;如果等待超时,则通过 SSE 返回系统繁忙信息。对话完成、超时或异常时,会自动释放 permit。
RagTraceNode用于标记 RAG 调用链中需要单独采集的业务节点,例如问题改写、意图识别、多路检索、模型路由和具体模型调用。切面在方法执行前创建节点记录,并保存节点名称、类型、父节点和开始时间;方法执行成功后记录执行耗时并将状态更新为 SUCCESS,执行异常时将状态更新为 ERROR。最后弹出当前节点,恢复上一层调用上下文。
文档入库
对于 RAG 系统而言,文件上传成功并不等于文档已经完成入库。上传阶段解决的是原始文件和文档元数据的持久化问题,而真正决定后续检索效果的,是文档解析、分块、增强、向量化和索引写入过程。
在本项目中,文档入库被拆分为两个阶段:
上传文档
→ 保存原始文件与文档元数据
→ 创建文档记录
启动分块
→ 投递异步任务
→ 解析文档
→ 文档分块
→ Embedding
→ 写入关系库与向量库
这样设计的原因是,上传文件通常只涉及网络传输和文件存储,而解析、分块和 Embedding 都属于耗时操作。如果将这些操作全部放在上传请求中同步完成,请求很容易因为大文件、模型响应变慢或者向量库抖动而超时。
文档上传
请求进入 Controller 后,首先查询目标知识库,确认知识库存在,随后根据文档来源处理原始文件,并将文件保存到对象存储中。
文件保存完成后,向数据库表插入一条文档记录,其中包括:
- 所属知识库 ID
- 文档名称
- 文件存储地址
- 文件类型和文件大小
- 文档来源类型
- 分块策略或 Pipeline ID
- 当前处理状态
上传阶段的核心目标不是直接生成向量,而是先建立下面的关联关系:
知识库
→ 文档元数据
→ 原始文件
原始文件和文档记录被保留下来后,即使后续分块策略发生变化,也可以重新读取原始文件并构建索引,而不需要用户再次上传。
异步启动分块
用户启动文档分块时,不会直接在请求线程中完成文档解析,而是构造一个文档分块任务事件,并通过 RocketMQ 事务消息投递分块任务。在事务消息的本地事务中,系统通过条件更新将文档状态修改为RUNNING,条件更新同时承担并发控制作用。如果同一文档已经处于RUNNING状态,本次更新不会成功,系统会拒绝重复启动,避免两个消费者同时对同一文档执行分块和索引重建。 使用事务消息主要是为了协调两个状态:
数据库中:文档已经进入 RUNNING
消息队列中:分块任务可以被消费
与在数据库事务提交后直接发送普通消息相比,事务消息降低了以下两种不一致情况发生的概率:
数据库已经更新为 RUNNING,但消息发送失败
消息已经被消费,但数据库状态没有更新
如果 RocketMQ 无法确定本地事务是否成功,会触发事务回查:通过查询数据库中的文档状态判断本地事务是否已经提交,只有文档存在且状态为RUNNING,才认为分块任务应当被提交。
MQ消费与用户上下文恢复
事务消息提交后,由文档分块任务 MQ 消费者消费分块事件。
HTTP 请求进入系统时,用户信息由UserContextInterceptor写入当前线程的UserContext。但是 RocketMQ 消费线程和原始 HTTP 请求线程没有直接关系,因此消费者无法自动获得原请求线程中的用户信息。为了解决这个问题,生产者会将操作人写入消息事件。消费者收到消息后,根据事件中的操作人重新构造LoginUser并写入UserContext:
HTTP 请求线程
→ 将 operator 写入 MQ 消息
MQ 消费线程
→ 从消息读取 operator
→ 重建 UserContext
→ 执行文档分块
→ finally 清理 UserContext
这样,后续文档状态更新和操作日志仍然可以记录实际操作人,同时通过finally清理线程上下文,避免 MQ 消费线程复用时产生用户数据污染。
文档处理模式
当前项目支持两种处理模式:
- CHUNK:固定分块模式
- PIPELINE:可配置流水线模式
固定分块模式
固定分块模式的链路相对直接:
读取原始文件
→ Tika 提取文本
→ 根据分块策略切分文本
→ 调用 Embedding 模型生成向量
→ 写入 Chunk 表和向量库
系统目前默认使用 Apache Tika 提取 PDF、Word、TXT 等文件中的文本。 文本提取完成后,系统根据文档配置选择具体的分块策略。不同策略负责解决不同类型文档的边界问题,例如固定长度分块强调实现简单,递归分块则尽可能在段落、句子等自然边界处切分文本。
分块结果被封装为VectorChunk,其中主要包含:
chunkId
chunkIndex
content
embedding
metadata
随后,ChunkEmbeddingService批量调用 Embedding 模型,为每个 Chunk 生成向量。
Pipeline 模式
固定分块模式适合处理规则相对统一的文档,但企业知识库中的文档来源、结构和质量可能存在较大差异。因此,项目还提供了可配置的 Pipeline。
各节点职责如下:
- Fetcher:从本地文件、HTTP 地址、对象存储等数据源获取原始内容。
- Parser:识别文件类型,并将二进制内容转换成结构化文本。
- Enhancer:在分块前对整篇文档执行摘要、关键词提取或元数据补全。
- Chunker:按照配置的策略切分文档。
- Enricher:对单个 Chunk 生成问题、摘要或补充上下文。
- Indexer:生成向量并写入指定的向量空间。
文档Pipeline执行引擎不会将这些节点硬编码成固定流程,而是读取Pipeline定义实体类构建节点之间的连接关系。执行前,系统会检查节点引用是否合法以及 Pipeline 是否存在环;执行过程中,还可以根据节点条件决定是否跳过某个处理步骤。每个节点执行后都会记录状态、耗时、输出摘要和异常信息。因此,相比固定分块模式,Pipeline 模式不仅支持更复杂的文档处理,也提供了更细粒度的Pipeline可观测性。
Chunk与向量的统一替换
文档处理完成后,系统需要同时更新两类数据:
关系数据库中的 Chunk 原文
向量数据库中的 Chunk Embedding
重新处理文档时,系统采用“删除旧数据,再写入新数据”的替换语义:
→ 删除旧的关系型 Chunk
→ 批量写入新的 Chunk
→ 删除旧的文档向量
→ 写入新的文档向量
→ 更新文档状态和 Chunk 数量
全部步骤成功后,文档状态被更新为SUCCESS;任一处理阶段出现异常,则记录分块日志并将文档状态更新为FAILED。
至此,一份原始文档完成了从文件到可检索知识单元的转换:
→ 原始文件
→ 文本
→ Chunk
→ Embedding
→ 向量索引
RAG问答
文档入库解决了“知识如何进入系统”的问题,而 RAG 问答链路需要解决的是另一个问题:当用户提出问题时,系统如何从大量知识中找到真正相关的内容,并将这些内容组织成大模型能够理解的上下文。
最简单的 RAG 实现通常只有三步:
用户问题
→ 向量检索
→ 拼接 Prompt
→ 调用大模型
这种方式能够完成基本的知识库问答,但在多轮对话、复杂问题、多知识库、实时数据和模型故障等场景下,很容易出现检索范围错误、问题语义不完整、召回内容不足或者模型调用失败等问题。
因此我的项目将在线问答拆分成一条完整的处理流水线:
会话记忆
→ 问题改写与拆分
→ 意图识别
→ 歧义引导或系统短路
→ 多通道检索与 MCP 调用
→ 结果去重与重排
→ Prompt 构建
→ 模型路由
→ SSE 流式输出
当流式请求通过前面的幂等校验和全局并发控制,并获得执行许可后,系统创建本次请求对应的conversationId和taskId,再构造StreamChatContext交给 Pipeline 执行。
StreamChatContext贯穿整个问答过程,用于保存:
- 用户原始问题
- 会话 ID 和任务 ID
- 当前用户 ID
- 是否启用深度思考
- 会话历史
- 问题改写结果
- 子问题及其意图
- SSE 回调对象
与在多个方法之间不断增加参数相比,使用上下文对象可以让 Pipeline 的各个阶段围绕同一份请求状态工作。
会话记忆
在多轮对话中,用户当前的问题往往依赖前文。 例如:
用户:公司的请假流程是什么?
助手:……
用户:那超过三天呢?
第二个问题中的“那”并没有明确指出讨论对象。如果直接使用“那超过三天呢”执行向量检索,Embedding 模型很难恢复完整语义。 因此,Pipeline 的第一步是加载历史会话,会话记忆由两部分组成:
历史摘要
+
最近若干轮原始消息
系统会并行加载会话摘要和历史消息,然后将摘要放在历史消息之前。摘要用于保留更早的对话信息,最近的原始消息则用于保存更精确的上下文。当历史轮数超过配置阈值后,系统会异步触发会话压缩,将较早的消息压缩成摘要,从而避免每次请求都把全部历史发送给大模型。
这是一种典型的长短期记忆组合处理,在上下文完整性和 Token 消耗之间做了平衡:
长期记忆:历史摘要,信息密度高
短期记忆:最近几轮原始消息,细节更准确
问题归一化、改写与拆分
加载会话历史后,Pipeline开始处理用户问题。这一阶段包含三个动作:
术语归一化
→ 结合历史进行问题改写
→ 将复杂问题拆分为多个子问题
术语归一化
不同用户可能使用不同名称描述同一个概念,例如简称、别名或者内部业务术语。系统首先将这些词语转换成统一表达,减少因为词汇差异造成的检索偏差。
问题改写
系统会把最近的用户和助手消息以及当前问题一起发送给 LLM,让模型将带有指代、省略或上下文依赖的问题改写成一条独立、完整的问题。
例如:
原始问题:那超过三天呢?
结合历史改写后:
员工请假超过三天时需要经过哪些审批流程?
改写后的问题不再依赖对话历史,可以直接用于后续意图识别和知识检索。
为了降低改写结果的随机性,请求使用较低的temperature和topP,同时关闭深度思考。因为这一阶段的目标不是创作内容,而是稳定地恢复用户问题的完整语义。
子问题拆分
用户的一句话可能同时包含多个独立问题:
请介绍 OA 系统的请假流程,并比较事假和年假的区别。
如果将整句话只生成一个向量,检索结果可能只覆盖其中一个主题。因此,改写服务要求 LLM 返回结构化结果:
{
"rewrite": "介绍 OA 系统的请假流程,并比较事假和年假的区别",
"sub_questions": [
"OA 系统中的请假流程是什么?",
"事假和年假有什么区别?"
]
}
后续系统会以子问题为基本单位并行执行意图识别和上下文构建。
如果 LLM 调用失败或者返回内容无法解析,系统不会直接终止问答,而是将归一化后的原问题作为唯一子问题继续执行:
LLM 改写成功
→ 使用改写结果和子问题
LLM 改写失败
→ 使用归一化问题作为单一子问题
这种降级策略保证了问题改写属于质量增强能力,而不是整个 RAG 链路的单点故障。
意图识别
问题改写完成后,IntentResolver会为每个子问题识别业务意图。意图节点以树形结构维护,每个叶子节点可以描述:
- 意图 ID
- 所属业务路径
- 意图描述
- 示例问题
- 意图类型
- 对应知识库
- MCP 工具 ID
- 自定义 Prompt 模板
- 检索 TopK
当前在线主链实际从 Redis 加载最新的意图树;缓存不存在时,再从数据库读取并回填缓存。随后,系统将所有叶子意图的路径、描述、类型和示例问题组织成 Prompt,交给 LLM 对当前子问题进行打分。
模型需要返回类似下面的结构:
[
{
"id": "oa_leave_process",
"score": 0.92
},
{
"id": "hr_leave_policy",
"score": 0.71
}
]
系统会过滤掉低于最低置信度的意图,并限制单个问题以及整次请求的意图总数,避免后续因为意图过多而访问大量知识库或工具。
当一个问题被拆成多个子问题时,各个子问题的意图识别会在线程池中并行执行:
子问题 1 → 意图识别 ┐
子问题 2 → 意图识别 ├→ 汇总意图
子问题 3 → 意图识别 ┘
如果总意图数超过上限,系统会优先为每个子问题保留一个最高分意图,再按照全局分数分配剩余名额。这样可以防止某个子问题占满意图配额,导致其他子问题完全没有检索方向。
歧义引导
意图识别并不一定总能得到唯一答案。例如用户只问:
如何申请权限?
系统可能同时识别出 OA、财务系统和数据平台的权限申请意图,并且几个意图的得分非常接近。此时如果直接选择最高分意图,可能在错误的知识库中检索。系统此时会检查:
- 当前是否只有一个子问题。
- 是否存在多个相近的高分知识库意图。
- 这些意图是否属于不同业务系统。
- 用户问题中是否已经明确写出系统名称。
满足歧义条件时,系统不会继续检索,而是直接通过 SSE 返回一个引导问题:
请确认你要咨询哪个系统的权限申请:
1. OA 系统
2. 财务系统
3. 数据平台
这相当于在高成本检索和生成之前增加一次澄清,让用户补充缺失条件。
歧义引导的核心思想是:
不确定时先询问
优于
带着错误假设继续检索
SYSTEM 意图短路
有些问题不依赖企业知识库,例如欢迎语、能力说明或者系统使用帮助。这类意图会被标记为SYSTEM。如果所有子问题都只命中 SYSTEM 意图,Pipeline 会跳过知识库检索和 MCP 调用,直接使用意图节点上的自定义 Prompt 或默认系统 Prompt 调用大模型:
SYSTEM 意图
→ 跳过检索
→ 构造系统 Prompt
→ 流式调用 LLM
这种短路机制可以减少无意义的向量检索,并降低问答延迟。
多通道知识检索
项目目前主要包含两个知识检索通道:1.意图定向检索。2.全局向量检索。
意图定向检索
当系统识别出置信度足够高的知识库查询意图时,系统只需要在相关知识库中执行检索:
用户问题
→ 命中 OA 请假意图
→ 只检索 OA 制度知识库
相比在所有知识库中搜索,这种方式具有三个优势:
- 缩小检索范围。
- 减少不相关 Chunk。
- 降低不同业务系统中相似术语造成的干扰。
如果一个子问题命中多个知识库查询意图,系统会并行检索这些意图所关联的知识库。
全局向量检索
意图识别也可能失败,或者最高意图分数不足以确定检索范围。此时系统在所有可用知识库的 Collection 中执行向量检索。
全局检索承担的是兜底职责:
无意图
或
意图最高分低于阈值
↓
启用全局向量检索
当意图分数处于某个中间区间时,意图定向检索和全局向量检索可能同时启用,并在线程池中并行执行:
意图定向检索 ┐
├→ 合并检索结果
全局向量检索 ┘
意图定向检索强调精确性,全局检索强调召回兜底,二者共同降低“意图识别错误导致完全检索不到内容”的风险。
多子问题之间同样会并行构建检索上下文,因此整个检索过程形成了多层并发:
子问题之间并行
↓
检索通道之间并行
↓
多个知识库之间并行
这种设计可以降低复杂问题在串行检索下不断累积的延迟,但也意味着线程池大小、超时控制和上下文传递必须受到统一管理。
MCP 实时数据调用
知识库适合保存制度、文档和相对稳定的业务知识,但无法及时回答实时数据问题,例如:
今天还有多少会议室可用?
当前订单状态是什么?
本月销售额是多少?
这类问题会被识别为 MCP 意图。系统根据意图节点上的mcpToolId找到对应的MCPToolExecutor,并从用户问题中提取工具调用参数,构造MCPRequest。如果一个子问题命中多个 MCP 工具,这些工具会并行执行。工具返回结果被格式化成动态数据上下文,供后续 Prompt 使用。
因此,项目中的上下文来源可以分成两类:
知识库 Context
→ 来自文档和向量库
→ 相对稳定的企业知识
MCP Context
→ 来自外部系统或实时接口
→ 动态业务数据
检索结果去重与重排
多通道检索提高了召回率,但也会带来重复结果。同一个 Chunk 可能同时出现在:
- 意图定向检索结果中。
- 全局向量检索结果中。
- 多个知识库或多个子问题的结果中。
系统去重时优先使用 Chunk ID 作为唯一标识;没有 ID 时,使用内容哈希作为兜底,保留分数最高的结果:
多通道结果
→ 按 Chunk ID 合并
→ 重复 Chunk 保留最高分
去重完成后,系统执行重排。向量相似度衡量的是语义距离,但相似度高并不一定意味着它最适合回答当前问题。Rerank 模型会同时读取原始问题和候选 Chunk,对它们进行更精细的相关性判断,并输出最终 TopK:
向量召回:从大量文档中快速找候选
Rerank:从候选中选择真正适合回答问题的内容
因此,完整检索过程不是简单的 TopK 向量查询,而是:
扩大候选召回
→ 多通道合并
→ 去重
→ Rerank
→ 截取最终 TopK
如果所有检索通道都没有返回有效内容,Pipeline 不会在没有知识依据的情况下继续调用模型,而是通过 SSE 返回“未检索到相关文档内容”,然后结束本次请求。这可以减少模型脱离知识库自由回答产生幻觉的概率。
Prompt 构建
检索完成后,系统需要把不同来源的信息组织成模型消息。
不同场景使用不同的系统 Prompt:
只有知识库内容
→ 企业知识库问答模板
只有 MCP 数据
→ 动态数据回答模板
同时存在知识库内容与 MCP数据
→ 混合上下文模板
最终消息按照下面的结构构建:
System Prompt
→ MCP 动态数据
→ 知识库文档证据
→ 会话摘要和历史消息
→ 当前用户问题
如果当前问题被拆分成多个子问题,系统会显式编号:
请基于上述文档内容回答以下问题:
1. OA 系统中的请假流程是什么?
2. 事假和年假有什么区别?
这样可以提醒模型逐项作答,降低复杂问题中漏答某个子问题的概率。
Prompt 构建并不是简单地把所有字符串拼接起来,而是需要明确区分:
- 模型应遵守的系统规则。
- 来自实时工具的数据。
- 来自知识库的文档证据。
- 用户历史对话。
- 当前真正需要回答的问题。
模型路由与故障转移
构造ChatRequest后,Pipeline不会直接绑定某个固定模型客户端,而是根据请求和模型配置选择合适的模型。路由器按照默认模型、优先级和健康状态选择候选模型。普通对话和深度思考请求会生成不同的候选模型列表。深度思考模式下,只保留声明支持 Thinking 的模型。如果某个模型启动请求失败、首包超时或没有返回有效内容,系统会取消当前调用并切换到下一个候选模型。
对于流式请求,故障转移不能只判断 HTTP 连接是否建立,因为连接成功不代表模型一定能够返回内容。因此,路由器增加了首包探测:
启动流式模型请求
→ 等待第一个有效数据包
→ 首包成功:确认模型可用
→ 启动失败:切换下一个模型
→ 首包超时:取消并切换下一个模型
→ 首包报错:取消并切换下一个模型
一旦某个模型已经成功返回首包,系统就将该模型标记为调用成功,并把后续数据交给 SSE 回调处理。已经开始向用户输出内容后,不再随意切换模型,否则可能产生内容重复或回答语义不连续的问题。系统还会记录模型调用的成功和失败状态。当某个模型连续失败达到阈值后,会暂时跳过该模型;等待一段时间后再允许试探调用,从而避免故障模型持续拖慢整个问答链路。
SSE 流式输出与任务取消
一次完整对话可能产生以下 SSE 事件:
meta
→ 返回 conversationId 和 taskId
message
→ 返回思考内容或回答内容
finish
→ 返回最终消息 ID 和会话标题
done
→ 表示本次 SSE 完成
模型输出的内容会被区分为两类:
think
response
回调对象一边将模型输出发送给客户端,一边累积完整的思考过程和回答正文。模型正常完成后,系统将助手消息持久化到会话记忆,发送finish和done事件,并注销当前流式任务。用户调用停止接口时,系统会:
在 Redis 写入取消标记
→ 通过 Redis Topic 广播 taskId
→ 找到任务所在的应用实例
→ 调用模型取消句柄
→ 保存已经生成的部分内容
→ 发送 cancel 和 done 事件
使用 Redis 广播的原因是,在多实例部署下,停止请求和实际执行流式任务的请求不一定落在同一个应用节点。单纯使用本地 Map 无法完成跨节点任务取消。至此,一次在线 RAG 问答完成了从用户问题到最终回答的完整闭环:
用户问题
→ 恢复对话语义
→ 识别问题意图
→ 确定检索范围
→ 获取文档和实时数据
→ 筛选有效证据
→ 构建结构化 Prompt
→ 选择可用模型
→ 流式返回并持久化回答
这条链路的重点并不是简单地“调用一次向量库和大模型”,而是在召回质量、回答延迟、系统成本和故障容错之间进行权衡。
