--- title: "00-聊天系统整体架构讲解" created: 2026-05-21 aliases: - 聊天系统整体架构讲解 tags: - 项目 --- # 聊天系统整体架构讲解 我们从用户按下"发送"那一刻开始,把整个聊天系统的链路完整走一遍。整条链路分 12 步,每一步都有明确的代码位置和设计意图,跟着这个顺序看下去,整个系统的骨架就清晰了。 ## **整体设计哲学** 先说大方向。这套聊天系统采用的是**分层响应式架构**,底层基于 Spring WebFlux 的 Flux/Mono 实现 SSE 流式推送。核心思路有三个,先立起来再讲细节: **第一,前端建立 SSE 长连接之后,后端才真正开始干活**。这个"才真正开始"不是随口说的,而是通过 `Flux.defer` 在代码层面严格保证的——副作用绑定订阅事件,这是响应式编程的核心思想。 **第二,不是所有问题都该交给 Agent 自由发挥**。知识问答追求"稳"和"可解释",让 Agent 自己探索反而不可控;而"今天天气怎么样"这种问题,知识库里根本没答案,必须让 Agent 调用搜索工具。所以系统设计了 5 种执行器,按场景精准分流——优先用最稳定的方式回答,只有确实需要自由探索时才启动 Agent。这是**策略模式 + 自动发现**在业务上的落地。 **第三,任何一步失败都不能留下悬空状态**。整条链路有 Redis 租约、数据库轮次、内存运行态注册表三种"占用",任何一种残留都会让后续请求受阻或数据混乱。所以从抢租约到 cleanup 释放,每一步都有对应的失败补偿——多层 try-catch-finally 不是为了好看,是为了保证每个失败点都有合适的兜底。 整条链路分成三大阶段:**准备阶段**(步骤 1-6)、**生成阶段**(步骤 7-10)、**收尾阶段**(步骤 11-12)。下面逐步展开。 --- ## **步骤 1:接收 HTTP 请求,校验参数** **代码位置**:`BusinessChatController.stream()` ```java @PostMapping(value = "/stream", produces = "text/event-stream;charset=UTF-8") public Flux stream(@Valid @RequestBody ChatRequestDto dto) { return businessChatService.openConversationStream(dto); } ``` Controller 就做两件事:接参数、转交 Service。但这一个方法里藏着四个值得讲的设计点: **返回值是** `Flux` **而不是普通 JSON**。聊天回答是流式生成的——模型每产出一小段文字就立刻推给前端,用户看到"边想边写"的效果,感知延迟从 30 秒降到 1-2 秒。如果写成 `ResponseEntity>` 同步等结果,模型生成 30 秒前端就以为系统死了。 `produces = "text/event-stream;charset=UTF-8"` **声明的是 SSE 协议**。SSE 比 WebSocket 简单——单向通信就够,不需要双向;比 long polling 高效——不需要反复重连。聊天这种"服务器推、客户端只接"的场景,SSE 是最匹配的协议。 `@Valid` **是最外层防御**。ChatRequestDto 里 question 和 chatMode 都标了 `@NotBlank`,question 为空连方法体都进不去直接 400 拒绝。这一层挡住的非法请求**不消耗任何业务资源**——不创建轮次、不抢租约、不调模型——这是对系统最大的保护。 **不再额外包装 ApiResponse**。正常 REST 都返 `{"code":0,"data":...}` 这种壳,这里不包装。原因是 SSE 不是单次响应而是事件流——事件结构由 StreamEventWriter 统一定义(type + content + timestamp + metadata),再套一层 ApiResponse 反而打破了事件流的标准结构。错误处理走 SSE 的 error 事件,不走 HTTP 状态码。**流式接口的协议规范由事件流自己定义,不沿用 REST 包装**——这是流式 API 的设计哲学。 ## **步骤 2:Flux.defer 延迟启动** **代码位置**:`BusinessChatService.openConversationStream()` ```java public Flux openConversationStream(ChatRequestDto request) { return Flux.defer(() -> openDeferredConversationStream(request)); } ``` 这一行看起来很简单,但设计意图非常关键。**Flux.defer 的语义是把内部逻辑的执行从"Flux 对象创建时"推迟到"订阅发生时"**——类比 JavaScript 的 Promise,`new Promise(executor)` 是 executor 立刻执行,而 `Flux.defer(supplier)` 是 supplier 在订阅时才执行,这是 Flux 和 Promise 一个非常关键的区别。 **为什么聊天场景必须 defer?因为这条链路有强副作用**。普通业务接口没这个问题——`SELECT * FROM user WHERE id=1` 即使客户端没收到结果也没副作用。但聊天链路要抢 Redis 锁、要创建数据库轮次、要注册到运行态——这些都是**改变系统状态**的操作。如果不用 defer,Flux 对象一创建副作用就跑起来了,但这时候 WebFlux 框架内部还要走一系列处理(写 header、绑定 sink),这期间客户端可能因为网络抖、客户端崩、超时无法订阅——结果就是租约已占、轮次已建但客户端没订阅,**资源直接泄漏**。 defer 让所有副作用**绑定到订阅事件**——订阅成功就执行副作用、订阅失败副作用根本没发生、系统状态干净。说白了,defer 把"创建管道"和"点火执行"拆开了。 ## **步骤 3:buildLaunchPlan 规范化参数** **代码位置**:`BusinessChatService.buildLaunchPlan()` 外部参数进来之后不能直接用,需要规范化成内部的不可变快照 `StreamLaunchPlan`。这一步本质上是**反腐层**——把外部不可控参数转成内部可信赖结构。具体做四件事: **normalizeQuestion**:trim 空白、校验非空,确保下游不处理空白问题。 **normalizeConversationId**:前端传了就 trim 后用、没传就生成 UUID。这个设计让"传与不传"有明确语义——传了是继续会话、不传是开启新会话。UUID 还做了 `replace("-", "")` 把短横线去掉,从 36 字符缩到 32 字符——更短省存储、URL 安全无需转义、作为 Redis key 不需要处理特殊字符。轻量惯例,贯穿全系统会让代码更整齐。 **parseRequiredChatMode**:把字符串转成 ChatQueryMode 枚举。这里有几个细节值得讲——把 Optional 和 Required 拆成两个方法(语义放方法名上比放参数里清晰)、`"ALL"` 等同空值(用户友好的别名归一化成 null)、`toUpperCase` 大小写不敏感(前端容错)、IllegalArgumentException 包装成可读错误信息("chatMode 非法: xxx" 比 "No enum constant ChatQueryMode.xxx" 友好得多)。 **resolveSelectedDocument**:根据聊天模式做模式相关校验: | 聊天模式 | selectedDocumentId 规则 | | --- | --- | | OPEN\_CHAT | 不允许传,传了报错 | | AUTO\_DOCUMENT | 不允许传,文档由系统自动路由 | | DOCUMENT | 必须传,且必须当前可检索 | 这里 DOCUMENT 模式的"当前可检索"校验特别值得讲——它不是 `documentMapper.selectById`,而是 `documentKnowledgeService.listRetrievableDocuments().stream().filter(...)`。**selectById 看的是"文档是否存在",listRetrievableDocuments 看的是"文档现在是否可被检索"**——文档可能存在但 indexStatus 是 BUILD\_FAILED、可能被用户下线、可能索引在重建中。把"业务可用"语义封装在一处,避免每次都重新写过滤条件。 OPEN\_CHAT 和 AUTO\_DOCUMENT 都不允许传——**严格拒绝比"忽略不处理"更好**。忽略的话用户以为传的有效、实际系统当没看见、造成困惑;拒绝的话用户立刻知道用错了。**让错误尽早暴露**是 API 设计的常用原则。 最后还会注入两个时间字段——`currentDate` 和 `currentDateText`,作为后续处理"今天/最新/本周"等相对时间表达的统一基准。这里特意用 `LocalDate.now(CHAT_ZONE_ID)` 显式指定 Asia/Shanghai 时区——不用默认时区是因为 JVM 默认时区受部署环境影响(Docker、操作系统设置),不同节点的"今天"可能不一样,**显式指定让全局一致**。 整个 buildLaunchPlan 的产出是一个不可变快照,所有字段都 final——后续所有环节都基于这个快照工作,外部参数从此和内部状态彻底解耦。在入口处就把非法参数全拦住了,后续的编排器和执行器就可以放心使用,不用再做重复校验。 ## **步骤 4:Redis 分布式租约抢占** **代码位置**:`BusinessChatService.claimConversationLease()` ```java private boolean claimConversationLease(StreamLaunchPlan launchPlan) { return redisLeaseManager.acquire( launchPlan.getLeaseKey(), // chat:running:{conversationId} launchPlan.getLeaseOwnerToken(), // 本次请求的唯一 token (UUID) CHAT_RUNNING_LEASE_TTL // 30 秒 ); } ``` 这一步保证**同一个会话在集群任意时刻只有一个生成任务在运行**。集群部署下用户快速点两次发送或前端重试,没有租约的话两个生成任务可能在不同节点上同时跑——同一个 conversationId 下两条轮次记录、两份答案文本、两份引用,前端状态机彻底乱套。 **三个参数的设计都有讲究**: - **leaseKey** 是 `chat:running:{conversationId}`,按会话粒度加锁不同会话彼此不影响。 - **leaseOwnerToken** 是 UUID,本次请求对租约的"身份证"。这个 token 防的是"误释放别人的锁"——考虑这个场景:节点 A 抢到锁(token=T1),任务卡住超过 30 秒,锁过期,节点 B 抢到锁(token=T2),节点 A 恢复后执行 release——**无 token 验证就释放了 B 的锁**,有 token 验证发现不匹配拒绝释放。这是 Redis 分布式锁的标准模式(Redlock 思想)。 - **TTL 30 秒**防止节点宕机后锁永久残留,但 30 秒不够支撑模型生成所以要续期。 **TTL 和续期间隔的比例 30:10 = 3:1 是经验法则**——允许最多 2 次续期失败仍不至于让锁过期,给 Redis 抖动、网络延迟留余量。为什么不直接 TTL 设 1 小时省掉续期?节点宕机时锁要等 1 小时才释放,用户 1 小时都没法重发。为什么续期间隔不是 25 秒(接近 TTL)?任何延迟(GC、网络抖)都可能错过续期,5 秒安全余量太小。 **续期失败立刻停止任务**——这是租约机制最有意思的设计。租约失效有两种原因:Redis 故障(没法判断别人是否抢了)、真有别人抢了(必须停止避免双写)。保守做法是不管哪种原因都停——代价是用户可能要重发,收益是避免数据混乱。**互斥一旦不确定就要回到安全态**。 实现细节上还有个值得注意的地方:续期失败时**先 dispose 定时器自身再调 stopTask**——不先停定时器的话 10 秒后又触发 renewLeaseOrStop、又检测续期失败、又调 stopTask、循环触发、日志大量重复。**先切断重复触发的源头再执行真正的停止逻辑**,是个工程细节。 抢占失败直接返回拒绝流"该会话当前正在执行中,请稍后再试",前端可以显示一个友好提示。 ## **步骤 5:Bootstrap——创建轮次记录 + 构建 TaskInfo + 注册运行态** **代码位置**:`conversationArchiveStore.startExchange()` + `createTaskInfo()` + `chatRuntimeRegistry.register()` Bootstrap 做三件强关联的事,必须**顺序执行且全部成功**: ```text 1. startExchange:数据库创建轮次记录 → 拿到 exchangeId 2. createTaskInfo:内存组装运行时上下文 → 用到 exchangeId 3. register:把 taskInfo 注册到运行态 → 用 conversationId 做 key ``` 先在归档层创建一条新的 exchange 轮次(后续无论成功失败都依赖它收尾落库),然后组装 TaskInfo 运行时上下文,最后注册到运行态注册表。 ### **TaskInfo:运行时的"万能上下文"** TaskInfo 是整条执行链路的核心数据载体,几乎所有组件都要从它身上拿东西。按用途分组它装了这些: - **身份标识(不可变)**:conversationId、exchangeId、traceId、leaseKey、leaseOwnerToken - **请求参数(不可变)**:question、chatMode、selectedDocumentId、currentDate - **状态机**:executionPlan(后期填充)、debugTrace、finalized(AtomicBoolean)、firstResponseTimeMs(AtomicLong) - **运行态集合(线程安全)**:sink、answerBuffer、thinkingSteps、references、usedTools - **基础设施**:runnableConfig、traceRecorder、eventMetadata - **资源句柄**:disposable(执行流订阅)、leaseRenewalDisposable(续期任务) **线程安全机制按需选择,不是一刀切**——这一点特别值得点出来。answerBuffer 用 StringBuffer(append 频繁,内部 synchronized 性能足够);thinkingSteps、references 用 `Collections.synchronizedList`(add 不频繁,简单互斥就行);usedTools 用 `ConcurrentHashMap.newKeySet()`(并发 add 高,需要无锁);finalized、firstResponseTimeMs 用 AtomicBoolean / AtomicLong(单变量原子操作走 CAS)。**不同字段不同诉求,匹配合适的并发原语**。 **finalized 通过 CAS 保证收尾只执行一次**这一点特别关键。任务可能正常完成、可能被用户停止、可能因为续期失败被停止、可能因为执行异常失败——三种收尾路径都会触发,必须保证只有一种生效。CAS 的 `compareAndSet(false, true)` 谁先抢到谁有资格执行收尾,后到的看到 false→true 已经发生了直接 return。 **RunnableConfig.context() 是个共享 Map,所有组件通过它读写运行态**。把 sink、metadata、thinkingSteps、references、usedTools、traceId、question、chatMode、currentDate、selectedDocumentId 这些都放进去,深层组件(工具、检索器、追踪器)能直接通过 context 取到自己需要的东西,**不需要在每个方法签名里传一堆参数**。这是规模化系统里减少参数爆炸的标准手法。 还有一个工程细节是 **debugTrace 预占位**——在 executionPlan 还没生成时先放一个空白 debugTrace 进 context。原因是 prepareExecutionPlan 过程中某些组件可能去 context 取 debugTrace,如果没放就是 null 就 NPE。先放空容器、组件取出来不会 NPE 只是写入空容器、prepareExecutionPlan 完成后再替换为真实的——**Null Object Pattern 的经典应用**。 ### **双重防御:为什么抢到租约还可能注册失败** ```java if (!chatRuntimeRegistry.register(taskInfo)) { failBootstrappedExchange(...); releaseLeaseQuietly(...); return BootstrapResult.rejected("该会话当前正在执行中,请稍后再试"); } ``` 新手看到这里会问——租约都抢到了,注册怎么还会失败?设想这个场景:节点 A 服务请求 1(cid=X)抢到 Redis 锁任务执行中,节点 B 服务请求 2(cid=X),Redis 主从延迟时——A 占的锁还没同步到从库、B 访问从库看到锁不存在也"抢到"、两个节点都以为自己拿到了锁。 这时候本地的 chatRuntimeRegistry 就站出来做最后防御。整个互斥体系是**双重防御**: - **Redis 租约**:跨节点互斥(应对正常情况) - **本地注册表**:节点内互斥(应对 Redis 短暂不一致) 注册失败时的补偿性收尾**顺序很关键**——先 `failBootstrappedExchange` 标记轮次失败再 `releaseLeaseQuietly` 释放租约。先释放租约的话,锁刚释放后续请求进来、看到失败的轮次也开始建新轮次、两个轮次叠加污染数据。**先标记失败确保数据库状态收敛,再释放锁,下一次进来面对的是干净状态**。 外面还有 try-catch 兜底——bootstrap 三步任意一步抛异常都进 catch 释放租约 + 标记失败。注意 exchangeView 是在 try 外声明的局部变量,catch 块才能访问到它判断是否要清理——**try-catch 跨作用域的常见 Java 技巧**。 ## **步骤 6:绑定 SSE 通道** **代码位置**:`BusinessChatService.bindClientChannel()` ```java private Flux bindClientChannel(TaskInfo taskInfo) { return taskInfo.sink().asFlux() .doOnSubscribe(ignored -> activateGeneration(taskInfo)) .doOnCancel(() -> stopTask(taskInfo, "客户端已取消请求")); } ``` 这一步把 TaskInfo 的内部 sink 转成前端可消费的 Flux。 `doOnSubscribe` **是整条链路的"点火开关"**——前端建立 SSE 连接那一刻,后端才触发 `activateGeneration` 开始生成。这是 Flux.defer 设计意图的最终落点——所有副作用都绑在订阅事件上。 `doOnCancel` **是断连兜底**——前端主动关闭连接(用户关掉浏览器、网络断了)时,后端立即停止当前任务,避免模型空跑浪费 token。 这里 sink 用的是 `Sinks.many().unicast().onBackpressureBuffer()`——**unicast 单播**保证一条流只服务当前订阅的前端连接(聊天场景一个会话只有一个观察者,不需要 multicast);**onBackpressureBuffer** 处理背压——前端消费慢时事件先缓存到 buffer 里,避免快速产 chunk 撑爆下游。 三者配合形成**生命周期完整的响应式管道**——defer 控启动时机、subscribe 触发副作用、cancel 兜底中断。 ## **步骤 7:激活生成——启动租约续期 + 执行流水线** **代码位置**:`BusinessChatService.activateGeneration()` 激活后做两件事:第一启动租约续期定时器(每 10 秒一次,续期失败自动停止);第二启动执行链路 `buildConversationExecution` 的订阅。两个 Disposable 都存到 TaskInfo 上,停止时统一 dispose。 激活的瞬间还会检查 `finalized` 标志——如果任务在刚启动后立即被其他线程结束(极端情况下停止请求和启动请求几乎同时到达),主动 dispose 避免重复执行。 执行链路的核心组装是 `buildConversationExecution`: ```java return Flux.defer(() -> { safeEmit(taskInfo.sink(), streamEventWriter.thinking("正在分析问题上下文。", ...)); return Mono.fromCallable(() -> prepareExecutionPlan(taskInfo)) .subscribeOn(Schedulers.boundedElastic()) .flatMapMany(plan -> { ConversationExecutor executor = conversationExecutorRegistry.get(plan.getMode()); return executor.execute(taskInfo); }); }) .publishOn(Schedulers.boundedElastic()) .doOnNext(chunk -> emitModelChunk(taskInfo, chunk)) .doOnError(error -> finishWithFailure(taskInfo, error)) .doOnComplete(() -> finishSuccessfully(taskInfo)); ``` 这段代码非常密,每一行都有讲究: **嵌套 Flux.defer**——外层 bindClientChannel 已经在 defer 里了为什么这里又来一个?外层是订阅时才点火,内层让"每次重新订阅都重新执行 thinking 事件 + prepareExecutionPlan"。理论上单次会话不会重订阅但 defer 把这种**重入安全性**当成基础保证,比依赖"不会重订阅"更稳。 `Mono.fromCallable` **+** `subscribeOn(boundedElastic)`——prepareExecutionPlan 是同步阻塞方法(数据库查询、检索、压缩),直接在响应式流里调会阻塞当前线程。Mono.fromCallable 把它包装成响应式异步任务,subscribeOn 指定它跑在 boundedElastic 调度器上。**boundedElastic 是 Reactor 专门为阻塞 IO 准备的弹性线程池**——线程数有上限不会爆、空闲线程会回收、适合数据库查询这种 IO 密集型任务。Schedulers.parallel 是 CPU 密集型不适合 IO;Schedulers.single 串行单线程并发量上不去——都不合适。 `flatMapMany` **不是** `map`——map 是 T → R 同步一对一,flatMapMany 是 T → `Flux` 异步一对多。prepareExecutionPlan 产出一个 plan,executor.execute 产出**多个 chunk**,是一对多的关系。如果用 map 会拿到 `Flux>` 嵌套 Flux 还要再 flatMap 解嵌套,不如 flatMapMany 直接搞定。 `publishOn(boundedElastic)` **切换下游线程**——模型每产生一个 chunk doOnNext 会执行 emitModelChunk。如果 doOnNext 在模型推送的线程上跑(可能是 LLM SDK 内部线程),长时间占用 LLM 线程会影响模型自身吞吐。publishOn 把 chunk 处理切换到 boundedElastic 上、释放 LLM 线程、让模型可以继续生成下一个 chunk。 `doOnNext` **/** `doOnError` **/** `doOnComplete` **三态闭环**——这是响应式流的三种终态:持续中触发 doOnNext、异常触发 doOnError、正常完成触发 doOnComplete。三种路径都通向收尾——成功走 finishSuccessfully、失败走 finishWithFailure,**保证任何一种路径系统状态都收敛**。 ## **步骤 8:编排器决策** **代码位置**:`ChatPreparationOrchestrator.prepare()` 这一步是整个系统的"大脑"。编排器不只是"看用户传了啥跑啥",而是"综合各种信号判断到底应该跑啥"——这是它叫"大脑"的原因。按顺序做这些事: **第一,装载会话记忆**——`summarizeHistory` 拿到 ConversationMemoryContext,里面有长期摘要(之前 N 轮对话的 LLM 压缩版)和最近窗口(最近几轮原文)。压缩 + 原文的两段式记忆是为了**远的轮次保留语义但省 token、近的轮次保留细节给问题改写用**。 **第二,构建历史上下文**——`buildPlanningHistory` 给编排决策用、`buildAnswerHistoryContext` 给最终答案生成用。两个上下文用途不同所以分开构造——编排时关心"用户之前问过什么",回答时关心"我之前是怎么答的"。 **第三,判断时效性**——`TimeSensitiveQueryHelper` 看用户问的是不是"今天/最新/本周"这类需要实时信息的问题,产出两个 boolean:`requiresCurrentDateAnchoring`(需要日期锚定)和 `requiresFreshSearch`(需要联网搜索最新事实)。两个是分开的——"今天周几"只需日期锚定不需联网,"今天的新闻"两者都要。 **第四,问题改写**——把用户原始问题改写成检索友好的表达,可能拆成多个子问题。这一步对 RAG 召回质量影响巨大——用户口语化的提问"那个上次说的方案怎么样了"对检索引擎几乎没用,必须改写成包含关键实体的明确表达。 **第五,知识路由**(AUTO\_DOCUMENT 模式)——通过 `KnowledgeRouteService` 确定候选文档范围,避免"从全部文档里大海捞针"。 **第六,主题路由**——通过 `DocumentQuestionRouter` 判断走图查询(结构化导航类)还是混合检索(开放性问答)。 **第七,生成执行计划**——确定 `ExecutionMode`,输出 `ConversationExecutionPlan`。 ### **三级路由的精妙之处** 编排器的路由设计是三级的,每一级解决不同的问题: - **Scope 路由**:确定查哪些文档(AUTO\_DOCUMENT 模式下系统自动选) - **Topic 路由**:在文档内定位到主题/章节 - **Mode 决策**:选择具体的执行路径(RAG 检索 / 图查询 / Agent 自由探索 / 澄清确认) 编排器甚至可以**修正用户的文档选择**——AUTO\_DOCUMENT 模式下用户没选文档,编排器分析问题后决定该查哪个文档。修正发生时还要做善后: ```java if (executionPlan.getSelectedDocumentId() != null && !Objects.equals(executionPlan.getSelectedDocumentId(), taskInfo.selectedDocumentId())) { conversationArchiveStore.refreshSessionScope(...); // 同步归档 putContextIfNotNull(taskInfo.runnableConfig(), ChatContextKeys.SELECTED_DOCUMENT_ID, ...); // 同步上下文 ... } ``` 不光更新执行计划本身,还要同步刷新归档的会话范围 + RunnableConfig 上下文——**三个数据源(执行计划、归档、运行上下文)保持一致**,防止后续组件读到旧值。 ### **buildAgentQuestion:不要让 Agent 裸接用户问题** 最后 `buildAgentQuestion` 把时间锚点 + 时效约束 + 历史摘要 + 原始问题拼成最终 prompt,结构大致是: ```text 系统时间信息: 当前日期是 2026-05-21,时区为 Asia/Shanghai。 [条件化时效约束] [相关会话背景(如有)] 用户问题: [原始问题] ``` **为什么不能裸传问题?裸问题会出问题**——用户问"今天天气怎么样"不告诉 Agent 今天是几号,它可能用训练数据里的旧日期回答;用户问"还记得我刚才说的吗"不给历史,Agent 直接懵。 包装的核心是**条件化注入**——`requiresCurrentDateAnchoring` 时给更严格的日期约束、`requiresFreshSearch` 时强制要求联网搜索并标注来源日期、有 historySummary 才注入历史。**不同问题有不同的 prompt 强度,精确匹配**。 每一步编排决策都通过 `ConversationTraceRecorder` 记录成 stage——MEMORY、ROUTE、REWRITE 等,每个 stage 都有开始时间、结束时间、状态、附加信息。后续可以直接看到"哪一步耗时多少、走了什么分支、为什么选了这个执行模式"——**整条决策路径完全可观测**。 ## **步骤 9:执行器执行** **代码位置**:`ConversationExecutorRegistry.get(mode)` → `executor.execute()` 编排器确定执行模式后,`ConversationExecutorRegistry` 用 EnumMap 做 O(1) 查找找到对应执行器: ```java @Component public class ConversationExecutorRegistry { private final Map executorMap = new EnumMap<>(ExecutionMode.class); public ConversationExecutorRegistry(List executors) { for (ConversationExecutor executor : executors) { executorMap.put(executor.mode(), executor); } } ... } ``` 这一段代码非常短但有三个值得讲的设计点: **EnumMap 不是 HashMap**——枚举做 key 时 EnumMap 是数组实现、查找 O(1) 且常数极小、内存紧凑,是 Java 集合里的小优化。 **Spring 自动注入** `List` 是 Spring 的"魔法"——所有实现了 ConversationExecutor 接口、被 @Component 标注的 Bean,Spring 会自动收集成 List 注入。**新增执行器只需写新类 implements ConversationExecutor 然后加 @Component——注册表代码完全不用改**。这是开闭原则(对扩展开放、对修改关闭)的教科书式实现。 **接口里有** `mode()` **方法让执行器自己声明处理哪种模式**——构造函数遍历 List 时 `executorMap.put(executor.mode(), executor)` 自动建索引。**执行器自己知道自己是什么、不需要外部告诉它**——职责更内聚。 5 种执行器各司其职: | 执行模式 | 执行器 | 做什么 | | --- | --- | --- | | RETRIEVAL | RagChatExecutor | 多通道检索→RRF 融合→重排→父块提升→Prompt 组装→模型流式生成 | | REACT\_AGENT | ReactAgentExecutor | 自主推理引擎,支持工具调用和联网搜索 | | GRAPH\_ONLY | GraphOnlyExecutor | 纯图查询,适合结构化导航类问题 | | GRAPH\_THEN\_EVIDENCE | GraphThenEvidenceExecutor | 先图查询定位再补充证据 | | CLARIFICATION | ClarificationExecutor | 文档范围歧义时,向用户确认 | ### **RagChatExecutor 的检索链路** RAG 执行器的检索链路值得展开说一下,它不是简单地"查向量库返结果",而是一条**多阶段精排管道**: 1. **多通道并行召回**——**向量检索**(PGVector,语义匹配,召回语义相似但用词不同的内容)+ **关键词检索**(ES,BM25 精确匹配,召回包含具体术语的内容)双通道并行。两者互补:向量召回"意思相近的",关键词召回"用词精确的"。 2. **RRF 秩融合**——两路结果合并不是简单拼接,而是用 Reciprocal Rank Fusion 算法做秩融合。每个文档在每路里的排名转成分数 `1/(k + rank)`、加总后重排——**不依赖原始打分的绝对值,只看相对排名,跨通道融合更稳定**。 3. **Cross-Encoder 重排**——用专门的重排模型对融合后的 Top-N 做精排,提升相关性。Bi-Encoder(向量召回用的)速度快但精度一般,Cross-Encoder 速度慢但精度高——**召回阶段用 Bi-Encoder 拉宽,精排阶段用 Cross-Encoder 收敛**。 4. **父块提升**——这是 Parent-Child 切块策略的关键。**Child 小块进向量库精准召回**(小块语义聚焦匹配准),**Parent 大块保留完整上下文不进向量库**——召回 Child 后通过 parentId 反查 Parent 给模型回答。**保证答案既有精准度又有信息量**。 5. **Prompt 组装**——`RagPromptAssemblyService` 把检索到的上下文 + 用户问题 + 时效约束拼成最终 prompt。 6. **模型流式生成**——通过 `ObservedChatModelService` 调用模型,每个 chunk 实时返回。 ### **ReactAgentExecutor 的开放探索** REACT\_AGENT 走的是另一条路——基于 ReAct 框架(Reasoning + Acting)让 Agent 自主推理。Agent 拿到问题后会思考(Thought)→ 决定要不要调工具(Action)→ 看工具返回(Observation)→ 继续思考……循环到能给出答案为止。 **ChatCheckpointManager 在这里很关键**——多轮工具调用之间需要状态传递(前一步 Observation 影响下一步 Thought),Checkpoint 把每一步的 Agent 状态快照保存下来。 可调用工具目前主要是 `TavilySearchTool` 联网搜索,应对"今天天气""最新新闻"这类时效性问题。 ## **步骤 10:流式输出** **代码位置**:`BusinessChatService.emitModelChunk()` 模型每产出一个 chunk,`emitModelChunk` 同时做三件事: ```java private void emitModelChunk(TaskInfo taskInfo, String chunk) { taskInfo.answerBuffer().append(chunk); if (taskInfo.firstResponseTimeMs().get() == 0L) { taskInfo.firstResponseTimeMs() .compareAndSet(0L, System.currentTimeMillis() - taskInfo.startTime()); } safeEmit(taskInfo.sink(), streamEventWriter.text(chunk, taskInfo.eventMetadata())); } ``` **第一,append 到 answerBuffer**——chunk 是模型流式输出的小片段(几个 token),最终落库时拼成完整答案。answerBuffer 是 StringBuffer 线程安全。 **第二,记首包耗时(TTFB)**——`firstResponseTimeMs` 用 AtomicLong + CAS 只记录第一次。**为什么要 CAS 不直接 set?多个线程同时进 emitModelChunk 时(理论上不会但保险起见)两个 set 都执行后者覆盖前者,CAS 保证只有第一次的值生效**。TTFB 是流式接口的核心 SLA 指标——用户看到第一个字的延迟、直接决定体验,必须精确测量。 **第三,推到前端**——`streamEventWriter.text` 把 chunk 包装成标准 SSE 事件、`safeEmit` 推到 sink。 **safeEmit 是关键的容错**——它捕获推送异常只记 warn 不抛。原因是 emit 失败抛异常 → doOnError 触发 → finishWithFailure,但实际上模型可能还在正常生成、一次推送失败不应影响整个任务。answerBuffer 里仍有完整数据、收尾落库不受影响。**把"瞬时推送失败"和"任务真正失败"解耦**。 SSE 事件协议是统一的,每个事件携带 `type` + `content` + `timestamp` + `conversationId` + `exchangeId`。事件类型是 5 种: | 类型 | 用途 | 推送时机 | 频次 | | --- | --- | --- | --- | | text | 模型输出正文 | 每个 chunk | 极高 | | thinking | 思考提示 | 阶段切换 | 低 | | error | 错误信息 | 失败时 | 一次 | | reference | 引用来源 | 收尾时 | 一次 | | recommend | 推荐追问 | 收尾时 | 一次 | 事件结构标准化的好处是前端只需要写一套事件路由代码,新增事件类型完全不破坏已有协议——**前后端协议解耦**。 ## **步骤 11:成功收尾** **代码位置**:`BusinessChatService.finishSuccessfully()` 模型输出完毕,`doOnComplete` 触发成功收尾。这里的顺序**不能随意打乱**,按 8 步严格执行: 1. **CAS 抢收尾权**:`finalized.compareAndSet(false, true)`,保证只执行一次 2. **快照数据**:冻结最终答案、去重引用列表(snapshotReferenceList + deduplicateReferences) 3. **开启 finalize/recommendation 阶段追踪** 4. **生成或提取推荐追问** 5. **补发引用和推荐事件**:必须在关闭 SSE 流之前 6. **关闭 SSE 流**(safeComplete) 7. **落库 COMPLETED 状态**:把完整快照 + thinkingSteps + tools + debugTrace 都落 8. **异步刷会话摘要 + cleanup 释放资源** ### **几个关键设计点** **CAS 抢权放第一**——保证整个收尾只有一次。停止/失败/成功三种路径都会动 finalized 这个标志,谁先抢到谁执行后面的,后到的直接 return。 **推荐追问的两种来源**值得讲——澄清模式(CLARIFICATION)下推荐就是编排阶段已经产出的澄清选项、直接复用避免再调推荐服务生成语义不一致的;常规模式下基于"原问题 + 最终答案 + 最近历史轮次"调 RecommendationService 生成、让推荐贴近本轮回答结果。**让推荐生成的输入包含"最终答案"是关键**——只看问题生成的推荐是泛泛的,看了答案才能生成"针对刚才回答的下一步追问"。 **事件顺序:正文 → 引用 → 推荐 → 关闭流**——前端按这个顺序稳定渲染,正文流结束后再看到引用和推荐、避免引用插在正文中间让前端状态机变复杂。**用户先看到答案、再看到来源、最后看到下一步建议**——符合阅读直觉。 **补发事件必须在关闭流之前**——关闭后 sink 已 complete,再 emit 不生效,前端收不到引用和推荐。 ### **三层 try-catch-finally 嵌套** 这是这段代码最值得讲的工程细节: ```java try { // 补发引用 + 推荐 } catch (RuntimeException e) { log.warn(...); // 只记 warn 不改判会话失败 } finally { try { safeComplete(sink); // 关流 } catch (...) { log.warn(...); } try { // 落库 COMPLETED + 完成 finalize stage } catch (...) { log.error(...); // 提升为 error } finally { safeRefreshConversationSummary(...); cleanup(...); } } ``` 每一层的设计意图是: - **最外层 try**:补发引用 + 推荐 → 失败只记 warn 不改判会话失败(**正文已成功,补发失败不能推翻"本轮成功"这一主结果**) - **第一层 finally**:safeComplete 关流 → 即使补发失败也要明确告诉前端"本轮结束" - **第二层 try**:refreshDebugTraceRuntimeStats + completeExchange 落库 → 失败提升日志为 error(**回答成功生成但收尾落库失败、需要重点排查**) - **最内层 finally**:safeRefreshConversationSummary + cleanup → 无论前面所有步骤成功失败都执行、保证资源最终释放 **这种多层嵌套不是为了好看,是为了保证每个失败点都有合适的兜底——正文成功不被补发失败覆盖、补发失败不阻塞关流、关流失败不阻塞落库、落库失败不阻塞清理**。每一层都比下一层"更必须执行"。 ### **cleanup:彻底释放资源** ```java private void cleanup(TaskInfo taskInfo) { if (leaseRenewalDisposable != null && !leaseRenewalDisposable.isDisposed()) { leaseRenewalDisposable.dispose(); // 1. 停租约续期定时器 } if (disposable != null && !disposable.isDisposed()) { disposable.dispose(); // 2. 停业务执行流 } releaseLeaseQuietly(...); // 3. 释放 Redis 租约 chatRuntimeRegistry.remove(...); // 4. 移出本地注册表 } ``` 四个动作分别释放四种资源——续期定时器、业务执行流、Redis 租约、本地注册表。任何一步失败都用 quietly 版本只记日志不抛异常——**cleanup 是终态不应阻塞,租约即使没主动释放 30 秒后也会自动过期**。 cleanup 的本质是让系统状态回到"这个 conversationId 完全干净、可以接受新请求"的初始态。 ## **步骤 12:异常收尾** **代码位置**:`BusinessChatService.finishWithFailure()` 同样通过 CAS 保证单次执行。流程: 1. CAS 抢收尾权(与 finishSuccessfully 互斥) 2. 提取异常信息(从异常链中找 `WebClientResponseException` 的响应体,给用户尽可能有用的提示而不是栈) 3. 发送 error 事件给前端 4. 关闭 SSE 流 5. 以 FAILED 状态落库(已生成的部分答案也要保留,方便后续分析) 6. finally 层 cleanup 释放资源 **异常信息提取的工程细节**——直接 `e.getMessage()` 拿到的可能是 `"500 INTERNAL_SERVER_ERROR"` 这种没用的,但 WebClientResponseException 的响应体里通常有模型 API 返回的具体错误(比如"超出 token 限制""API key 无效")。从异常链里挖到这个响应体、提取关键信息给前端,**给用户的提示就有指导意义而不是天书**。 注意一个原则:**落库失败不影响前端已收到的答案**——用户已经看到了回复,只是后台归档可能不完整,这属于可接受的降级。**用户体验优先**。 --- ## **降级策略一览** 系统在多个环节都做了防御性设计,遇到问题时优雅降级而不是直接崩: - **检索无结果**:返回预设的 `noEvidenceReply`,不会给用户一个空白回复 - **子问题超时**:忽略该子问题继续处理其他子问题,**部分结果总比完全无结果好** - **续期失败**:主动停止会话释放所有资源,防止幽灵任务 - **落库失败**:日志 error 但不影响前端已收到的答案,用户体验优先 - **补发引用/推荐失败**:只记 warn 不改判会话失败,不推翻"正文已成功"这一主结果 - **emit 推送失败**:safeEmit 吞异常只记 warn,不让单次推送失败影响整体 - **编排器异常**:直接 throw,由外层 catch 统一走失败收尾,保证租约释放 **整体哲学是"主结果优先"**——回答生成成功是核心目标,附加信息(引用、推荐、调试追踪、归档落库)失败都不能推翻这个主结果。 ## **状态转换** **轮次状态**:`CREATED` → `COMPLETED` / `FAILED` / `STOPPED`,一旦进入终态不可逆转。 **租约状态**:`ACQUIRED` → `RENEWING` → `RELEASED` / `EXPIRED`。正常流程是 ACQUIRED → RENEWING → RELEASED(cleanup 时主动释放);异常流程是续期失败后 EXPIRED 触发自动停止;极端情况是节点宕机后 30 秒 TTL 自动 EXPIRED。 **TaskInfo.finalized**:`false` → `true`(CAS 单向跳变),保证三种收尾路径只有一种生效。 ## **会话记忆体系** 最后提一下支撑整条链路的记忆体系,它分三层,每层应对不同的诉求: - **长期摘要**:通过 LLM 压缩历史对话生成的摘要,注入 AgentQuestion 前缀。**作用是"远记忆"——让 Agent 知道之前聊了什么主题但不占太多 token**。 - **最近对话窗口**:保留最近几轮的原文。**作用是"近记忆"——保证短期上下文准确性,问题改写时参照最近表达**。 - **Checkpoint**:ReAct Agent 的状态快照,支持多轮工具调用间的状态传递。**作用是"过程记忆"——单次 Agent 执行内的中间状态**。 为什么要分三层?长期摘要太粗(细节丢了)、最近窗口太窄(远的看不到)、Checkpoint 只在单次 Agent 内有效(跨轮次没用)——**任何单一记忆机制都不够,组合起来才覆盖全场景**。 每一轮收尾时 `safeRefreshConversationSummary` 异步刷新长期摘要——把刚结束的轮次纳入 --- **企业级项目导航**:⬅️ [[05-龙虾OpenClaw深度解析:从对话到执行的跨越|05-龙虾OpenClaw深度解析:从对话到执行的跨越]] | 00-聊天系统整体架构讲解 | ➡️ [[01-文档上传到RAG检索完整链路讲解|01-文档上传到RAG检索完整链路讲解]]