Ragent AI 企业级 Agentic RAG平台源码深度学习教程
目标读者:有 Java/Spring 基础、想快速读懂 Ragent 源码的后端开发者。
阅读时长:精读约 3~5 小时,配合断点调试效果最佳。项目一句话定位:Ragent 是一个面向 Agentic RAG 演进的生产级 Java AI 应用平台(Spring Boot 3.5.7 + Java 17 + React 18 + Vite 5)。它不是一个调通 OpenAI 的 Demo,而是把企业落地 RAG 时真正会踩的坑——多格式文档解析、多通道检索融合、模型路由容错、会话记忆控成本、分布式限流、全链路追踪、来源溯源——都用工程化方式解决了。
第 0 章:先建立全局认知(10 分钟)
🚀 什么是 Ragent AI?
Ragent 是一个面向 Agentic RAG 演进的生产级 Java AI 应用平台,覆盖从文档入库到智能问答的完整链路。
- 混合检索:向量、关键词、知识图谱、联网搜索并行召回,支持去重、RRF 融合与 Rerank。
- 问题理解:支持查询词映射、问题重写与拆分、树形意图识别和多知识库路由。
- 模型与工具:支持模型档位、首包探测、熔断降级,以及 MCP 工具发现、提参与校验。
- 会话记忆:最近 N 轮消息结合持久化摘要,控制 Token 成本并保留关键上下文。
- 流量保护:Redis 公平排队与分布式并发控制,避免突发请求压垮模型服务。
- 知识闭环:提供可编排入库 Pipeline、远程刷新、回答溯源、用户反馈、Trace 和管理后台。

生产落地智能体会踩的坑,这里都有对应方案,一套经过真实场景锤炼的工程实践,系统补全 RAG / Agent / MCP 等知识,面试写进简历聊得起来。
把 Ragent 想象成一套**“企业级 RAG 问答中台”**。它的工程价值不在于"能调 LLM",而在于把生产环境必然出现的脏活和坑都填平了:
| 生产痛点 | Ragent 的解法 | 代码落点(后文详述) |
|---|---|---|
| 文档格式混乱(PDF/Word/网页/图片) | framework/core 统一解析 + 切块;入库流水线可编排 |
core/、ingestion/engine |
| 单路向量召回率低 | 4 路并行检索 + RRF 融合 + Rerank | rag/core/retrieval |
| 口语化、多轮依赖 | 查询改写/拆分 + 会话记忆 + 意图树 | rewrite/、memory/、intent/ |
| 模型供应商会挂 | 模型档位路由 + 首包探测 + 三态熔断 | infra-ai/.../RoutingLLMService |
| 高并发打爆模型 | Redis ZSET 公平排队 + Lua 抢占 + 过期信号量 | rag/service/ratelimit |
| 答案不可信 | 来源引用、原文预览、用户点赞/点踩 | rag/core/source、messages 表 |
| 线上难排查 | 全链路 Trace(AOP + TTL 透传)+ 审计日志 | framework/trace、audit/ |
顶层目录职责速记:
| 目录 | 文件数 | 一句话职责 |
|---|---|---|
bootstrap |
496 | 业务核心:DDD 分层,承载所有 Controller/Service/DAO 与 RAG 编排 |
framework |
39 | 公共契约 + 基础设施(追踪、幂等、分布式 ID、MQ、异常、缓存) |
infra-ai |
50 | AI 接入层(LLM/Embedding/Rerank/VLM),统一 OpenAI 协议 |
mcp-server |
8 | 独立 MCP 工具服务(默认 9099 端口) |
frontend |
141 | React 问答 UI + 管理后台 |
resources |
32 | SQL 建表、yaml 配置、提示词模板、Lua 脚本 |
docs |
6 | 架构图(PNG/SVG)与少量说明 |
assets |
26 | README 引用的架构图 |
记忆口诀:framework 管地基,infra-ai 管模型,bootstrap 管业务,mcp-server 管工具。 四方互不越界、单向依赖。
0.1 数据模型鸟瞰(读 resources/database/schema_pg.sql)
RAG 系统的持久化核心是"文档→切片→向量+来源→消息"这条主线,建表脚本在 resources/database/schema_pg.sql(PostgreSQL + pgvector;另有 schema_table.sql 为 MySQL 版)。关键表:
| 表 | 职责 | 关键字段 |
|---|---|---|
t_user |
系统用户 | username/password/role(admin/user) |
t_conversation |
会话列表 | conversation_id/user_id/title/last_time(唯一约束 conversation_id+user_id) |
t_conversation_summary |
会话摘要(与消息分离存储) | conversation_id/last_message_id/content(配合第 7 章记忆压缩) |
t_message |
消息表 | conversation_id/role/content/thinking_content/sources/message_status(思考链与来源引用都落这里) |
t_message_source |
消息来源引用(升级脚本 260722_01 新增) |
把答案引用的文档来源持久化,供前端来源面板 |
t_knowledge_document |
知识库文档 | doc_id/name/status/collection |
t_knowledge_chunk |
文档切片 | doc_id/chunk_index/text/embedding(向量存 pgvector 列) |
t_knowledge_vector_collection |
向量集合(升级脚本 260703 新增) |
支持多集合/多知识库隔离检索 |
t_biz_change_log |
业务变更审计(升级脚本 260709) |
配合 @EnableLogInstance 审计 |
理解这条主线后,再看
bootstrap的 7 个业务域就清晰了:knowledge管文档/切片,ingestion把文档变成 chunk,rag用 chunk 检索并产出 message,user管账号,audit管变更日志。
第 1 章:技术栈与构建系统
1.1 构建与版本(读 pom.xml 必看)
- 根
pom.xml:packaging=pom的 Maven 聚合工程,直属子模块bootstrap / framework / infra-ai / mcp-server。 - 父 BOM:
spring-boot-dependencies3.5.7;Java 17(maven.compiler.release=17)。 - 关键插件(改代码前必须知道,否则本地编译不过):
spotless-maven-plugin:绑定到compile阶段,强制每个 Java 文件带 Apache-2.0 License 头(文件顶部 16 行注释块)。新增文件务必先加头,否则mvn compile失败。maven-surefire-plugin:通过<argLine>注入 Mockito 5 inline mock agent,保证单测中 mock final 类/方法可用。lombok.config:根目录存在,统一 Lombok 配置(如@RequiredArgsConstructor生成构造器注入)。
1.2 依赖速查表

| 领域 | 技术 | 用途 |
|---|---|---|
| Web | Spring Boot Web (WebFlux/Servlet) | REST + SSE |
| ORM | MyBatis-Plus 3.5.14 (spring-boot3) | 持久层;<parameters>true</parameters> 保留参数名 |
| 向量库 | Milvus 2.6.6 或 pgvector(二选一) | vector/ 两套实现 |
| 文档解析 | Apache Tika 3.2.3 | core/ 解析 PDF/Word/HTML 等 |
| 认证 | Sa-Token 1.43.0(+ redis 模板) | user/ 用户与权限 |
| 分布式 | Redisson 4.0.0(限流/锁)、RocketMQ 2.3.5(事务消息) | 流量保护、数据一致性 |
| 对象存储 | AWS S3 SDK / 阿里云 OSS SDK | storage/ 抽象 |
| MCP | modelcontextprotocol SDK 1.1.2 | mcp-server + rag/core/mcp |
| AI 调用 | OkHttp 4.12(HTTP 客户端) | 全部走 OpenAI 兼容协议 |
| 工具库 | Hutool(IdUtil.getSnowflakeNextIdStr() 等)、Caffeine(本地缓存) |
贯穿各模块 |
1.3 启动类与全局配置
- 后端入口:
bootstrap/.../RagentApplication.java- 注解:
@SpringBootApplication+@EnableScheduling(文档定时刷新/过期清理)+@EnableLogInstance(bizlog 审计 SDK)+@MapperScan(精确扫描user、knowledge、ingestion、rag、audit5 个 DAO 包)。 - 端口:
9090;上下文路径:/api/ragent(见application.yml)。
- 注解:
- 前端:
frontend/(Vite dev server,默认 5173,通过VITE_API_BASE_URL指向后端)。

第 2 章:四模块架构与依赖方向
项目是一个 Maven 多模块工程(父 pom.xml 的 <modules> 声明了 bootstrap、framework、infra-ai、mcp-server 四个子模块,技术基线为 Spring Boot 3.5.7 + Java 17)。理解这套分层,是读源码前最重要的一课:它不是"一个应用拆了几个包",而是刻意用"依赖方向"把"会变的东西"和"不变的东西"隔开。

┌──────────────────────────────┐
│ frontend │ React 18 + Vite
│ 问答 UI / 管理后台 / SSE │
└──────────────┬───────────────┘
│ HTTP + SSE
▼
┌──────────────────────────────┐
│ bootstrap │ DDD 业务编排
│ controller→service→dao │
│ rag/ knowledge/ ingestion/ │
└───────┬───────────────┬──────┘
│ │
depends ▼ ▼ depends
┌────────────────────┐ ┌───────────────────┐
│ infra-ai │ │ framework │
│ LLM/Emb/Rerank/VLM│◀─┤ 契约+追踪+缓存+ID │
│ OpenAI 兼容协议 │ │ 幂等+MQ+异常 │
└────────────────────┘ └───────────────────┘
│
│ HTTP/MCP 调用(独立进程,不编译依赖)
▼
┌────────────────────┐
│ mcp-server │ MCP 工具服务 (9099)
└────────────────────┘
2.1 四个模块各自的职责与关键落点
framework(地基,谁都依赖它,它不依赖任何人):这是整个项目的"标准库 + 统一语言"。它不写任何业务,只沉淀横切能力,包结构就是它的能力清单:convention(贯穿全项目的统一数据契约)、trace(RAG 链路追踪 SPI)、context(用户/租户上下文,含 TTL 传递)、cache(缓存抽象)、distributedid(分布式 ID)、idempotent(幂等注解 + AOP 切面)、mq(消息队列封装)、exception / errorcode(统一异常与错误码)、web(Web 层通用封装)、database(数据源/多租户相关)。bootstrap 和 infra-ai 都依赖它,但 framework 反过来不认识它们——这就是"地基不向上依赖"的铁律。
infra-ai(模型接入层):把"调模型"这件最容易变、最易出错的事隔离在此。对外只暴露稳定的接口形态(LLMService / EmbeddingService / RerankService / VLM 等),对内用"OpenAI 兼容协议"统一对接各家供应商,并把路由、熔断、首包探测等容错逻辑藏在这里。它依赖 framework(用统一契约和追踪),但不依赖 bootstrap。
bootstrap(业务编排层,应用主入口):DDD 分层,controller → service → dao。它把 framework 的横切能力和 infra-ai 的模型能力,"组装"成真正的产品功能:RAG 问答(rag/)、知识库(knowledge/)、文档入库(ingestion/)、对话管理、限流等。它是唯一依赖上述两者的模块,也是唯一和前端、数据库直接打交道的地方。
mcp-server(独立工具服务):见第 8 章。它是编译期完全独立、运行时单向被调用的另一个 JVM,仅依赖 spring-boot-starter-web + MCP SDK,不依赖另外三个模块。
2.2 依赖方向铁律
bootstrap → framework + infra-aiinfra-ai → frameworkmcp-server编译期独立:仅依赖spring-boot-starter-web+ MCP SDK;运行时由bootstrap的McpClientAutoConfiguration通过 HTTP/MCP 协议主动调用(方向是bootstrap → mcp-server,单向),可单独部署、独立扩缩容。
这种单向依赖不是文档约定,而是编译期强制的:framework 的源码里不可能 import 到 bootstrap 的类,否则就形成了循环依赖,Maven 直接编不过。所以"地基不向上依赖"既是设计原则,也是物理约束。
2.3 为什么这么分:依赖倒置的实战样板
分层的终极目的,是让"最容易变的外部依赖"被接口包住,核心业务不跟着抖。framework 提供稳定抽象,infra-ai 实现"具体供应商怎么调",bootstrap 只依赖抽象。三个最值得背下来的实例:
- 模型可换:
RoutingLLMService屏蔽了"用哪家模型、哪家挂了切哪家"。bootstrap调用的是LLMService接口,从不知道背后是 OpenAI 还是自建模型——换供应商只改infra-ai,主链路零改动。 - 向量库可换:
VectorStoreService接口屏蔽了"用 Milvus 还是 pgvector"。知识库建索引、检索的代码写一次,底下存储可替换。 - 对象存储可换:S3 / 阿里云 OSS 通过
framework或infra-ai的抽象封装,业务侧不感知。
这背后的思想就是依赖倒置(DIP):高层模块(bootstrap)不依赖低层模块(infra-ai)的细节,两者都依赖抽象(framework 里的接口/契约)。它和 Spring 的"面向接口编程 + 依赖注入"是一回事,但在本项目里被放大到了"模块级别"。读懂这一层,你就读懂了整个工程为什么"好改、好测、好扩"。
第 3 章:建议的学习路线图

不要从 main 方法一行行啃。Ragent 的代码量不小(仅 bootstrap 就有近 500 个 Java 文件),正确的读法是由底层契约 → 模型接入 → 业务编排 → 工具 → 前端,由静态结构到动态链路:先建立"统一语言",再看"能力怎么被实现",最后看"能力怎么被串成一次问答"。下面每条都标了"读什么 + 看的时候盯什么"。
-
README.md+assets/*.png:先看官方架构图(module-layering、chain、multi-channel-retrieval、model-routing-failover 四张)。目的不是记细节,而是先把"有哪些零件"装进脑子,后面读代码时能对号入座。 -
framework/convention(统一语言):读ChatMessage/ChatRequest/RetrievedChunk/SourceRef/Result/GroundingChunk。这些是贯穿全项目的"名词",几乎每个模块的方法签名里都能见到。重点看:消息体怎么表示角色( user / assistant / system )、检索片段RetrievedChunk带哪些字段(来源、分数、原文)、SourceRef如何关联溯源。先认字,后读书。 -
framework/trace(链路追踪 SPI):读RagTraceContext/RagTraceNode/RagStreamTraceSupport。理解"一次问答的每一步(检索、重排、生成)如何被记录成树状节点"。这是后面读bootstrapAOP 埋点的钥匙——你会发现很多追踪是靠注解 + 切面自动织入的,而这里的 SPI 定义了"可以被织入什么"。 -
infra-ai(模型接入):读LLMService(接口)→RoutingLLMService(路由 + 熔断 + 首包探测)→AbstractOpenAIStyleChatClient(OpenAI 兼容协议的统一 HTTP 封装)。盯着看:接口如何定义"发一条对话、收一个流";路由服务如何用"档位 + 熔断器三态"做容错;抽象基类如何把"不同供应商只是 baseUrl / 鉴权头不同"这件事收敛掉。 -
bootstrap/rag/core/retrieval(检索核心):核心中的核心——多通道检索 + 后处理责任链。读完你会明白"为什么单路向量召回率低":VectorSearchChannel/KeywordSearchChannel/ 图谱/网页通道并行跑,FusionPostProcessor(RRF 融合)与DeduplicationPostProcessor(去重)串成责任链,最后接 Rerank。这一块是面试高频区。 -
bootstrap/rag/core其余子包:intent/(意图树,IntentTreeFactory)、memory/(会话记忆,DefaultConversationMemoryService)、prompt/(Prompt 装配)、source/(溯源,SourcesAssembler)、graph/(知识图谱同步)、rewrite/(查询改写,MultiQuestionRewriteService)。建议按"一次问答真实会经过的顺序"读,而不是按字母序。 -
bootstrap/ingestion+knowledge:文档怎么进系统。IngestionEngine是一条可编排的 DAG 流水线(解析 → 切块 → 向量化 → 入库),knowledge管知识库元信息与权限。读懂它,你就理解了"问答之前,数据是从哪来的"。 -
bootstrap/ragController/Service(把能力串成一次问答):重点顺着RAGChatController→RAGChatServiceImpl→StreamChatPipeline(或等价编排类)读一遍。这是第 4 章要深拆的主链路,这里先建立"谁调谁"的骨架。 -
mcp-server:独立工具服务如何暴露、被调用(见第 8 章)。理解了第 2 章的"单向调用",这里会很顺。 -
frontend:看 UI 如何消费 SSE 流式事件。重点chatStore.ts(Zustand 状态)与useStreamResponse.ts(真实流解析),对应后端payload.type的各类事件(meta / think / response / finish / done / error)。后端吐什么、前端怎么接,在这一步闭环。
一句话心法:前 3 步是"认字和地基",4–7 步是"零件",第 8 步是"装配",9–10 步是"对外接口 + 表现层"。按这个顺序,你不会在 bootstrap 里被某个 RetrievedChunk 或 RagTraceNode 卡住——因为那些"名词"你早就认识了。
第 4 章:核心链路深度拆解(一次问答的全程)
这是全教程最重要的一章。建议你真正运行时,在 RAGChatController.chat() 打条件断点,发起一次问答,顺着调用栈走一遍——比读任何图都深。

4.1 入口:RAGChatController.chat()(SSE)
@IdempotentSubmit(
key = "T(com.nageoffer.ai.ragent.framework.context.UserContext).getUserId()",
message = "当前会话处理中,请稍后再发起新的对话"
)
@GetMapping(value = "/rag/v3/chat", produces = "text/event-stream;charset=UTF-8")
public SseEmitter chat(@RequestParam String question,
@RequestParam(required = false) String conversationId,
@RequestParam(required = false, defaultValue = "false") Boolean deepThinking) {
SseEmitter emitter = new SseEmitter(ragDefaultProperties.getSseTimeoutMs());
ragChatService.streamChat(question, conversationId, deepThinking, emitter);
return emitter;
}
要点:
- 用
@IdempotentSubmit(来自framework.idempotent)按用户 ID 防重:同一用户并发发问会被拦截,返回友好提示。这是幂等 AOP 在业务入口的典型用法。 - 返回
SseEmitter,超时时间来自RAGDefaultProperties.getSseTimeoutMs()(配置项,注意本地调大避免断流)。 GET+text/event-stream:前端用EventSource或 fetch 流读(见第 9 章)。
4.2 编排:RAGChatServiceImpl.streamChat()
public void streamChat(String question, String conversationId, Boolean deepThinking, SseEmitter emitter) {
String actualConversationId = StrUtil.isBlank(conversationId) ? IdUtil.getSnowflakeNextIdStr() : conversationId;
String taskId = IdUtil.getSnowflakeNextIdStr();
StreamCallback callback = callbackFactory.createChatEventHandler(emitter, actualConversationId, taskId);
chatQueueLimiter.enqueue(question, actualConversationId, emitter,
() -> traceRunner.run(question, actualConversationId, taskId, callback, traceAware -> {
StreamChatContext ctx = StreamChatContext.builder()
.question(question)
.conversationId(actualConversationId)
.taskId(taskId)
.deepThinking(Boolean.TRUE.equals(deepThinking))
.userId(UserContext.getUserId())
.callback(traceAware)
.build();
chatPipeline.execute(ctx);
}));
}
这一小段把全流程的关键设计一次性亮了出来,请逐行体会:
- 会话 ID 与任务 ID:空会话自动用雪花 ID 新建(
IdUtil.getSnowflakeNextIdStr(),来自 Hutool;workerId 由framework/distributedid的 Lua 脚本在 Redis 分配,避免多实例重复)。 StreamCallback:把 SSEemitter包装成回调,让下游 core 代码只管"产出事件",不关心怎么发到 HTTP(解耦,也便于单测)。chatQueueLimiter.enqueue(...):高并发保护入口。用 Redis ZSET 公平排队 + Lua 原子抢占 + 过期信号量 + Pub/Sub 唤醒(详见第 7 章限流小节)。保证模型服务不被瞬时流量打爆。traceRunner.run(...):包裹全链路追踪。内部把callback包成traceAware(实现了framework.trace.RagStreamTraceSupport),保证 SSE 跨线程时 Trace 上下文不丢。StreamChatContext:贯穿整条流水线的上下文对象(Builder 模式),持有 question / conversationId / taskId / deepThinking / userId / callback。chatPipeline.execute(ctx):真正执行问答流水线,把前面所有 core 能力串起来。
学习方法:把
RAGChatServiceImpl当作"目录页",顺着它持有的 5 个依赖(chatPipeline/chatQueueLimiter/callbackFactory/traceRunner/taskManager)逐个点进去读,就是一条完整的学习路径。
4.3 流水线内部(按执行顺序)

streamChat()
└─ chatQueueLimiter.enqueue() // 分布式限流排队
└─ traceRunner.run() // 开启链路追踪 + Trace 透传
└─ chatPipeline.execute(ctx)// 构建 StreamChatContext
├─ [Memory] 加载会话历史 + 摘要(memory/)
├─ [Rewrite] 查询改写/拆分(rewrite/)
├─ [Intent] 意图识别树判定分支(intent/)
├─ [Retrieval] 多通道并行召回(retrieval/channel/)
│ ├─ VectorSearchChannel (Milvus / pgvector)
│ ├─ KeywordSearchChannel (Elasticsearch)
│ ├─ GraphSearchChannel (LightRAG 知识图谱)
│ └─ WebSearchChannel (You.com 联网)
├─ [PostProcess] 责任链(postprocessor/)
│ ├─ DeduplicationPostProcessor 去重
│ ├─ FusionPostProcessor RRF 融合
│ ├─ RerankPostProcessor Rerank 重排
│ └─ ChannelAttributionPostProcessor 来源归因
├─ [Prompt] RAGPromptService 装配上下文 + 来源引用(prompt/)
├─ [LLM] RoutingLLMService.streamChat() 模型路由 + 熔断(infra-ai)
├─ [Source] SourcesAssembler 落库引用来源(source/)
└─ [Trace] RagTraceAspect 记录每个节点耗时/IO(framework.trace 落地)
SSE 事件流输出顺序(前端对应事件类型):
onMeta → onThinking(think) → onMessage(response) → onReject → onFinish → onDone
另有 onTitle / onCancel / onError
关键认知(面试常问):
- 检索是核心:多通道 + 融合 + 重排是 Ragent 区别于玩具 RAG 的关键。单路向量检索语义好但易漏关键词,关键词检索互补,知识图谱补关系,联网补时效——四路融合召回率最高。
- 模型调用全程容错:候选模型列表、首包探测、三态熔断,任何一个挂了自动切下一个,用户无感。
- 流式 + 异步 + TTL 透传:SSE 跨多个线程池(限流、检索、LLM),
framework/trace用TransmittableThreadLocal保证 Trace 上下文不丢。
第 5 章:framework 模块精读
根包:com.nageoffer.ai.ragent.framework,39 个文件、10 个包。它是"地基",不依赖任何业务模块。
怎么读这个模块:先读 convention/(领域统一语言,全项目都在用),再读 trace/(链路追踪,是后面 bootstrap AOP 的消费对象),其余包(idempotent/mq/distributedid/cache/context)随用随查即可,不必一次性读完。下面 5.1~5.3 是必读项,5.4 是速览。
5.1 convention:AI 领域统一契约(最该先读)
这是全项目的"统一语言",所有模块都在用:
| 类 | 文件 | 职责要点 |
|---|---|---|
ChatMessage |
convention/ChatMessage.java |
对话消息抽象。含 Role(SYSTEM/USER/ASSISTANT)、MessageStatus 枚举;除文本 content 外,还承载 thinkingContent(思维链)、sources(来源引用)、retrievedChunks(命中片段)、replyToMessageId(多轮关联) |
ChatRequest |
convention/ChatRequest.java |
大模型统一入参:messages + 采样参数 temperature/topP/topK/maxTokens + 能力开关 thinking/enableTools,Builder 模式 |
RetrievedChunk |
convention/RetrievedChunk.java |
向量检索命中:id/text/score + 富化后的 docId/chunkIndex/docName |
SourceRef |
convention/SourceRef.java |
文档级来源引用,按文档去重赋号(前端来源面板/预览用);含 sourceType(file/url/feishu) |
GroundingChunk |
convention/GroundingChunk.java |
仅 docName+text,给推荐追问生成做 grounding,不参与回答上下文 |
Result<T> |
convention/Result.java |
全局返回体:code/message/data/requestId,SUCCESS_CODE="0" |
数据流记忆:RetrievedChunk → 去重赋号成 SourceRef(给前端展示)/ 取最高分截断成 GroundingChunk(给追问生成) → 二者挂载到 ChatMessage 落库。
5.2 trace:RAG 链路追踪 SPI(工程含量最高的部分)
| 类 | 类型 | 职责 |
|---|---|---|
RagTraceContext |
final 工具类 | 基于 TransmittableThreadLocal 持有 TRACE_ID / TASK_ID / NODE_STACK(Deque)。API:pushNode / popNode / currentNodeId / depth / clear |
RagTraceNode |
注解 | @Target(METHOD),标记链路节点(name() / type()) |
RagStreamTraceSupport |
接口(SPI) | 跨线程 stream 节点追踪,内嵌 StreamSpan(detach / finishSuccess / finishError) |
必看亮点(面试加分):NODE_STACK 重写了 TTL.copy() 做深拷贝(new ArrayDeque<>(parentValue))。原因:默认 copy() 返回父值引用,并发子任务会共用同一 Deque,导致父子节点 ID 串挂、trace 层级紊乱。这是教科书级并发坑修复案例——务必在 RagTraceContext 里找到这段并读注释。
5.3 其余包速览(按需读)
exception/errorcode:三层异常体系 + 错误码规约(全局异常处理在web/包)。idempotent:幂等 AOP(@IdempotentSubmit/@IdempotentConsumer),防重提交 / 防重消费。mq+mq.producer:RocketMQ 生产者抽象与事务消息(保障入库分块/删除可靠执行)。context:Spring 容器 & 登录用户上下文(UserContext.getUserId(),业务里随处可取当前用户)。distributedid:Redis 协调的 Snowflake(lua/snowflake_init.lua分配 workerId/datacenterId)。database:MyBatis-Plus 字段自动填充(如createTime/updateTime);cache:Redis Key 前缀序列化;config:自动装配。
第 6 章:infra-ai 模块精读
根包:com.nageoffer.ai.ragent.infra,50 个文件。边界要记牢:它只管"模型接入",不含 RAG 业务、不含向量库(向量库的 Embedding 生产方和 Rerank 打分方在这里,但向量库与 RAG 编排都在 bootstrap);Prompt 管理也在 bootstrap,不在 infra-ai。
6.1 统一语言:四类能力 + Routing 实现
四类同级能力(接口 + Routing* 路由实现):
| 能力 | 接口 | 路由实现 |
|---|---|---|
| LLM 对话 | LLMService |
RoutingLLMService(@Primary) |
| Embedding | EmbeddingService |
RoutingEmbeddingService |
| Rerank | RerankService |
RoutingRerankService |
| VLM 视觉 | VlmService |
RoutingVlmService(仅入库期图生文) |
6.2 协议统一:全部走 OpenAI 兼容
这是关键设计——没有专门的 OpenAI/Claude 客户端类。所有厂商通过抽象基类承载 99% 逻辑:
AbstractOpenAIStyleChatClient/AbstractOpenAIStyleEmbeddingClient:构造 OpenAI 格式 body,含doChat()(同步) /doStreamChat()(SSE 流式)。- 具体子类平均仅 ~45 行,只覆写
provider()+ 加 trace 注解。 - 支持的
ModelProvider:OLLAMA(本地) /BAI_LIAN(阿里百炼/通义) /SILICON_FLOW(硅基流动) /AI_HUB_MIX(聚合网关,可间接接 OpenAI/Claude) /NOOP(占位)。
6.3 精读顺序与要点
AbstractOpenAIStyleChatClient:- 按
timeoutMs档位派生客户端并复用连接池(syncClientByTimeoutMap),避免每次新建 OkHttp 客户端。 doStreamChat():用 OkHttp 发请求,OpenAIStyleSseParser解析流,区分content(正文) 与reasoning_content(思考链)。
- 按
OpenAIStyleSseParser:把 SSE 文本流解析成结构化事件(思考/正文/结束/错误),是流式正确性的核心。RoutingLLMService(核心容错逻辑):- 构造时把候选模型列表注入(按档位/优先级排序)。
streamChat():按档位 + 健康状态路由到具体ChatClient。- 首包探测:发送后等待首个 token,超时或异常则降级到下一个候选模型。
- 依赖
model/ModelHealthStore(三态熔断:健康/半开/熔断,根据错误率自动切换)。
model/包(路由三件套如何协作):一次streamChat调用内部,ModelSelector先按档位(Tier)+ 健康状态筛选出候选列表(已剔除熔断中的模型),交给ModelRoutingExecutor,由它依次尝试每个候选——发请求后启动"首包探测"计时,首个成功返回 token 的即被采用,前面失败的自动跳过下一个;ModelHealthStore维护三态(健康/半开/熔断),根据错误率自动切换,下次ModelSelector就不会再选到它。配置校验则在启动时校验候选模型参数合法性,避免运行时才报错。http/包:ModelUrlResolver解决"请求打到哪个 URL"——优先级为模型自定义 URL > provider 默认 host + 配置的 endpoint;同时定义了统一的异常体系(把各家 SDK 异常归一化),让上层路由能一致地处理"超时/限流/空响应"并触发降级。
练手切入点:想新增一个模型供应商(如 DeepSeek),只需新建
DeepSeekChatClient extends AbstractOpenAIStyleChatClient,覆写provider()并加ModelProvider.DEEP_SEEK,再在 yaml 配候选列表即可,无需改动任何编排代码。
第 7 章:bootstrap 模块精读
启动类已在第 1 章说明。根包 com.nageoffer.ai.ragent,采用 DDD 风格"业务域优先" 分层。共 7 个业务域:
rag/ (231 文件) ★核心域:对话、检索、意图、记忆、MCP、向量
knowledge/ (68) 知识库/文档/切片管理 + 定时刷新 + MQ
ingestion/ (62) 文档摄入流水线(可编排 DAG)
core/ (56) 文档解析 + 切块(底层能力,无 Controller)
user/ (20) 用户/认证/权限(Sa-Token)
audit/ (12) 业务变更日志(bizlog-sdk)
admin/ (10) 运营看板 Dashboard
每个域内部统一分层:controller/(request/vo/) → service/(impl/) → dao/(entity/mapper/handler/) + config/、enums/、mq/。
7.1 rag/core/ 子能力地图(按推荐阅读顺序)
| 子包 | 职责 | 关键类 | 设计模式 |
|---|---|---|---|
retrieval/ |
多通道检索编排 | RetrievalEngine、MultiChannelRetrievalEngine、RetrievalBudget |
策略/责任链 |
retrieval/channel/ |
4 路检索通道 | SearchChannel(SPI) + Vector/Keyword/Graph/Web 实现 |
策略 |
retrieval/postprocessor/ |
结果后处理 | Deduplication → Fusion(RRF) → Rerank → ChannelAttribution |
责任链 |
vector/ |
向量存储双实现 | VectorStoreService(接口) → MilvusVectorStoreService/PgVectorStoreStoreService |
策略 |
vector/decorator/ |
写入时同步索引 | GraphSyncingVectorStoreService、KeywordSyncingVectorStoreService |
装饰器 |
vector/strategy/ |
并行检索 | AbstractParallelRetriever → CollectionParallelRetriever/IntentParallelRetriever |
模板方法 |
intent/ |
树形意图识别 | IntentTreeFactory(21KB)、DefaultIntentClassifier、IntentResolver、IntentNodeRegistry |
工厂/注册表 |
mcp/ |
MCP 客户端集成 | McpClientAutoConfiguration、LLMMcpParameterExtractor(17KB)、McpToolRegistry |
注册表 |
memory/ |
会话记忆 + 摘要 | DefaultConversationMemoryService、JdbcConversationMemorySummaryService |
— |
prompt/ |
Prompt 模板装配 | RAGPromptService、DefaultContextFormatter、PromptTemplateLoader |
模板方法 |
storage/ |
对象存储抽象 | S3ObjectStorageClient、OssObjectStorageClient |
策略 |
guidance/ |
歧义引导 | IntentGuidanceService、AmbiguityLLMChecker |
— |
aop/ |
链路追踪落地 | RagTraceAspect(消费 framework.trace 的 @RagTraceNode) |
AOP |
表格里
rewrite/graph/source/prompt/memory是核心子能力,下文 7.3~7.6 逐一展开(其余如storage/intent/mcp已在第 7 章对应小节讲过或见模式章)。
7.3 查询改写:rewrite/ —— 把"口语化问题"变成"可检索问题"
为什么需要:用户提问往往含糊、带上下文指代(“它”、“这个”)、或多问题揉在一起。直接拿去检索,召回质量差。MultiQuestionRewriteService(@RagTraceNode 标记)做两件事:
- 改写(Rewrite):用 LLM 把口语化/指代性问题规范成独立、自包含的问题(结合会话记忆里的上文)。
- 拆分(Split):若一个问题实际包含多个子问,拆成多条查询,各自检索后合并——提高复杂问题的覆盖度。
它依赖 PromptTemplateLoader 加载改写 Prompt(见 RAGConstant.QUERY_REWRITE_AND_SPLIT_PROMPT_PATH),用 LLMService 调用(按 Tier 选模型),并用 LLMResponseCleaner 清洗模型返回的 JSON(容错解析,避免模型输出非标准格式导致崩溃)。QueryTermMappingService 则负责把改写后的查询词映射回原文关键术语,辅助关键词通道召回。
7.4 知识图谱:graph/ —— 用 LightRAG 补"关系/实体"召回
为什么需要:向量检索擅长语义相似,但弱于"实体之间的关系路径"(如"A 公司的 CEO 是谁"这类多跳关系)。LightRagClient 接入 LightRAG 服务,把文档抽取成实体-关系图,检索时走图路径召回相关片段。
LightRagClient:@ConditionalOnProperty控制是否启用(未配置则整个图谱通道不加载,体现"可插拔"),用 OkHttp 调 LightRAG HTTP API,内部用 Jackson 解析 JSON、用正则处理返回格式。GraphQueryService:把用户问题转成图查询、调用LightRagClient、把结果包装成RetrievedChunk喂给后续融合。- 写入侧的同步由 7.5 的装饰器(
GraphSyncingVectorStoreService)负责——这里讲"读",装饰器讲"写",两者配合形成图谱的读写闭环。
7.5 引用溯源:source/ —— 让答案"可信任、可点击"
为什么需要:RAG 答案必须能溯源,否则用户无法判断"这是模型编的还是文档里的"。SourcesAssembler 是核心:
- 把检索片段(KB 命中)按文档去重、按相关度排序赋号(如
[1][2]),补齐来源类型(file/url/feishu)与外部链接。 - 产出文档级
SourceRef列表,既用于 SSE 下发、前端来源面板展示,也作为"行内引用角标"的唯一编号源(答案里[1]对应哪篇文档)。 - 摘录超长时按
EXCERPT_MAX_LENGTH=100截断加省略号。 CitationContextEnricher:把来源编号注入到 Prompt,让模型回答时显式标注引用;GroundingChunksAssembler:取最高分片段截断成GroundingChunk,专供"推荐追问生成"做 grounding,不参与回答上下文。
这一套是 Ragent 区别于玩具 RAG 的关键工程点之一——答案不但要准,还要可证伪、可点击溯源。
7.6 会话记忆与 Prompt 装配(补完)
memory/(DefaultConversationMemoryService+JdbcConversationMemorySummaryService):取最近 N 轮原文 + 更早历史的 LLM 持久化摘要(详见第 4 章 4.3 与面试题 Q8),多轮用replyToMessageId关联。prompt/(RAGPromptService+DefaultContextFormatter+PromptTemplateLoader):把"系统提示 + 检索上下文(带来源引用)+ 会话历史 + 当前问题"按模板装配成最终ChatRequest。DefaultContextFormatter决定检索片段如何格式化成上下文(含来源角标),PromptTemplateLoader从资源文件加载模板,与代码解耦、可热改。
7.7 数据模型与持久化(读 resources/database/schema_pg.sql)
Ragent 的持久化主线是"文档 → 切片 → 向量+来源 → 消息",核心表见第 0 章 0.1。与 bootstrap 业务域的对应关系:
knowledge/域管t_knowledge_document/t_knowledge_chunk(入库后落库);ingestion/把文档变成 chunk 并写t_knowledge_chunk的embedding(pgvector 列);rag/检索 chunk、生成t_message(含thinking_content/sources/message_status),并把引用来源写t_message_source;- 会话摘要落
t_conversation_summary(与消息分离,避免长会话拉全量历史); - 业务变更落
t_biz_change_log(配合@EnableLogInstance审计)。
读源码时,遇到某个 Service 不确定的"数据从哪来/存到哪",直接按表名反查
dao/mapper最快——这是理解 bootstrap 的捷径。
7.2 重点精读:检索编排(最有价值的一段)
MultiChannelRetrievalEngine:调度 4 个 SearchChannel,每个通道用独立线程池并行执行,最后汇总。建议重点看:
- 各通道如何并行(
CompletableFuture/ 线程池隔离),总超时如何控制(RetrievalBudget)。 SearchChannelSPI:新增一路检索只需实现该接口 +@Component,框架自动发现。
postprocessor/FusionPostProcessor:实现 RRF(Reciprocal Rank Fusion) 融合——把多路召回按各自的排名(rank)加权融合,公式 score = Σ 1/(k + rank),比简单拼接召回率更高。这是 RAG 工程里最经典的融合算法,值得把代码抄一遍理解。
后处理责任链顺序:Deduplication(跨通道去重)→ Fusion(RRF 融合)→ Rerank(用 RerankService 精排)→ ChannelAttribution(标注每个 chunk 来自哪路,供来源展示)。
7.3 意图树:IntentTreeFactory + IntentNodeRegistry
IntentTreeFactory(21KB):在启动时构建一棵意图树(根节点 + 子节点,如"知识库问答 / 闲聊 / 工具调用 / 拒答"等),每个节点关联分类提示词与下游处理策略。DefaultIntentClassifier:用 LLM 把用户问题归类到树上的某个叶子节点。IntentResolver:根据命中的节点决定走哪条处理分支。IntentNodeRegistry:运行期按节点 ID 快速查节点(接口极简,仅getNodeById(String))。registry 是个经典注册表模式落地。
public interface IntentNodeRegistry {
IntentNode getNodeById(String id);
}
7.4 会话记忆:DefaultConversationMemoryService
- 控制 Token 成本:取最近 N 轮 + 对更早的对话做持久化摘要(摘要也由 LLM 生成,存
JdbcConversationMemorySummaryService)。 - 多轮
replyToMessageId关联,保证上下文连续。
7.5 分布式限流:rag/service/ratelimit/ChatQueueLimiter
(对应第 4 章 chatQueueLimiter.enqueue)
- 用 Redis ZSET 公平排队 + Lua 脚本原子抢占 + 过期信号量 + Pub/Sub 唤醒,保护下游模型服务。
- 这是高并发 AI 系统的标配保护手段,代码在
ratelimit/包,建议结合resources/下的 Lua 脚本一起读。
7.6 入库流水线:ingestion/engine/IngestionEngine
- 文档摄入是可编排的 DAG 流水线:解析 → 切块 → Embedding → 向量写入 → (装饰器)同步图谱/关键词索引。
IngestionEngine负责调度这条流水线;失败时通过 RocketMQ 事务消息保证分块/删除可靠执行(见framework/mq)。- 配合
knowledge/的定时刷新与文档状态管理。
7.7 MCP 客户端:rag/core/mcp/McpClientAutoConfiguration
- 通过 yaml 配置远程 MCP Server 地址(如
rag.mcp.servers[0].url: http://localhost:9099)。 LLMMcpParameterExtractor(17KB):用 LLM 从用户问题抽取工具调用参数,做 Schema 校验后调用。McpToolRegistry:统一注册/查找 MCP 工具,供意图树"工具调用"分支使用。
第 8 章:mcp-server 模块精读
mcp-server 是 Ragent 的独立 MCP 工具服务:一个单独的 Spring Boot 应用,进程内嵌 MCP Java SDK(io.modelcontextprotocol),向 LLM / Agent 暴露一组可被直接调用的"工具"(天气、工单、销售线索、You.com 联网搜索)。它与 bootstrap 是编译期完全独立、运行时单向调用的两个 JVM:mcp-server 的 pom.xml 不依赖 bootstrap / framework,部署后可独立启停;bootstrap 通过 McpClientAutoConfiguration 在启动期连接它的 MCP 端点,把工具清单拉回主链路供 Agent 调度。bootstrap → mcp-server 是单向依赖,反过来不成立。这种隔离的好处是:工具服务的崩溃、升级、扩容都不会牵连问答主链路。
8.1 模块结构与启动
模块根目录为 mcp-server/src/main/java/com/nageoffer/ai/ragent/mcp/,共 8 个 Java 文件(含 2 个测试):
| 文件 | 作用 |
|---|---|
McpServerApplication.java |
启动类,@SpringBootApplication,监听端口 9099(与 bootstrap 的 8080 区分,避免同机冲突) |
config/McpServerConfig.java |
MCP Server 基础配置:装配 McpSyncServerExchange、McpServerFeatures 等 SDK 组件,并收集所有 @Bean 形式的 SyncToolSpecification 注册进服务器 |
executor/WeatherMcpExecutor.java |
天气工具(最短,建议最先读) |
executor/TicketMcpExecutor.java |
工单 / ITSM 查询工具 |
executor/SalesMcpExecutor.java |
销售线索 / CRM 查询工具 |
executor/YouComSearchMcpExecutor.java |
You.com 联网搜索工具 |
McpServerConfig 的关键职责不是"手写工具注册",而是把 Spring 容器里所有 McpServerFeatures.SyncToolSpecification 类型的 Bean 自动汇总到 MCP Server。也就是说,新增工具 = 新增一个 @Component 类、在里面 @Bean 产出一个 SyncToolSpecification 即可,配置类无需改动。这是典型的"约定优于配置 + 注册表模式":配置是稳定的骨架,工具是热插拔的插件。
8.2 一个 MCP 工具的"标准范式"
四个 executor 虽然对接的外部系统各异(天气 API、内部工单库、CRM、You.com),但结构上完全同构。以 WeatherMcpExecutor 为模板,一个工具由三部分组成:
① 工具类声明:@Component 注册进 Spring;常见再叠加 @ConditionalOnProperty 之类的条件注解,控制"有配置才启用"(见 8.5)。
② 工具元数据(buildTool()):用 MCP SDK 的 McpSchema.Tool 描述工具对外暴露的"能力目录"——名称、自然语言描述、以及用 JsonSchema 描述的入参结构。入参里有 type、description,必填字段放进 required 列表,枚举值用 enum。这段 JSON Schema 不是给自己看的,是给模型看的:模型据此决定何时调、传什么参。所以描述写得好坏直接决定工具是否被误用。
③ 处理函数(@Bean 产出的 SyncToolSpecification):把"工具元数据"和"处理逻辑"绑定在一起。处理逻辑签名固定为 (exchange, request) -> handleCall(request),从 request.arguments() 取参,最终返回 CallToolResult(成功 isError=false 带 TextContent;失败 isError=true)。
WeatherMcpExecutor 的骨架(节选)体现了这个范式:
@Component
public class WeatherMcpExecutor {
private static final String TOOL_ID = "weather_query";
@Bean
public McpServerFeatures.SyncToolSpecification weatherToolSpecification() {
return new McpServerFeatures.SyncToolSpecification(buildTool(),
(exchange, request) -> handleCall(request));
}
private Tool buildTool() {
Map<String, Object> properties = new LinkedHashMap<>();
properties.put("location", Map.of("type", "string",
"description", "城市或地区,如 '北京' / '上海浦东新区'"));
JsonSchema inputSchema = new JsonSchema("object", properties,
List.of("location"), null, null, null);
return Tool.builder().name(TOOL_ID)
.description("查询指定地区的实时天气与未来预报")
.inputSchema(inputSchema).build();
}
CallToolResult handleCall(CallToolRequest request) { /* 取参→校验→调用→格式化→返回 */ }
}
8.3 四个工具各自的工程细节
WeatherMcpExecutor(天气):入参只有 location。handleCall 里做基础判空后调用天气数据源,把返回的温湿度 / 风力 / 预报格式化为纯文本塞进 TextContent。它的价值不在业务逻辑复杂度,而在"范式样板"——读它你能最快看懂一个 MCP 工具的全貌。
TicketMcpExecutor(工单):体量最大(约 19KB),面向内部 ITSM / 客服工单系统。它通常不止一个动作,而是通过入参里的路由字段在一份代码里区分"查工单列表 / 查工单详情 / 改状态 / 催办"等多个子能力,用 switch 分流到不同私有方法。实现里强调对内部系统的鉴权与超时:调用内部 HTTP / RPC 时带连接超时、读超时,失败时返回 isError=true 且不把内部系统错误原文直接透传给模型(防信息泄露),而是归一为"查询失败,请稍后重试"之类的安全提示。结果格式化会按工单字段(单号、标题、状态、负责人、更新时间)拼出结构化文本,方便模型抽取。
SalesMcpExecutor(销售线索):体量次之(约 16KB),面向 CRM / 销售线索库。和工单工具同构:用 action 路由"查线索列表 / 查客户详情 / 查商机阶段"等子能力;对外部 CRM 调用同样做了超时与异常归一。它的入参往往带 keyword、stage、limit 等筛选条件,返回时按"客户名—联系人—阶段—金额"这样的销售口径拼文本。两个业务工具(Ticket / Sales)的"一工具多 action"设计,避免了为每个子能力都建一个 @Component,是对"工具即能力分组"的务实取舍。
YouComSearchMcpExecutor(联网搜索):对接 You.com Search API(GET https://ydc-index.io/v1/search,X-API-Key 鉴权),返回带链接和摘录的网页 / 新闻结果。它完整示范了一个外部 HTTP 工具该有的工程纪律:
- 入参
query、count(默认 5、最大 20)、freshness(day / week / month / year 枚举); - 调 API 用 JDK 原生
HttpClient,connectTimeout(10s)、timeout(10s); - 响应非 200 时不回显响应体(避免泄露账号 / 密钥信息),只抛"异常状态码";
- 结果格式化把
results.web与results.news两段合并、再统一截断到count,使count对外表达"返回总条数上限",与直觉一致、也省 token;摘录优先取description,缺失回退第一条snippet; - 防御式读取所有可选字段,任一段缺失都不崩。
8.4 请求如何从模型流到工具再回来
完整的"Agent 调 MCP 工具"链路是跨进程的:
bootstrap的 Agent 主链路判断需要外部工具,生成CallToolRequest(含工具名 + 入参 JSON);bootstrap侧的 MCP 客户端(McpClientAutoConfiguration装配的McpSyncClient/McpAsyncClient)通过 MCP 传输(stdio 或 SSE / HTTP)把请求发到9099端口;mcp-server的McpServerExchange路由到对应SyncToolSpecification的处理函数(即各 executor 的handleCall);- executor 调外部系统、格式化为
TextContent,返回CallToolResult; - 结果经原路回传
bootstrap,作为上下文喂回 LLM,模型据此继续生成最终回答。
理解这点很关键:mcp-server 是无状态工具执行器,它不持有会话、不做 RAG、不碰向量库;它只负责"给定入参 → 产出文本结果"。所有编排、记忆、融合都在 bootstrap。职责切得干净,所以是"工具服务"而非"另一个大脑"。
8.5 条件注册与"工具存在即等于可用"
mcp-server 里多数 executor 用 @ConditionalOnProperty 控制注册,例如 YouComSearchMcpExecutor 上挂 @ConditionalOnProperty(name = "YDC_API_KEY"):只有当环境变量 YDC_API_KEY 存在时才把这个工具 Bean 注册进 MCP Server。其设计意图与 bootstrap 通道"无 Key 即不启用"对齐——工具清单是给 LLM 消费的能力目录,登记一个缺 Key 不可用的工具只会诱导模型调用后失败、污染清单。因此这里遵循"工具存在 ⟺ 可用"原则:handleCall 内部仍对 Key 失效等运行期边界做兜底校验,但"启不启用"由条件注解在启动期决定。新增工具时,若该工具依赖某外部密钥,应仿照加 @ConditionalOnProperty,避免向模型暴露不可用的能力。
8.6 服务级重复 vs 模块耦合的取舍
注意 YouComSearchMcpExecutor 与 bootstrap 里的 YouComWebSearchChannel 是有意重复两套 You.com 调用逻辑,而非抽公共模块。原因在前文强调过:mcp-server 是零内部依赖、可独立部署的服务,抽公共模块会打破"编译期独立",引入对 bootstrap / framework 的反向依赖。因此这里按"服务级重复"处理——重复量小、边界清晰,换来的是部署与演进的完全解耦。代价是:You.com 契约(端点 / 参数 / 响应结构)变更时,两处需同步修改。这是"复用 vs 隔离"之间一个非常真实、可讲给面试官的架构权衡。
8.7 学习路径与扩展方式
建议阅读顺序:McpServerApplication → McpServerConfig(理解"工具如何挂到 MCP 服务")→ WeatherMcpExecutor(最短,吃透"工具 = 元数据 + 处理函数"范式)→ YouComSearchMcpExecutor(学外部 HTTP 工具的工程纪律)→ TicketMcpExecutor / SalesMcpExecutor(学"一工具多 action"的业务工具写法)。
扩展一个新工具的步骤(无需改主链路):
- 在
executor/下新建@Component类; - 写一个
@Bean返回McpServerFeatures.SyncToolSpecification,把buildTool()(工具元数据)与handleCall()(处理逻辑)绑在一起; - 若依赖外部密钥,加
@ConditionalOnProperty做条件注册; - 启动
mcp-server(9099 端口),bootstrap 侧McpClientAutoConfiguration会自动发现并把工具拉回主链路——主链路代码零改动。
这正是 MCP 协议的核心卖点:工具以标准化协议暴露,调用方与提供方解耦,能力可热插拔。对外部现成的 MCP Server,则在 bootstrap 的 yaml 增加 rag.mcp.servers 条目接入即可,mcp-server 模块本身无需改动。
第 9 章:前端结构与 SSE 联调
技术栈:React 18 + Vite 5 + TypeScript + Tailwind + Radix UI + Zustand + react-markdown + @antv/g6(图谱)。
9.1 目录与依赖要点
- 状态管理
zustand;表单react-hook-form + zod;表格@tanstack/react-table;Markdownreact-markdown + remark-gfm + rehype-raw/sanitize(注意做了 XSS 净化);虚拟列表react-virtuoso;图表recharts;代码高亮react-syntax-highlighter。 - 27 个页面级 TSX,覆盖用户问答界面与管理后台两套控制台。
9.2 前后端 SSE 协议(真实字段,来自 chatStore.ts)
前端 chatStore.ts 用 createStreamResponse 消费 SSE,回调对应后端事件类型:
| 前端回调 | 后端事件类型 (payload.type) |
作用 |
|---|---|---|
onMeta |
meta |
携带 conversationId + taskId(首包,建立会话/任务) |
onThinking |
think |
思考链增量(appendThinkingContent) |
onMessage |
response |
回答正文增量(appendStreamContent) |
onReject |
reject |
拒答/安全拦截内容增量 |
onFinish |
finish |
收尾:title/sources/messageStatus/messageId |
onTitle |
title |
会话标题 |
onCancel |
cancel |
用户停止,追加"(已停止生成)" |
onDone |
done |
流关闭,清理 streaming 状态 |
onError |
error |
错误,消息置 error 态 |
关键前端逻辑(值得读 chatStore.ts 全文件):
sendMessage:先本地构造userMessage+ 空的assistantMessage(状态streaming),再发起 SSE;URL 形如${API_BASE_URL}/rag/v3/chat?question=...&conversationId=...&deepThinking=...。appendStreamContent/appendThinkingContent:用 Zustand 的set做不可变更新,按streamingMessageId精准定位正在流式渲染的消息。onFinish:把后端回传的sources、messageStatus、messageId合并到消息;首条消息用title自动命名会话。cancelGeneration:调/rag/v3/stop?taskId=...后端停止任务(对应RAGChatController.stop→taskManager.cancel)。loadRecommended:回答完成后可异步请求推荐追问(用GroundingChunk做 grounding),状态机idle/loading/ready/error防止重复请求。
9.3 前端 SSE 解析的真实实现(读 hooks/useStreamResponse.ts)
chatStore.ts 里的 createStreamResponse 最终调用 useStreamResponse.ts,这段对理解"前端怎么稳健地吃 SSE"很有价值:
- 原生
fetch+ReadableStream手动解析(不是EventSource),因为要支持AbortSignal中断(用户点"停止")。readSseStream用TextDecoder把字节流累积到buffer,按行切分(\r?\n)。 - 标准 SSE 行协议解析:
event:行定事件名、data:行收集数据、空行触发一次dispatchEvent、以 :开头的行是心跳/注释直接跳过**(如:keep-alive)。 - 事件分发:
dispatchEvent把data:累积的 JSON 解析后,按eventName路由到onMeta/onThinking/onMessage/onFinish/onDone/onCancel/onReject/onTitle/onError回调(注意message事件里还要按payload.type==="think"再区分思考链与正文)。 - 重试与退避:
streamWithRetry默认retryCount=2、retryDelayMs=600,失败按指数退避(retryDelayMs * 2^attempt)重试;但signal.aborted时立即抛出,不重试(用户主动取消不算故障)。 createStreamResponse返回{ start, cancel },cancel通过AbortController.abort()中断读取,对应后端的taskManager.cancel。
这套实现比"直接用 EventSource"更可控:既能中断、又能重试、又能精确处理思考链与正文混在同一
message事件里的分情况。前端面试聊"流式渲染"时可以拿它当范例。
9.4 前端状态组织鸟瞰
- 全局会话状态集中在
stores/chatStore.ts(Zustand),核心数据结构是messages数组,每条消息带status(streaming/done/error/rejected)、content、thinkingContent、sources、title等字段,与后端t_message表字段一一对应。 conversations列表、recommendedQuestions状态也在此 store;UI 组件(ChatPage.tsx等)只读 store、发 action,不直接管 SSE——关注点分离很清晰。- 表单/表格等用
react-hook-form + zod、@tanstack/react-table,管理后台复用这些基建。
联调建议:打开浏览器 DevTools → Network,过滤
chat请求,看 SSE 事件流;对照chatStore.ts的 handlers 与useStreamResponse.ts的解析,理解每个事件如何驱动 UI 状态变化。这是理解"流式 AI 应用"的最佳实践。
第 10 章:设计模式落点地图(详细版)
设计模式不是"炫技",而是把反复出现的工程问题,用已被验证的结构化解法封装起来。本章对每个模式给出:【概念】→【项目落点(真实代码)】→【为什么用这个模式】→【设计意图/取舍】。读完你会明白:Ragent 的优雅在于"模式精准对应痛点",而非堆砌模式。
10.1 策略模式(Strategy)—— 让可互换的算法彼此独立
【概念】 定义一系列算法,把它们封装起来,使它们可互相替换,且算法的变化不影响使用方。核心是"面向接口编程 + 运行时多态"。
【项目落点】 SearchChannel 接口是典型策略:
public interface SearchChannel {
String getName();
boolean isEnabled(SearchContext context); // 按知识库维度动态开关
SearchChannelResult search(SearchContext context);
SearchChannelType getType();
}
MultiChannelRetrievalEngine 只依赖这个接口,启动时收集所有 SearchChannel Bean(4 个实现:Vector/Keyword/Graph/Web),运行时并行调度。同类还有 VectorStoreService(Milvus vs pgvector)、ObjectStorageClient(S3 vs OSS)。
【为什么】 检索策略是最易变的部分——今天加 ES 关键词检索,明天某个知识库不需要图谱。如果 if-else 堆在一起,会越来越乱。策略模式让"每一种检索方式"自成一个类,开关、权重都内聚在自己身上(如 VectorSearchChannel.isEnabled() 读 SearchChannelProperties)。
【设计意图】 配合"SPI 自发现"(Spring 扫描 @Component),新增一路检索 = 写一个类,编排代码零改动(开闭原则)。isEnabled(context) 还支持按上下文动态启停,避免"无图谱库硬跑图谱通道"的浪费。
10.2 模板方法模式(Template Method)—— 把"骨架"和"步骤"分离
【概念】 在抽象类中用 final 方法定义算法的骨架(固定步骤与顺序),把其中可变的"步骤"声明为抽象方法交给子类实现。父类控制流程,子类只填差异。
【项目落点】 AbstractParallelRetriever<T>:
public final List<RetrievedChunk> executeParallelRetrieval(String question, List<T> targets, int topK) {
float[] queryVector = retrieverService.embedAndNormalize(question); // 只算一次,所有任务共享
// 1. 并行提交到线程池(CompletableFuture.supplyAsync)
// 2. 逐个 join 收集,try/catch 统计 successCount/failureCount(单目标失败不影响整体)
// 3. 跨目标按 score 统一降序(见下方不变式说明)
allChunks.sort((a, b) -> Float.compare(scoreOf(b), scoreOf(a)));
// 4. 打印统计日志
}
骨架里的可变点(createRetrievalTask / getTargetIdentifier / getStatisticsName)留给子类。CollectionParallelRetriever(按知识库集合并行)和 IntentParallelRetriever(按意图并行)就是两个子类。
【为什么】 意图检索、全局检索都要"预计算 query 向量 → 对多目标起 Future → 收集 → 排序"。这段骨架在每处重复,差异只在"单个目标怎么查"。模板方法消除重复,且保证所有检索器行为一致(比如都做了单目标容错)。
【设计意图】 两个细节体现功力:① query 向量只算一次(embedAndNormalize 在循环外),省掉 N 次 Embedding 调用;② 出口统一排序的不变式(代码注释明确):各目标并行返回的子列表只在自身内部有序,addAll 拼接后跨目标名次 = 拼接顺序,而目标集合可能无序(如 HashSet)。若不统一排序,下游 RRF 融合的"名次基准"会失真、截断误砍高分。所以在通道出口兑现"该通道视角下的全局相关性排序"这一不变式,把脏活堵在边界上。
对比策略模式:这里"流程固定、仅步骤不同"→ 用模板方法;若是"整个算法整体可互换"→ 用策略。本项目
AbstractParallelRetriever用模板方法,而SearchChannel用策略,区分很清晰。
10.3 适配器模式(Adapter)—— 把异构外部系统"翻译"成统一接口
【概念】 将一个类的接口转换成客户期望的另一个接口,使原本不兼容的类可以一起工作。常用于"接入第三方、但又不希望业务代码耦合其 SDK"。
【项目落点】 infra-ai 的 AbstractOpenAIStyleChatClient:所有模型厂商(OpenAI/通义/硅基流动/Ollama/Claude 网关)都兼容 OpenAI 的 /v1/chat/completions 协议,于是项目只有一个抽象基类统一负责构造请求体、连接池复用(按 timeoutMs 档位缓存 OkHttp 客户端)、用 OpenAIStyleSseParser 解析流(同时处理 content 与 reasoning_content 思考链)。具体厂商子类平均仅覆写 provider() + 加 @RagTraceNode 注解,约 45 行。
同构的还有 RoutingEmbeddingService(见下)和 AbstractOpenAIStyleEmbeddingClient。
【为什么】 若每家居然写一套客户端,维护噩梦,且切换时要改业务代码。适配器让"异构外部系统"统一成内部接口。
【设计意图】 关键洞察是协议标准化(OpenAI 兼容生态)是降低集成成本的最大杠杆。不兼容的模型只需重写 doChat/doStreamChat 两个钩子,不影响其他 99% 逻辑——适配器 + 模板钩子的组合,既统一又灵活。
10.4 路由/代理 + 容错(Routing + Fallback)—— 多候选的"智能选路"
严格说这不是 GoF 23 经典模式,而是代理(Proxy)+ 策略 + 状态机的组合,在 AI 系统里极常用,单列一节。
【项目落点】 RoutingLLMService(@Primary,业务无感注入)与 RoutingEmbeddingService:
@Service
@Primary
public class RoutingEmbeddingService implements EmbeddingService {
private final ModelSelector selector;
private final ModelRoutingExecutor executor;
private final Map<String, EmbeddingClient> clientsByProvider; // provider -> 客户端 查表
...
@Override
public List<Float> embed(String text) {
return executor.executeWithFallback( // 关键:带降级执行
ModelCapability.EMBEDDING,
selector.selectEmbeddingCandidates(), // 候选列表(按优先级)
this::resolveClient,
(client, target) -> client.embed(text, target));
}
}
RoutingLLMService 同理:ModelSelector 选候选 → ModelHealthStore 三态熔断(健康/半开/熔断)→ ModelRoutingExecutor.executeWithFallback 按首个成功/首包探测降级。
【为什么】 模型供应商 SLA 不可靠,直接暴露给用户会体验崩坏。路由层把"选哪个、挂了换哪个"封装起来,业务代码只调 LLMService.streamChat(...),完全不感知背后在多家模型间切换。
【设计意图】 这是微服务容错(Resilience4j 那一套)在 LLM 场景的落地,只是触发条件从 HTTP 状态码变成"首包延迟/空响应"。clientsByProvider 是个**查表(Map)**而非 if-else,新增 provider 无需改分支。
10.5 装饰器模式(Decorator)—— 透明地给对象加职责
【概念】 动态地给一个对象添加额外职责,就增加功能来说,比继承更灵活。装饰器与被装饰者实现同一接口,内部持有被装饰者的引用。
【项目落点】 vector/decorator/GraphSyncingVectorStoreService implements VectorStoreService,内部持有真正的 delegate:
public class GraphSyncingVectorStoreService implements VectorStoreService {
private final VectorStoreService delegate;
// indexDocumentChunks → delegate.indexDocumentChunks(...); 然后 syncGraph(...)
}
所有写方法都是"先调 delegate,再同步图谱"。调用方(检索/入库)完全无感。同构还有 KeywordSyncingVectorStoreService。
【为什么】 写入向量是高频主链路。产品要"写入时同步更新知识图谱",若去改每个调用 VectorStoreService 的地方,侵入大、易错。装饰器把增强"包"在接口外面,对调用方透明。
【设计意图】 两个工程判断:
- best-effort 容错:
syncGraph()把异常吞掉只记 warn,图谱失败不影响向量主链路——图谱是"增强能力"而非"必须能力"。 - 粒度取舍:图谱抽取是文档级(跨 chunk 合并实体),所以只在
indexDocumentChunks(全量分块拼全文)和deleteDocumentVectors(按 docId 清)同步;子文档级的updateChunk/deleteChunkById不同步(靠整文重摄刷新)。作者想清楚了自己系统的真实写入路径,没被"理论上都该同步"绑架。
装饰器 vs 继承:装饰器可运行时组合、可叠加多个(向量 + 图谱 + 关键词层层包),比继承灵活;且对调用方透明。这正是
VectorStoreService接口存在的价值——只有面向接口,装饰才插得进去。
10.6 责任链模式(Chain of Responsibility)—— 一串可插拔的处理器
【概念】 使多个对象都有机会处理请求,将这些对象连成一条链,并沿链传递请求,直到有对象处理它。在"加工流水线"场景,更多是"每个节点都处理并返回给下一个"。
【项目落点】 retrieval/postprocessor/ 下每个处理器实现 SearchResultPostProcessor,按 getOrder() 串成链:Deduplication(order=1) → Fusion(RRF) → Rerank → ChannelAttribution。链头 DeduplicationPostProcessor:
@Override
public List<RetrievedChunk> process(List<RetrievedChunk> chunks,
List<SearchChannelResult> results,
SearchContext context) {
// 按 key 做集合去并:同一 Chunk 多路命中时保留首次出现的实例
// 不比较各通道原始分——BM25/余弦/图谱分跨量纲不可比,名次统一交由下游 RRF 融合赋分
Map<String, RetrievedChunk> chunkMap = new LinkedHashMap<>();
for (SearchChannelResult result : results) {
for (RetrievedChunk chunk : result.getChunks()) {
chunkMap.putIfAbsent(generateChunkKey(chunk), chunk);
}
}
return new ArrayList<>(chunkMap.values());
}
【为什么】 后处理是一串可增减、可调序的步骤。写死在一个方法里,调顺序要改代码;某步失败还要决定中断还是跳过。责任链让步骤可插拔、顺序可配(每个处理器自带 getOrder()),且单步失败可局部容错(如 Rerank 挂了可降级用融合结果)。
【设计意图】 去重这一环有个精妙点:注释明确写道不在这里比较跨通道分数——因为 BM25/余弦/图谱分跨量纲不可比,直接比分会出错;正确做法是"先去重保首次出现,名次统一交给下游 RRF 融合赋分"。这体现了"每个处理器只做一件事、且清楚自己职责的边界"。另外 generateChunkKey 用 SHA-256 而非 String.hashCode()——注释指出 32 位哈希碰撞概率不可忽略(“Aa” 与 “BB” 碰撞),碰撞会误删内容不同的 Chunk,故改用内容摘要。这是教科书级的"看似小事、实为隐蔽 bug"的规避。
10.7 注册表模式(Registry)—— 运行时按 ID 查找组件
【概念】 维护一个"名字/ID → 组件实例"的映射表,提供注册、查找、注销。调用方不关心组件从哪来(本地还是远程),只按 ID 取用。
【项目落点】 McpToolRegistry:
public interface McpToolRegistry {
void register(McpToolExecutor executor);
void unregister(String toolId);
Optional<McpToolExecutor> getExecutor(String toolId);
List<Tool> listAllTools();
boolean contains(String toolId);
int size();
}
启动时所有工具 @PostConstruct 注册进 DefaultMcpToolRegistry(内部 Map<toolId, executor>),调用方 O(1) 取执行器。IntentNodeRegistry.getNodeById(String) 同理(极简接口),意图树节点按 ID 快速查。
【为什么】 意图节点、MCP 工具在运行时被多处触发,若每次遍历扫描,低效且耦合。注册表把"查找逻辑"从调用方剥离,是插件化/可扩展系统的标配,常与 SPI、工厂配合。
【设计意图】 工具可能来自本地 @Tool 或远程 MCP Server,注册表让调用方对来源无感——这正是"依赖倒置"在组件查找层面的体现。
10.8 工厂模式(Factory)+ 构建者(Builder)—— 封装复杂对象的创建
【概念】 工厂把"对象怎么造"收拢到一处,调用方只声明"要什么",不关心"怎么造"。构建者则专门用于构造参数多的对象,链式可读。
【项目落点】
IntentTreeFactory(21KB):集中构建意图树——读节点定义、装配分类 prompt、挂接处理策略,对外只暴露一棵建好的树。业务侧只getNodeById取用。callbackFactory.createChatEventHandler(emitter, conversationId, taskId):把 SSE emitter 封装成StreamCallback的细节收进工厂。StreamChatContext.builder():用构建者装配贯穿流水线的上下文对象(question/conversationId/taskId/deepThinking/userId/callback)。framework里还有distributedid、cache等的工厂式配置类。
【为什么】 意图树、流式回调、分块策略的"创建逻辑"若散在业务里,职责混乱。工厂隔离复杂创建;构建者让多参数对象可读性高、且不易漏设必填项。
【设计意图】 与前面模式配合:工厂产出"策略对象"(如意图树节点),注册表存它,责任链/路由消费它——模式之间不是孤立的,而是协同成一套"可扩展架构"。
10.9 事件回调 / 观察者(Callback / Observer)—— 解耦"事件生产"与"消费"
【概念】 当某事件发生时,通知所有关心它的回调,而不必在生产者里硬编码消费者。SSE 流式输出就是典型应用。
【项目落点】 StreamCallback 把 SSE emitter 包成回调,下游 core 代码(检索、LLM)只管"产出事件"(onThinking/onMessage/onFinish…),不关心怎么发到 HTTP。前端 chatStore.ts 用 createStreamResponse 注册一组回调(onMeta/onThinking/onMessage/onReject/onFinish/onDone/onError),后端每个 SSE 事件对应一个回调驱动 UI 状态变化(详见第 9 章)。
【为什么】 问答是异步流式、跨多线程池的。如果核心逻辑里直接写 emitter.send(...),会把"业务编排"和"HTTP 协议"耦合,且无法单测。回调把二者解耦:core 层产出语义事件,传输层负责送达。
【设计意图】 回调 + 接口,使"首包探测"(LLM 客户端发请求后等首个 token)也能以回调形式通知路由层降级,而不阻塞主流程。
10.10 AOP(面向切面)—— 横切关注点统一织入
【概念】 把散布在多处、与业务无关的"横切逻辑"(日志、事务、权限、幂等、追踪)抽出来,通过切面统一织入,业务代码保持纯净。
【项目落点】 三处典型:
RagTraceAspect(rag/core/aop):用@RagTraceNode(name, type)注解标记方法,AOP 自动pushNode/popNode记录每个链路节点耗时(消费framework.traceSPI)。@IdempotentSubmit(framework.idempotent):IdempotentSubmitAspect在标记方法前后加分布式锁:
@Around("@annotation(com.nageoffer.ai.ragent.framework.idempotent.IdempotentSubmit)")
public Object idempotentSubmit(ProceedingJoinPoint joinPoint) throws Throwable {
if (evalEnabled) return joinPoint.proceed(); // 压测开关:绕过幂等
String lockKey = buildLockKey(joinPoint, idempotentSubmit);
RLock lock = redissonClient.getLock(lockKey);
if (!lock.tryLock()) {
throw new ClientException(idempotentSubmit.message()); // 重复提交直接拦
}
try {
return joinPoint.proceed();
} finally {
lock.unlock();
}
}
buildLockKey 支持两种维度:注解里配了 key(如 UserContext.getUserId(),SpEL 解析)→ 按用户锁;否则按 path + userId + 参数MD5 锁。还有 @IdempotentConsume 用于 MQ 消费防重。
- 审计日志:
@EnableLogInstance+ bizlog-sdk,业务变更自动落audit/表。
【为什么】 幂等、追踪、审计都是"横切"在业务方法上的,如果每个方法手写,重复且易漏。AOP 一处定义、处处生效。
【设计意图】 evalEnabled 开关很有意思——压测时可绕过幂等(因为压测会故意并发同请求),说明作者考虑到了"防护逻辑本身在特殊场景下要能关掉"。锁 key 的设计也体现"幂等维度可配置":同一用户连发(配 key=userId)vs 同请求参数连发(默认 MD5),是两种不同语义。
10.11 依赖倒置(DIP / IoC)—— 全项目的总纲
【概念】 高层模块不应依赖低层模块,二者都应依赖抽象;抽象不应依赖细节,细节应依赖抽象。配合 Spring IoC 容器完成"运行时注入实现"。
【项目落点】 四模块依赖单向:bootstrap → framework + infra-ai,且 bootstrap 只依赖 LLMService / VectorStoreService / SearchChannel 等接口。VectorStoreService 定义了 indexDocumentChunks/updateChunk/deleteDocumentVectors/...,MilvusVectorStoreService 与 PgVectorVectorStoreService 是两套实现,由配置决定注入哪个。mcp-server 更彻底——编译期零依赖,只通过 MCP 协议通信,可独立部署扩缩容。
【为什么】 如果业务编排直接 new MilvusVectorStoreService() 或硬编码某模型,未来替换成本极高。DIP 把"易变的外部依赖"隔离到独立模块/实现类,靠接口和配置解耦。
【设计意图】 这是贯穿全项目的"总纲"——前文所有模式(策略/装饰器/适配器/注册表)都是 DIP 在局部的具体落地。理解这一点,你就理解了 Ragent 架构为什么"好改":所有变化都被关进了接口背后的实现类里。
10.12 模式协同全景图
这些模式不是孤立的,而是协同成一套"可扩展架构":
┌─────────────────────────────┐
│ DIP / IoC(总纲:依赖接口) │
└──────────────┬──────────────┘
│ 接口背后是
┌──────────────┬───────────┼───────────┬──────────────┐
▼ ▼ ▼ ▼ ▼
策略 Strategy 装饰器 Decorator 适配器 Adapter 注册表 Registry 模板方法 Template
(SearchChannel)(同步图谱包装) (统一模型客户端) (McpTool/IntentNode)(并行检索骨架)
│ │ │ │ │
└──────┬───────┴─────┬─────┴─────┬─────┴──────┬───────┘
▼ ▼ ▼ ▼
工厂 Factory 产出上面这些对象 → 责任链 Chain 串起后处理 →
回调 Callback 解耦流式输出 → AOP 统一织入幂等/追踪/审计
一句话收尾:Ragent 的优雅不在于用了多少"高大上"模式,而在于每个模式都精准对应一个真实工程痛点——可变外部依赖用策略/装饰器/适配器隔离,重复骨架用模板方法收敛,运行时查找用注册表,异步上下文用 TTL 修正,稀缺资源用限流+熔断+幂等保护,横切逻辑用 AOP 抽离。读源码时带着"作者在这里遇到了什么问题"去想,比背模式定义收获大得多。
10.13 模式 → 源码索引速查表
| 模式 | 概念一句话 | 项目真实落点 | 解决的问题 |
|---|---|---|---|
| 策略 Strategy | 可互换算法各自成类 | SearchChannel、VectorStoreService、ObjectStorageClient |
检索/存储/模型策略易变,需独立替换 |
| 模板方法 Template | 父类定骨架、子类填步骤 | AbstractParallelRetriever、AbstractOpenAIStyleChatClient |
消除并行检索/模型调用的重复骨架 |
| 适配器 Adapter | 异构系统翻译为统一接口 | AbstractOpenAIStyleChatClient(各厂商子类) |
多模型厂商接入零改业务 |
| 路由+容错 Routing | 多候选智能选路+降级 | RoutingLLMService/RoutingEmbeddingService + ModelHealthStore |
模型供应商不稳定,需无感降级 |
| 装饰器 Decorator | 透明加职责 | GraphSyncingVectorStoreService、KeywordSyncingVectorStoreService |
写入时附加同步,不动主流程 |
| 责任链 Chain | 一串可插拔处理器 | postprocessor/*(去重→RRF→Rerank→归因) |
后处理步骤可调序、可容错 |
| 注册表 Registry | ID→实例 映射表 | McpToolRegistry、IntentNodeRegistry |
运行时按 ID 查找组件,来源无感 |
| 工厂+构建者 Factory/Builder | 封装复杂创建 | IntentTreeFactory、callbackFactory、StreamChatContext.builder() |
复杂对象创建不污染业务 |
| 回调/观察者 Callback | 事件生产消费解耦 | StreamCallback、前端 chatStore.ts 回调组 |
流式异步跨线程,编排与传输解耦 |
| AOP | 横切逻辑统一织入 | RagTraceAspect、@IdempotentSubmit、@EnableLogInstance |
追踪/幂等/审计不侵入业务 |
| 依赖倒置 DIP | 依赖抽象而非细节 | 四模块单向依赖 + 接口化服务 | 换模型/存储/工具不动业务 |
第 11 章:调试、上手与练手任务
11.1 本地跑起来
- 启动基础设施(Milvus 或 pgvector、Redis、MySQL、可选 RocketMQ、可选 Elasticsearch)。
- 启动
RagentApplication(9090)。 - 启动
McpServerApplication(9099)。 - 前端
cd frontend && npm i && npm run dev(5173,配置VITE_API_BASE_URL指向后端)。 - 本地可用 Ollama 跑本地模型(
ModelProvider.OLLAMA),避免 API Key 依赖。
11.2 调试技巧
- 在
RAGChatController.chat()打条件断点,发起一次问答,顺着RAGChatServiceImpl→StreamChatPipeline→ 各rag/core/*走一遍(对应第 4 章)。 - 在
RoutingLLMService.streamChat()看模型如何降级;在FusionPostProcessor看 RRF 融合数值。 - 读测试:30 个测试文件、84 个
@Test,重点覆盖模型路由、检索预算、结果去重、会话摘要、入库 Pipeline、MCP——测试是理解"正确行为"的最短路径。
11.3 练手任务(由易到难)
- 新增一路检索通道:实现
SearchChannel接口并加@Component(Spring 会自动把它收集进MultiChannelRetrievalEngine的通道列表,无需手动注册),再用 yaml 的SearchChannelProperties配置其启用与权重,验证被并行调度。 - 接入新模型:新建
XxxChatClient extends AbstractOpenAIStyleChatClient,覆写provider(),yaml 加候选列表。 - 新增 MCP 工具:在 mcp-server 加一个
@Tool,前端问答触发"工具调用"意图分支。 - 调优检索:改
RetrievalBudget的并发/超时,观察召回质量与延迟权衡。
11.4 编译注意
spotless要求每个 Java 文件有 Apache-2.0 License 头(文件顶部 16 行);新增文件务必先加,否则compile失败。- 新增 DAO 后若需扫描,检查
RagentApplication的@MapperScan是否覆盖对应包。
第 12 章:面试与讲解切入点(实战题 + 参考答案)
下面 12 道题均源于本项目真实代码,分「RAG 检索」「模型容错」「并发与限流」「可观测与链路」「工程架构」五个维度。每题给出【考点】【参考答案(结合本项目实现)】【可追问】三部分,便于你面试时展开讲。
维度一:RAG 检索
Q1. 为什么 RAG 要用多路检索,而不是只做向量检索?本项目怎么做的?
【考点】 检索召回率、不同检索范式的互补性、融合策略。
【参考答案】
单路向量检索(dense)语义好,但有两个固有短板:① 对精确关键词、专有名词、ID 不敏感;② 容易漏掉字面匹配但语义距离远的相关片段。本项目用 4 路并行检索通道 互补:
VectorSearchChannel(Milvus/pgvector,语义召回)KeywordSearchChannel(Elasticsearch,字面召回)GraphSearchChannel(LightRAG 知识图谱,关系/实体路径召回)WebSearchChannel(You.com 联网,时效召回)
四路在 MultiChannelRetrievalEngine 里各自独立线程池并行执行(避免慢通道拖累整体),结果汇总后进入后处理责任链:DeduplicationPostProcessor(跨通道去重)→ FusionPostProcessor(RRF 倒数排名融合 score=Σ 1/(k+rank))→ RerankPostProcessor(用 RerankService 精排到 contextTopK)→ ChannelAttributionPostProcessor(标注来源通道,供前端展示)。
关键判断点:RRF 是"无监督融合",不依赖模型、稳定;Rerank 是"有监督精排",质量高但慢且贵,放在融合之后只对融合后的 Top-N 跑,兼顾质量与成本。这就是生产级 RAG 的标准做法。
【可追问】 RRF 的 k 怎么选?Rerank 的候选 K 过大/过小分别有什么问题?(答:K 太大引入噪声拉低精排质量、且增加 Rerank 调用成本;本项目用 RetrievalBudget.contextTopK() 做预算控制。)
Q2. Rerank 之后为什么还要做"通道归因"?有什么工程价值?
【考点】 可观测性、成本优化、数据驱动的通道权重调优。
【参考答案】
RerankPostProcessor 的 logAttribution() 是个很值得讲的点:它对比 Rerank 前后各通道候选数量,重点统计"图谱证据存活率"(进入 Rerank 的图谱 chunk 有多少活到最后 Top-K)。
它的工程价值在于用数据决定去留:如果图谱通道大量进入候选却几乎不存活,说明它当前是"纯成本"——既占名额、又增加延迟,但没贡献。这时应下调 fusion.channel-weights.graph 或先优化其长证据的可排性,再决定保留还是下掉。这体现了"检索系统要可观测、可量化调优",而不是拍脑袋配权重。
【可追问】 如果多个通道权重都调过还是召回差,下一步怎么排查?(答:看 Deduplication 是否误杀、看 RetrievalBudget 预算是否过小、看 chunk size/overlap 是否合理、看 embedding 模型是否匹配业务语料。)
维度二:模型容错
Q3. 大模型 API 经常超时/限流/宕机,你怎么保障问答服务"不掉链子"?
【考点】 多模型冗余、路由策略、熔断降级、首包探测。
【参考答案】
本项目在 infra-ai 抽象了一层 RoutingLLMService(被 @Primary 注入,业务代码无感知)。核心机制:
- 候选模型列表:构造时注入多个候选(如 Ollama 本地 + 阿里百炼 + 硅基流动),按档位/优先级排序。
- 健康熔断:
ModelHealthStore维护三态(健康/半开/熔断),根据错误率自动切换,避免持续打挂的模型拖累。 - 首包探测:发出请求后等待首个 token,若超时或异常,立即降级到下一个候选模型,用户侧无感。
- 协议统一:所有厂商复用
AbstractOpenAIStyleChatClient(99% 逻辑共用,子类仅 ~45 行),切换供应商零改业务代码。
延伸:这种"路由 + 熔断 + 降级"思想和微服务里的 Resilience4j 一致,只不过触发点是 LLM 的"首包延迟/空响应"而非 HTTP 状态码。
【可追问】 降级时正在生成的 SSE 流怎么处理?上下文/温度要重新设吗?(答:降级本质是新发一次请求,复用同一 ChatRequest 与上下文,仅换底层 client;前端通过 onError/onReject 兜底。)
Q4. 你们接入了那么多模型厂商,怎么避免为每个厂商写一套客户端?
【考点】 适配器模式、协议标准化、OpenAI 兼容生态。
【参考答案】
关键认知:主流模型基本都兼容 OpenAI 的 /v1/chat/completions 协议。所以本项目没有"OpenAIClient / ClaudeClient / 通义Client"这种平行类,而是只有一个抽象基类 AbstractOpenAIStyleChatClient,统一负责:
- 构造 OpenAI 格式请求体(messages / temperature / stream)
- 连接池复用(按
timeoutMs档位syncClientByTimeout缓存 OkHttp 客户端) OpenAIStyleSseParser解析流,同时支持content(正文)和reasoning_content(思考链)
具体厂商子类平均仅覆写 provider() 返回枚举 + 加一个 @RagTraceNode 注解即可。支持 OLLAMA / BAI_LIAN / SILICON_FLOW / AI_HUB_MIX / NOOP。新增 DeepSeek 只需加一个子类 + yaml 候选列表。
【可追问】 不兼容 OpenAI 协议的模型(如某些私有模型)怎么办?(答:基类的 doChat/doStreamChat 留了可覆写点,子类可重写请求/解析逻辑,不影响其他 99%。)
维度三:并发与限流
Q5. SSE 流式问答 + 高并发,你怎么保护下游模型不被打爆?
【考点】 分布式限流、公平排队、信号量、与 SSE 的协作。
【参考答案】
入口在 RAGChatController.chat() 调 chatQueueLimiter.enqueue(...)。本项目用 FairDistributedRateLimiter(ratelimit/ 包),核心设计:
- 基于 Redis ZSET 做公平排队(先到先得),不是简单计数器,避免"突发流量把后面全拒"的饥饿问题。
- Lua 脚本原子抢占许可(保证分布式下并发安全);过期信号量 + 租约(
lease-seconds=600兜底释放,防止某连接卡死永久占坑)。 - Pub/Sub 唤醒:许可释放时通知等待者,不用忙等。
- 配置在
RAGRateLimitProperties:max-concurrent=50、max-wait-seconds=20(排队超时才 reject)。 - reject 也要优雅:
handleReject()会落一条REJECTED状态的 assistant 消息并向前端发meta→reject→finish→done事件,保证前端永远收到DONE、不会一直转圈。
【可追问】 为什么不用 Sentinel/Resilience4j?(答:SSE 是长连接 + 异步,需要"按用户公平排队 + 排队可超时 + 释放时唤醒"的语义,通用限流组件不好直接套;且 reject 路径要和业务状态联动落库。)
【可追问】 限流器本身挂了怎么办?(答:globalEnabled=false 时走直通分支,仅受线程池容量保护,最坏情况线程池拒绝再 reject,不会阻塞主流程。)
Q6. 幂等怎么做的?用户手抖连点两次会怎样?
【考点】 幂等设计、防重提交、与限流的区别。
【参考答案】
RAGChatController.chat() 上有 @IdempotentSubmit(key = "T(...UserContext).getUserId()", message="当前会话处理中,请稍后再发起新的对话")。这是 framework.idempotent 的 AOP——按用户维度加分布式锁/标记,同一用户并发发问直接被拦截并友好提示,不会开两条流水线抢同一个 SSE。
注意它和 Q5 限流的区别:限流是"保护模型"(全局并发数),幂等是"保护用户"(防同一用户重复触发)。两者正交,一个在入口防重、一个在队列控量。
【可追问】 用户第一个问题还没答完,又发第二个不同问题,能并行吗?(答:当前按 userId 加幂等键,会拦;若业务需要"同会话串行、不同会话并行",可把 key 改成 conversationId。这是个真实的设计取舍点。)
维度四:可观测与链路
Q7. 流式问答跨了限流线程池、检索线程池、LLM 异步回调,你怎么追踪一次请求的完整链路?
【考点】 跨线程上下文传递、TransmittableThreadLocal、AOP 埋点。
【参考答案】
framework.trace.RagTraceContext 用 TransmittableThreadLocal(TTL) 持有 TRACE_ID / TASK_ID / NODE_STACK。之所以必须用 TTL 而不是 ThreadLocal,是因为 SSE 全程跨多个线程池,普通 ThreadLocal 在异步切换时丢失。
这里有个教科书级坑点:NODE_STACK 是个 Deque(记录调用栈节点),默认 TTL 的 copy() 返回的是父值引用,并发子任务会共用同一个 Deque,导致父子节点 ID 串挂、层级紊乱。所以本项目重写了 copy() 做深拷贝(new ArrayDeque<>(parentValue)),每个子线程拿到独立副本。
埋点落地在 rag/core/aop/RagTraceAspect,用 @RagTraceNode(name, type) 注解标记节点,AOP 自动 pushNode/popNode 并记录耗时/IO。SSE 这种跨线程流还实现了 RagStreamTraceSupport(内嵌 StreamSpan 的 detach/finishSuccess/finishError),保证流结束才关闭节点。
【可追问】 为什么不用 SkyWalking/OpenTelemetry 的 Java Agent?(答:RAG 的"节点"是业务语义——检索、融合、重排、LLM——Agent 看不见这些,需要自定义 span;且要和控制台 UI 的 trace 视图打通,自研 SPI 更可控。)
Q8. 会话记忆怎么控制 Token 成本?长对话不会爆上下文吗?
【考点】 上下文窗口管理、记忆压缩、持久化摘要。
【参考答案】
DefaultConversationMemoryService 做了两级控制:
- 最近 N 轮直接拼进 prompt(近期上下文最相关);
- 更早的历史做持久化摘要:用 LLM 把老对话压缩成一段 summary,存
JdbcConversationMemorySummaryService,需要时召回,而不是把全部历史塞进去。
这样把"无限历史"变成"N 轮原文 + 1 段摘要",Token 量可控。配合 MemoryProperties(如 titleMaxLength 控制标题长度)等预算配置。
【可追问】 摘要会丢失细节吗?怎么权衡?(答:会,但摘要是可接受的损失;关键事实可强制保留在系统提示词或知识库里,不依赖对话记忆。这是 RAG + 记忆分层的标准取舍。)
维度五:工程架构
Q9. 向量写入时要同步更新知识图谱,你怎么在不改调用方的前提下加这个能力?
【考点】 装饰器模式、AOP 之外的行为增强、best-effort 容错。
【参考答案】
本项目用装饰器模式:GraphSyncingVectorStoreService implements VectorStoreService,内部持有真正的 delegate。所有 indexDocumentChunks / deleteDocumentVectors 等方法都是"先调 delegate,再同步图谱"。
亮点:
- 一处覆盖全部写调用点,调用方(检索/入库)完全无感,符合开闭原则。
- 同构的还有
KeywordSyncingVectorStoreService,证明这套装饰器是可复制的扩展范式。 - best-effort 容错:
syncGraph()把异常吞掉只记 warn,图谱失败不影响向量主链路——因为图谱是"增强能力"而非"必须能力"。这是个很成熟的设计判断。 - 粒度取舍:图谱抽取是文档级(跨 chunk 合并实体),所以只在
indexDocumentChunks(全量分块拼全文)和deleteDocumentVectors(按 docId 清)同步;子文档级的updateChunk/deleteChunkById不同步(Phase1 靠整文重摄刷新)。说明作者想清楚了自己系统的真实写入路径,没有被"理论上要同步"绑架。
【可追问】 如果图谱和向量要强一致怎么办?(答:装饰器改成事务型,要么都用同一事务/原子操作,要么引入补偿任务对账;但本项目明确图谱是增强,故选择最终一致 + 不强一致,更简单可靠。)
Q10. 文档入库(解析→切块→Embedding→写向量)怎么保证不丢数据、不重复?
【考点】 数据一致性、事务消息、可靠性。
【参考答案】
入库是 ingestion/engine/IngestionEngine 调度的可编排 DAG 流水线。一致性保障在 framework.mq + framework.idempotent:
- 分块、删除这类"最终要落库"的动作通过 RocketMQ 事务消息保证可靠执行——本地事务和 MQ 半消息要么都成、要么回滚,避免进程崩溃导致"文档状态变了但分块没写"。
- 配合
@IdempotentConsumer防重消费,避免消息重试造成重复分块。
【可追问】 大文档切块很多,一条消息装不下/要并行怎么办?(答:按文档分片发多条事务消息,消费端按 docId 聚合;或用"先落任务表、worker 拉取"的可靠队列模式。本项目用 MQ + 幂等消费覆盖。)
Q11. 你的模块划分依据是什么?为什么把 AI 接入单独拆成 infra-ai?
【考点】 分层架构、依赖倒置、单一职责、可替换性。
【参考答案】
四模块:framework(地基:契约/追踪/幂等/ID/MQ)、infra-ai(AI 接入:LLM/Embedding/Rerank/VLM)、bootstrap(业务编排)、mcp-server(独立工具)。依赖单向:bootstrap → framework + infra-ai,infra-ai → framework,mcp-server 编译期零依赖、只通过 MCP 协议通信。
拆 infra-ai 的核心原因是依赖倒置:业务层只依赖 LLMService / VectorStoreService 等接口,具体用 Milvus 还是 pgvector、用哪家模型,都是实现类 + 配置。这样:
- 换模型:加子类 + yaml,不动业务;
- 换向量库:换
VectorStoreService实现; - mcp-server 独立部署,可单独扩缩容、单独升级工具。
【可追问】 这种拆法在小项目里是不是过度设计?(答:是 trade-off。但 RAG 系统天然要对接多家模型/多种存储,把"易变的外部依赖"隔离到独立模块,长期收益远大于初期成本。)
Q12. 怎么保证这套系统"可扩展"——加一个新检索通道、一个新工具要改多少代码?
【考点】 扩展点设计、SPI/注册表模式、配置驱动。
【参考答案】
本项目大量使用接口 + 注册表 + 配置隔离,典型路径:
- 新检索通道:实现
SearchChannelSPI 加@Component,MultiChannelRetrievalEngine自动发现;权重在SearchChannelProperties(10KB 配置类)里配,不用改编排代码。 - 新工具:mcp-server 加一个
@Tool;或在 bootstrap yaml 加rag.mcp.servers条目接入外部 MCP Server,McpToolRegistry自动注册。 - 新模型:见 Q4。
这背后是"对修改封闭、对扩展开放"的落地:扩展成本 = 写一个实现类 + 一行配置,核心编排零改动。
【可追问】 注册表模式在运行时怎么防冲突/怎么热更新?(答:注册发生在 Spring 启动时(@Component 扫描 / @PostConstruct 注册),冲突可在注册时校验重复 ID 抛错;热更新需结合配置中心 + 动态注册,本项目走重启生效,足以覆盖大多数场景。)
面试收尾话术
如果面试官问"这个项目你最大的收获是什么",建议从工程权衡角度讲(比罗列功能更有分量):
- RAG 不是调通 LLM,而是把检索、融合、重排、容错、限流、溯源整条链路工程化;
- 所有"易变外部依赖"(模型/向量库/工具)必须隔离到独立模块,靠接口和配置解耦;
- 增强能力(图谱/联网)要 best-effort,绝不拖垮主链路;
- 可观测性不是附加项——没有通道归因、没有链路追踪,就没法调优和排障。
第 13 章:设计模式与工程问题深度讲解
本章不从"模式名词"出发,而从"工程问题"出发:先讲这个项目真实遇到的痛点,再讲它用什么设计模式/架构手段解决,最后给出代码落点。这样阅读源码时,你能在每一处都回答"作者为什么这样写"。
13.1 问题一:要同时支持向量/关键词/图谱/联网检索,且随时增减 —— 怎么组织?
工程问题
检索策略是系统里最易变的部分。今天用 Milvus,明天可能加 ES 关键词检索;某个知识库不需要图谱,要能关掉。如果把这些逻辑 if-else 堆在一个方法里,会越来越不可维护。
解法:策略模式(Strategy)+ SPI 自发现
项目抽象出 SearchChannel 接口:
public interface SearchChannel {
String getName(); // 通道名,用于日志/监控
boolean isEnabled(SearchContext context);// 按上下文动态开关(如某知识库无图谱则跳过)
SearchChannelResult search(SearchContext context);
SearchChannelType getType();
}
四个实现(Vector/Keyword/Graph/WebSearchChannel)各自封装一种策略。MultiChannelRetrievalEngine 只依赖接口,启动时把所有 SearchChannel Bean 收进列表,运行时并行调度。新增一路检索 = 写一个实现类 + 加 @Component,编排代码零改动(开闭原则)。isEnabled(context) 还支持按知识库维度动态启用/关闭,避免"无图谱库硬跑图谱通道"的浪费。
考点:策略模式 vs 简单工厂;SPI 自发现(Spring 扫描)比手动 switch 更利于扩展。
13.2 问题二:并行检索的"骨架代码"高度重复 —— 怎么消除?
工程问题
意图检索、全局检索都要"预计算 query 向量 → 对每个目标起 Future → 收集结果 → 统计成败 → 排序"。这段骨架在每个检索器里重复,差异只在"单个目标怎么查"。
解法:模板方法模式(Template Method)
AbstractParallelRetriever<T> 把不变骨架写在 executeParallelRetrieval(用 final 锁定),把变化点留给三个抽象方法:
public final List<RetrievedChunk> executeParallelRetrieval(String question, List<T> targets, int topK) {
float[] queryVector = retrieverService.embedAndNormalize(question); // 只算一次,共享给所有任务
// 1. 并行提交到线程池(CompletableFuture.supplyAsync)
// 2. 逐个 join 收集,try/catch 统计 successCount/failureCount(单目标失败不影响整体)
// 3. 关键:跨目标按 score 统一降序(见下方不变式说明)
allChunks.sort((a, b) -> Float.compare(scoreOf(b), scoreOf(a)));
// 4. 打印统计日志
}
两个值得背的细节:
- query 向量只算一次(
embedAndNormalize在循环外),所有并行任务共享,避免 N 个目标 N 次 Embedding 调用——这是真金白银的成本优化。 - 出口统一排序的不变式(代码 100-104 行注释写得很清楚):各目标并行返回的子列表只在"自身内部"有序,
addAll拼接后跨目标名次等于拼接顺序;而目标集合本身可能无序(如HashSet)。若不统一排序,下游 RRF 融合的"名次基准"会失真,导致截断时误砍高分。所以在通道出口兑现"该通道视角下的全局相关性排序"这一不变式,把脏活堵在边界上。
子类只需实现 createRetrievalTask(单个目标怎么查)、getTargetIdentifier、getStatisticsName。CollectionParallelRetriever(按知识库集合并行)和 IntentParallelRetriever(按意图并行)就是这么来的。
考点:模板方法 vs 策略——前者"父类定流程、子类填步骤",后者"整体可互换";这里流程固定、仅步骤不同,所以用模板方法。
13.3 问题三:接了七八家模型厂商,难道每家居然写一套客户端?
工程问题
OpenAI、通义、硅基流动、Ollama、Claude……每家 SDK 不同,若各写一套客户端,维护噩梦,且切换时要改业务代码。
解法:适配器模式(Adapter)+ 协议标准化
关键洞察:主流模型基本都兼容 OpenAI 的 /v1/chat/completions 协议。于是项目只有一个抽象基类 AbstractOpenAIStyleChatClient,统一负责构造请求体、连接池复用(按 timeoutMs 档位缓存 OkHttp 客户端)、用 OpenAIStyleSseParser 解析流(同时处理 content 正文与 reasoning_content 思考链)。
具体厂商子类平均仅覆写 provider() 返回枚举 + 加 @RagTraceNode 注解,约 45 行。不兼容的模型可重写 doChat/doStreamChat 两个钩子,不影响其他 99%。
考点:适配器让"异构外部系统"统一成内部接口;协议标准化(OpenAI 兼容)是降低集成成本的最大杠杆。
13.4 问题四:向量写入时要顺带同步图谱/关键词索引,怎么加而不改调用方?
工程问题
写入向量是高频主链路。现在产品要"写入时同步更新知识图谱"。如果去改每一个调用 VectorStoreService 的地方,侵入大、易错。
解法:装饰器模式(Decorator)
GraphSyncingVectorStoreService implements VectorStoreService,内部持有真正的 delegate:
public class GraphSyncingVectorStoreService implements VectorStoreService {
private final VectorStoreService delegate; // 真实实现
// indexDocumentChunks → delegate.indexDocumentChunks(...); 然后 syncGraph(...)
}
所有写方法都是"先调 delegate,再同步图谱"。调用方(检索/入库)完全无感。同构的还有 KeywordSyncingVectorStoreService。
两个工程判断(面试加分):
- best-effort 容错:
syncGraph()把异常吞掉只记 warn,图谱失败不影响向量主链路——因为图谱是"增强能力"而非"必须能力"。 - 粒度取舍:图谱抽取是文档级(跨 chunk 合并实体),所以只在
indexDocumentChunks(全量分块拼全文)和deleteDocumentVectors(按 docId 清)同步;子文档级的updateChunk/deleteChunkById不同步(靠整文重摄刷新)。说明作者想清楚了自己系统的真实写入路径,没被"理论上都该同步"绑架。
考点:装饰器 vs 继承——装饰器可运行时组合、可叠加多个(向量+图谱+关键词可层层包),比继承灵活;且对调用方透明。
13.5 问题五:检索结果要去重、融合、重排、归因,步骤还会调整顺序 —— 怎么编排?
工程问题
后处理是一串可增减、可调序的步骤。写死在一个方法里,调顺序要改代码;某步失败还要决定是中断还是跳过。
解法:责任链模式(Chain of Responsibility)
postprocessor/ 下每个处理器(去重 → RRF 融合 → Rerank → 通道归因)实现统一接口,按顺序串成链。FusionPostProcessor 用 RRF 倒数排名融合(score=Σ 1/(k+rank))把多路召回合并;RerankPostProcessor 再对融合后的 Top-N 用 RerankService 精排;ChannelAttributionPostProcessor 给每个 chunk 标注来源通道。
价值:步骤可插拔、顺序可配;单步失败可局部容错(如 Rerank 挂了可降级用融合结果)。RRF 与 Rerank 的分工(见 13.1 的 Q1)本身也是"无监督融合前置、有监督精排后置"的成本权衡。
考点:责任链 vs 管道(Pipeline)——这里每一步消费上一步的输出、产出给下一步,且强调"顺序与可插拔",是经典责任链。
13.6 问题六:意图树、MCP 工具散落各处,调用方怎么快速找到?
工程问题
意图节点、MCP 工具在运行时被多处触发,若每次都遍历扫描,低效且耦合。
解法:注册表模式(Registry)
McpToolRegistry 接口统一了工具的注册/注销/查找:
public interface McpToolRegistry {
void register(McpToolExecutor executor);
void unregister(String toolId);
Optional<McpToolExecutor> getExecutor(String toolId);
List<Tool> listAllTools();
boolean contains(String toolId);
int size();
}
启动时所有工具 @PostConstruct 注册进 DefaultMcpToolRegistry(内部 Map<toolId, executor>),调用方按 ID O(1) 取执行器,无需感知工具从哪来(本地 @Tool 还是远程 MCP Server)。IntentNodeRegistry 同理(getNodeById(String) 极简接口),意图树节点按 ID 快速查。
考点:注册表把"查找逻辑"从调用方剥离,是插件化/可扩展系统的标配;常与 SPI、工厂配合使用。
13.7 问题七:复杂对象(意图树、流式回调、分块策略)怎么创建才不污染业务?
工程问题
意图树要在启动时构建一棵带分类提示词和下游策略的树;流式回调要按 SSE emitter 包装;分块策略多样。这些"创建逻辑"若散在业务里,职责混乱。
解法:工厂模式(Factory)
IntentTreeFactory(21KB)集中负责意图树的构建——读取节点定义、装配分类 prompt、挂接处理策略,对外只暴露一棵构建好的树。业务侧只 getNodeById 取用,不关心怎么造的。callbackFactory.createChatEventHandler(emitter, conversationId, taskId) 同理,把"SSE emitter → StreamCallback"的封装细节收进工厂。
考点:工厂隔离复杂创建;与构建者(StreamChatContext.builder())配合,让业务代码只声明"要什么",不关心"怎么造"。
13.8 问题九:流式问答跨了多个线程池,链路上下文怎么不丢?
工程问题(最硬核的一处)
一次问答经历了限流线程池、检索线程池、LLM 异步回调,普通 ThreadLocal 在 CompletableFuture 切换线程时直接丢失,导致 Trace ID、节点栈全断。
解法:TransmittableThreadLocal(TTL)+ 深拷贝修正
framework.trace.RagTraceContext 用 TTL 持有 TRACE_ID / TASK_ID / NODE_STACK(Deque 记录调用栈节点)。但有个教科书级坑:NODE_STACK 是 Deque,TTL 默认 copy() 返回父值引用,并发子任务会共用同一 Deque,导致父子节点 ID 串挂、trace 层级紊乱。本项目重写了 copy() 做深拷贝(new ArrayDeque<>(parentValue)),每个子线程拿独立副本。
埋点落地在 rag/core/aop/RagTraceAspect,用 @RagTraceNode(name, type) 注解方法,AOP 自动 pushNode/popNode 记耗时。SSE 跨线程流还实现 RagStreamTraceSupport(内嵌 StreamSpan.detach/finishSuccess/finishError),保证流结束才关节点。
考点:ThreadLocal 在异步下的失效;TTL 的 copy() 语义;并发下可变集合共享引用导致的隐蔽 bug——这是本项目工程含量最高的细节,面试讲这个最容易出彩。
13.9 问题十:模型会挂、会限流,问答服务怎么"无感降级"?
工程问题
单个模型供应商 SLA 不可靠,直接暴露给用户会体验崩坏。
解法:路由 + 三态熔断 + 首包探测(组合模式)
RoutingLLMService(@Primary 注入,业务无感)维护候选模型列表,按档位/健康度路由:
ModelHealthStore三态(健康/半开/熔断),按错误率自动切换;- 发出请求后等首包 token,超时/异常立即降级到下一个候选;
- 所有厂商复用
AbstractOpenAIStyleChatClient(见 13.3)。
考点:这本质是微服务容错的思路(Resilience4j 那一套),只是触发条件从 HTTP 状态码变成"首包延迟/空响应"。讲清楚"为什么不用普通 try-catch 包一层"——因为要"按健康度动态选路 + 跨请求共享熔断状态 + 首包粒度探测",需要专门的路由层。
13.10 问题十一:高并发 SSE 怎么保护模型、又怎么防用户重复提交?
工程问题
模型是稀缺资源,突发流量会打爆;同时用户手抖连点会产生重复流水线。
解法:限流(保护模型)+ 幂等(保护用户),两个正交维度
- 限流:
chatQueueLimiter.enqueue(...)用FairDistributedRateLimiter——Redis ZSET 公平排队 + Lua 原子抢占 + 租约 + Pub/Sub 唤醒;reject 时仍落REJECTED状态消息并发meta→reject→finish→done,保证前端永远收到 DONE(不会无限转圈)。 - 幂等:
RAGChatController.chat()上@IdempotentSubmit(key=userId),同用户并发发问直接拦下并友好提示。
考点:限流(全局并发数,保护下游)与幂等(同用户去重,保护体验)职责不同、必须分开设计。可追问"不同会话要不要并行"——答案是当前按 userId 拦,若要"同会话串行、不同会话并行"应把 key 换成 conversationId,这是个真实的设计取舍。
13.11 问题十二:多模块之间怎么做到"换模型/换向量库不动业务"?
工程问题
业务编排(bootstrap)如果直接 new MilvusVectorStoreService() 或硬编码某模型,未来替换成本极高。
解法:依赖倒置(DIP)+ 接口隔离 + 配置驱动(架构级模式)
四模块依赖单向:bootstrap → framework + infra-ai,且 bootstrap 只依赖 LLMService/VectorStoreService/SearchChannel 等接口。VectorStoreService 接口定义了 indexDocumentChunks/updateChunk/deleteDocumentVectors/...,MilvusVectorStoreService 与 PgVectorStoreService 是两套实现,由配置决定注入哪个。mcp-server 更彻底——编译期零依赖,只通过 MCP 协议通信,可独立部署、扩缩容。
考点:这是贯穿全项目的"总纲"模式。所有前述模式(策略/装饰器/注册表/适配器)都是 DIP 在局部的具体落地——把"易变的外部依赖"隔离到独立模块/实现类,靠接口和配置解耦。
13.12 模式总览表(按"问题→解法"索引)
| 工程问题 | 设计模式/手段 | 代码落点 |
|---|---|---|
| 检索策略易变、需增减 | 策略 + SPI 自发现 | retrieval/channel/SearchChannel |
| 并行骨架重复 | 模板方法 | vector/strategy/AbstractParallelRetriever |
| 多模型厂商接入 | 适配器 + 协议标准化 | infra-ai/.../AbstractOpenAIStyleChatClient |
| 写入时附加同步(图谱/关键词) | 装饰器 | vector/decorator/*SyncingVectorStoreService |
| 后处理步骤可调序 | 责任链 | retrieval/postprocessor/* |
| 运行时按 ID 查找组件 | 注册表 | mcp/McpToolRegistry、intent/IntentNodeRegistry |
| 复杂对象创建 | 工厂 + 构建者 | IntentTreeFactory、callbackFactory、StreamChatContext.builder() |
| 跨线程链路上下文丢失 | TTL + 深拷贝修正 | framework/trace/RagTraceContext |
| 模型不稳定 | 路由 + 熔断 + 首包探测 | infra-ai/.../RoutingLLMService |
| 高并发打爆模型 | 分布式限流(公平排队) | rag/service/ratelimit/* |
| 用户重复提交 | 幂等 AOP | framework/idempotent + @IdempotentSubmit |
| 换模型/存储不动业务 | 依赖倒置 + 接口隔离 | 四模块依赖方向 + VectorStoreService |
一句话总结:Ragent 的优雅不在于用了多少"高大上"的模式,而在于每个模式都精准对应一个真实工程痛点——可变外部依赖用策略/装饰器/适配器隔离,重复骨架用模板方法收敛,运行时查找用注册表,异步上下文用 TTL 修正,稀缺资源用限流+熔断+幂等保护。读源码时带着"作者在这里遇到了什么问题"去想,比背模式定义收获大得多。
附录:必读文件清单(P0–P3)
| 优先级 | 文件 | 为什么 |
|---|---|---|
| P0 | README.md + assets/*.png |
全局认知 + 架构图 |
| P0 | framework/convention/ChatMessage.java |
领域统一语言 |
| P0 | framework/trace/RagTraceContext.java |
并发透传亮点(TTL 深拷贝) |
| P0 | bootstrap/.../RagentApplication.java |
启动入口与模块划分 |
| P0 | bootstrap/.../rag/controller/RAGChatController.java |
SSE 入口 + 幂等 |
| P0 | bootstrap/.../rag/service/impl/RAGChatServiceImpl.java |
全流程编排"目录页" |
| P1 | infra-ai/.../RoutingLLMService.java |
模型路由容错 |
| P1 | infra-ai/.../AbstractOpenAIStyleChatClient.java |
协议统一封装 |
| P1 | bootstrap/.../rag/core/retrieval/MultiChannelRetrievalEngine.java |
多通道并行检索 |
| P1 | bootstrap/.../rag/core/retrieval/postprocessor/FusionPostProcessor.java |
RRF 融合算法 |
| P2 | bootstrap/.../rag/core/intent/IntentTreeFactory.java |
意图树 |
| P2 | bootstrap/.../rag/core/memory/DefaultConversationMemoryService.java |
会话记忆 |
| P2 | bootstrap/.../rag/core/mcp/McpClientAutoConfiguration.java |
MCP 集成 |
| P2 | bootstrap/.../ingestion/engine/IngestionEngine.java |
入库 DAG 流水线 |
| P2 | frontend/src/stores/chatStore.ts |
SSE 事件消费 + 状态机 |
| P3 | mcp-server/.../McpServerApplication.java |
独立工具服务 |
| P3 | frontend/src/pages/ChatPage.tsx |
问答 UI 主页面 |
读 P0+P1 即可建立完整认知;P2 深入核心算法;P3 扩展到工具服务与前端。配合官方在线文档
nageoffer.com/ragent与本地断点调试,学习效果最佳。