--- title: "04-白话讲解" created: 2026-05-21 aliases: - 白话讲解 tags: - 项目 --- # 白话讲解 ## **一、把这条链路放在系统里的位置** 讲这条链路之前先把它和上一段衔接一下。前面我们花了好几大段讲"入料阶段"——文档怎么从原始文件变成可被检索的结构化资产:上传、解析、切块、向量化、索引落库。**入料阶段的终点是 PGVector 和 ES 里躺着一份份准备好被检索的 chunk**。 **这一段是"使用阶段"**——用户终于来了,对着输入框敲了一句话点了发送。从那一刻开始,到他屏幕上一个字一个字蹦出答案,再到答案完整收尾、引用展示、推荐追问出现——这中间发生的所有事,是这条链路的全部内容。 讲这条链路之前先立一个核心认知:**这是一条响应式 + 流式的链路,不是传统的 REST 请求**。所有设计都是围绕"流"展开的——返回值是 Flux 不是 JSON、副作用要绑定到订阅事件、收尾要保证只触发一次。如果按传统 REST 接口的思维去看,会处处觉得别扭——但只要把"流"这个核心抓住,后面所有奇怪的设计都能讲通。 整条链路可以切成三大块: ```text 准备阶段:Controller → defer → 启动计划 → 抢租约 → bootstrap → 绑通道 生成阶段:激活 → 编排 → 选执行器 → 模型输出 → 推送 收尾阶段:CAS 抢权 → 补发事件 → 关闭流 → 落库 → 清理 ``` ## **二、模块划分:先把房子的承重墙看清楚** 讲具体流程之前先看一眼后端的模块结构,整条链路才有"它走在哪条管子里"的空间感。 后端是个多模块 Maven 工程,关键四个:**super-agent-common** 装通用基础设施(枚举、工具类、统一响应封装);**super-agent-id-generator-framework** 管分布式 ID 生成;**super-agent-redisson-framework** 装 Redis 能力(分布式锁、会话租约管理 RedisLeaseManager);**super-agent-business-chat** 是聊天系统的全部业务代码——本次讲的所有东西都在这个模块里。 聊天业务模块内部按职责分包,一句话一个: - **controller** 只放 BusinessChatController,**唯一职责是接 HTTP 请求然后把球踢给 Service**,不做任何业务。 - **dto / vo** 装请求和响应的数据传输对象——dto 是前端送进来的,vo 是返给前端的。 - **service** 是核心,里面装着 BusinessChatService(总控)、TaskInfo(运行态上下文)、StreamLaunchPlan(启动计划)、ChatRuntimeRegistry(运行态注册表)、ConversationArchiveStore(归档存储)、ConversationMemoryService(会话记忆)、ChatCheckpointManager 等十几个组件。 - **rag** 是 RAG 链路的核心,下面再分四个子包:**rag/executor** 装 5 种执行器(RagChatExecutor、ReactAgentExecutor、GraphOnlyExecutor、GraphThenEvidenceExecutor、ClarificationExecutor)、**rag/service** 装编排器和检索引擎、**rag/retrieve/channel** 装多路检索通道(向量、关键词等)、**rag/model** 装执行计划和检索上下文这些数据模型。 - **support** 装 SSE 事件格式化(StreamEventWriter)、Sink 安全发送(SinkEmitHelper)、时效性判断(TimeSensitiveQueryHelper)、上下文键常量(ChatContextKeys)这些支撑类。 - **tool** 装 Agent 可调用的外部工具,比如 TavilySearchTool 联网搜索。 - **manage** 装文档管理子系统(虽然在同一个 Maven 模块里但逻辑独立)。 整体的分层架构是这样的——**接入层(controller)→ 编排层(service)→ 决策层(rag/service)→ 执行层(rag/executor)→ 检索层(rag/retrieve)→ 支撑层(support)→ 持久层(data/mapper)→ 基础设施层(common/redisson/id-generator)**。每一层只依赖它下面的层,不反向依赖。 这个分层有一个核心原则:**BusinessChatService 是"轻总控",rag 子包是"厚组件"**。BusinessChatService 依赖了 14 个组件(Agent、checkpoint、归档、运行态、记忆、编排器、执行器注册表、SSE writer、推荐、租约、追踪、检索观察、阶段基准、文档知识),但它自己只负责把这些组件串起来——具体的活儿全在被它依赖的组件里。**总控薄、组件厚——薄的总控好读、厚的组件可独立测试和替换**,这是一个非常工程化的取舍。 ## **三、入口:Controller 为什么返回 Flux** 入口在 BusinessChatController.stream(): ```java @PostMapping(value = "/stream", produces = "text/event-stream;charset=UTF-8") public Flux stream(@Valid @RequestBody ChatRequestDto dto) { return businessChatService.openConversationStream(dto); } ``` 这一个方法就藏了四个值得单独讲的设计点。 **第一,返回 Flux 不返回 JSON**。传统 REST 接口里这里会写 `ResponseEntity>`——同步等结果、统一包装。但聊天场景这么写直接走不通——一个完整回答可能要 30 秒,同步等 30 秒前端用户以为系统死了。流式输出的本质是把"等结果"变成"持续推送"——模型每生成一个字就推一次,**总耗时还是 30 秒,但感知延迟降到 1-2 秒**。 **第二,produces 用 text/event-stream 这是 SSE 协议的 MIME 类型**。SSE 比 WebSocket 简单——单向通信就够,不需要双向;比 long polling 高效——不需要反复重连。聊天这种"服务器推、客户端只接"的场景,SSE 是最匹配的协议。 **第三,@Valid 在最外层挡掉非法请求**。ChatRequestDto 里 question 和 chatMode 都标了 @NotBlank,请求进不到方法体就被 400 拒了。**非法请求不消耗任何业务资源——不创建轮次、不抢租约、不调模型——这是对系统最大的保护**。 **第四,不再额外包装 ApiResponse**。正常 REST 都返 `{"code":0,"data":...,"msg":""}` 这种壳,这里不包装。原因是 SSE 不是单次响应而是事件流——事件结构由 StreamEventWriter 统一定义(type + content + timestamp + metadata),再套一层 ApiResponse 反而打破了事件流的标准结构。错误处理走 SSE 的 error 事件,不走 HTTP 状态码。**流式接口的协议规范由事件流自己定义,不沿用 REST 包装**——这是一个值得单独点出来的设计哲学。 --- ## **四、Flux.defer:延迟启动的真实意图** Controller 调到 BusinessChatService.openConversationStream(),第一行就是个 defer: ```java public Flux openConversationStream(ChatRequestDto request) { return Flux.defer(() -> openDeferredConversationStream(request)); } ``` 这一层 defer 包装非常关键,但它的真实意图新手很容易看不出来。 **先讲没有 defer 会怎样**。直接 return openDeferredConversationStream(request) 的话,这个方法内部要做的事——buildLaunchPlan、claimConversationLease、bootstrapConversation、bindClientChannel——会在 **Controller return 之前** 全部执行。问题来了: ```text Controller 把 Flux 对象 return 给 Spring WebFlux WebFlux 框架内部还要走一系列处理:写 header、绑定 sink 等 这期间客户端可能因为各种原因(网络抖、客户端崩、超时)无法订阅 但租约已经占了、轮次已经创建了 → 资源泄漏 ``` **Flux.defer 的语义是把内部逻辑的执行从"Flux 对象创建时"推迟到"订阅发生时"**。类比 JS 的 Promise——`new Promise(executor)` 是 executor 立刻执行,而 `Flux.defer(supplier)` 是 supplier 在订阅时才执行。这是 Flux 和 Promise 一个非常关键的区别。 **为什么聊天场景必须 defer?因为这条链路有强副作用**。普通业务接口没这个问题——`SELECT * FROM user WHERE id=1` 即使客户端没收到结果也没副作用,重新调一次无害。但聊天链路要抢 Redis 锁、要创建数据库轮次、要注册到运行态——这些都是**改变系统状态**的操作。如果客户端最终没订阅成功,这些副作用就全部浪费了。 defer 让所有副作用**绑定到订阅事件**——订阅成功就执行副作用、订阅失败副作用根本没发生、系统状态干净。这是响应式编程的核心思想:**把副作用推迟到最晚的合适时机,让系统状态可控**。 ## **五、buildLaunchPlan:把外部参数变成内部可信赖的数据结构** 订阅触发后,第一件事是构建启动计划。这一步本质上是**反腐层**——把外部不可控的请求参数转成内部稳定可信赖的数据结构。 四个核心子动作,每个都有讲究: **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 的产出是一个**不可变快照**——StreamLaunchPlan 对象,所有字段都 final——后续所有环节都基于这个快照工作,外部参数从此和内部状态彻底解耦。 --- ## **六、Redis 分布式租约:同会话全局只有一个生成任务** 启动计划构建完,下一步是抢 Redis 分布式租约。这一段我专门拎出来讲,因为它是这条链路里**保证正确性的关键基础设施**。 ### **6.1 为什么需要租约** 集群部署下,同一个用户对着同一个会话快速点了两次发送、或者前端因为网络抖动重试了一次——如果没有任何互斥机制,**两个生成任务可能在不同节点上同时跑**。结果会是什么?同一个 conversationId 下有两条轮次记录、两份答案文本、两份引用——前端页面状态机彻底乱套,数据也污染了。 租约就是用来解决这个:**按 conversationId 加全局锁,保证任意时刻同会话只有一个生成任务在跑**。 ### **6.2 抢占的三个参数** ```java redisLeaseManager.acquire(leaseKey, leaseOwnerToken, CHAT_RUNNING_LEASE_TTL); ``` - **leaseKey**:`chat:running:{conversationId}` —— 按会话粒度加锁,不同会话彼此不影响。 - **leaseOwnerToken**:UUID —— 本次请求对租约的"身份证",续期和释放时要带上做验证。 - **TTL**:30 秒 —— 防止节点宕机后锁永久残留。 **leaseOwnerToken 防误释放**这一点特别要讲。考虑这个场景: ```text 节点 A 抢到锁(token=T1) 节点 A 卡住,任务超过 30 秒 锁过期,节点 B 抢到锁(token=T2) 节点 A 恢复,执行 release(leaseKey) 无 token 验证:释放了 B 的锁! → 节点 C 又能抢到 → 数据混乱 有 token 验证:发现 token 不匹配,拒绝释放 ``` 加了 token 验证之后,A 想释放只能释放 T1 标记的锁——T1 已经过期、当前锁是 T2,验证失败、拒绝释放。这是 **Redis 分布式锁的标准模式(Redlock 思想)**。 ### **6.3 续期机制:TTL = 3 × 续期间隔** TTL 30 秒不够用怎么办?模型生成可能一两分钟。所以要续期: ```java Flux.interval(CHAT_RUNNING_LEASE_RENEW_INTERVAL).subscribe(ignored -> renewLeaseOrStop(taskInfo)); ``` 每 10 秒续期一次。**TTL 和续期间隔的比例 30:10 = 3:1 是经验法则**——允许最多 2 次续期失败仍不至于让锁过期,给 Redis 抖动、网络延迟留余量。 为什么 TTL 不直接设 1 小时省掉续期?因为节点宕机时锁要等 1 小时才释放——用户 1 小时都没法重发。 为什么续期间隔不是 25 秒(接近 TTL)?任何延迟(GC、网络抖)都可能错过续期,5 秒安全余量太小。 为什么不是 1 秒一次?Redis 压力大,10 秒已经足够频繁。 **续期失败立刻停止任务**——这是租约机制最有意思的一个设计: ```java if (!renewed) { leaseRenewalDisposable.dispose(); // 先停定时器自身,避免重复触发 stopTask(taskInfo, "会话租约已失效,已停止生成"); } ``` 为什么不让任务继续跑?**租约失效有两种原因——Redis 故障(没法判断别人是否抢了)、真有别人抢了(必须停止避免双写)**。保守做法是不管哪种原因都停——代价是用户可能要重发,收益是避免数据混乱。**互斥一旦不确定就要回到安全态**——这是分布式系统里的安全策略。 注意这里**先 dispose 定时器再 stopTask** 的顺序——不先停定时器的话,10 秒后又触发 renewLeaseOrStop、又检测续期失败、又调 stopTask、循环触发、日志大量重复。**先切断重复触发的源头再执行真正的停止逻辑**。 ## **七、bootstrapConversation:三步联动** 抢到租约之后进入 bootstrap,要做三件强关联的事: ```text 1. startExchange:数据库创建轮次记录 → 拿到 exchangeId 2. createTaskInfo:内存组装运行时上下文 → 用到 exchangeId 3. register:把 taskInfo 注册到运行态 → 用 conversationId 做 key ``` 每一步都依赖前面的结果,必须**顺序执行且全部成功**才算 bootstrap 完成。 ### **7.1 双重防御:为什么抢到租约还可能注册失败** ```java if (!chatRuntimeRegistry.register(taskInfo)) { failBootstrappedExchange(...); releaseLeaseQuietly(...); return BootstrapResult.rejected("该会话当前正在执行中,请稍后再试"); } ``` 新手看到这里会问——租约都抢到了,注册怎么还会失败?设想这个场景: ```text 节点 A 服务请求 1(conversationId=X) → 抢到 Redis 锁 → 任务执行中 节点 B 服务请求 2(conversationId=X) → Redis 主从延迟时... A 占的锁还没同步到从库 B 访问从库 → 看到锁不存在 → 也"抢到" 两个节点都以为自己拿到了锁 ``` 这时候本地的 chatRuntimeRegistry 就站出来做最后防御。**chatRuntimeRegistry 是节点本地的内存注册表**,对**同节点上的并发**做最终防御。整个互斥体系是双重的: - **Redis 租约**:跨节点互斥(应对正常情况) - **本地注册表**:节点内互斥(应对 Redis 短暂不一致) ### **7.2 失败时的补偿性收尾** 注册失败时要做两件事,**顺序很关键**: ```java failBootstrappedExchange(...); // 1. 先标记轮次失败(数据库状态推进) releaseLeaseQuietly(...); // 2. 再释放租约 ``` 为什么这个顺序?**先释放租约的话,锁刚释放后续请求进来、看到失败的轮次也开始建新轮次、两个轮次叠加污染数据**。先标记失败、确保数据库状态收敛、再释放锁、下一次进来面对的是干净状态。 外面还有 try-catch 兜底——bootstrap 三步任意一步抛异常都进 catch 释放租约 + 标记失败。注意 exchangeView 是在 try 外声明的局部变量,catch 块才能访问到它判断是否要清理——**try-catch 跨作用域的常见 Java 技巧**。 ### **7.3 TaskInfo:运行时的"万能上下文"** createTaskInfo 这一步组装了整条链路的核心数据载体——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 频繁、StringBuffer 内部 synchronized 性能足够;thinkingSteps、references 用 `Collections.synchronizedList`——add 不频繁、简单互斥就行;usedTools 用 `ConcurrentHashMap.newKeySet()`——并发 add 高、需要无锁;finalized、firstResponseTimeMs 用 AtomicBoolean / AtomicLong——单变量原子操作走 CAS。**不同字段不同诉求,匹配合适的并发原语**。 **finalized 用 AtomicBoolean 通过 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 的经典应用**。 --- ## **八、bindClientChannel + activateGeneration:点火** bootstrap 完成后,bindClientChannel 把 TaskInfo 的内部 sink 转成可返回给前端的 Flux: ```java return taskInfo.sink().asFlux() .doOnSubscribe(ignored -> activateGeneration(taskInfo)) .doOnCancel(() -> stopTask(taskInfo, "客户端已取消请求")); ``` **doOnSubscribe 是整条链路的"点火开关"**——前端建立 SSE 连接的那一刻,后端才真正开始干活。这是 Flux.defer 设计意图的最终落点——所有副作用都绑在订阅事件上。 **doOnCancel 处理前端断开**——用户关掉浏览器、网络断了,前端不再订阅了,后端立刻调 stopTask 中断生成、不让模型空跑浪费 token。 activateGeneration 做两件事:启动租约续期定时器、启动执行链路订阅。两个 Disposable 都存到 TaskInfo 上,停止时统一 dispose。 --- ## **九、buildConversationExecution:把整条响应式流串起来** 这是整条链路最核心的一段代码——把"分析中提示 → 编排执行计划 → 选执行器 → 消费 chunk → 收尾"串成一条响应式流: ```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 里了为什么这里又来一个?外层 defer 是订阅时才点火,内层 defer 让"每次重新订阅都重新执行 thinking 事件 + prepareExecutionPlan"。理论上单次会话不会重订阅但 defer 把这种**重入安全性**当成基础保证,比依赖"不会重订阅"更稳。 **Mono.fromCallable + subscribeOn(boundedElastic)**——prepareExecutionPlan 是同步阻塞方法(数据库查询、检索、压缩),直接在响应式流里调会阻塞当前线程。Mono.fromCallable 把它包装成响应式异步任务,subscribeOn 指定它跑在 boundedElastic 调度器上。**boundedElastic 是 Reactor 专门为阻塞 IO 准备的弹性线程池**——线程数有上限不会爆、空闲线程会回收、适合数据库查询这种 IO 密集型任务。Schedulers.parallel 是 CPU 密集型(线程数 = 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(append answerBuffer + safeEmit 到 sink)。如果 doOnNext 在模型推送的线程上跑(可能是 LLM SDK 内部线程),长时间占用 LLM 线程会影响模型自身的吞吐。publishOn 把 chunk 处理切换到 boundedElastic 上、释放 LLM 线程、让模型可以继续生成下一个 chunk。 **doOnNext / doOnError / doOnComplete 三态闭环**——这是响应式流的三种终态:持续中触发 doOnNext、异常触发 doOnError、正常完成触发 doOnComplete。三种路径都通向收尾——成功走 finishSuccessfully、失败走 finishWithFailure,**保证任何一种路径系统状态都收敛**。 --- ## **十、prepareExecutionPlan:编排器是大脑** ChatPreparationOrchestrator 是整条链路最聪明的部分——它决定这次问答到底走哪条路。prepare() 方法内部按顺序做四件事: **第一,装载会话记忆**——summarizeHistory 拿到 ConversationMemoryContext,里面有长期摘要(之前 N 轮对话的压缩版)和最近窗口(最近几轮原文)。压缩+原文的两段式记忆是为了**远的轮次保留语义但省 token、近的轮次保留细节给问题改写用**。 **第二,构建历史上下文**——buildPlanningHistory 给编排决策用、buildAnswerHistoryContext 给最终答案生成用。两个上下文用途不同所以分开构造。 **第三,判断时效性**——TimeSensitiveQueryHelper 看用户问的是不是"今天/最新/本周"这类需要实时信息的问题,产出两个 boolean:requiresCurrentDateAnchoring(需要日期锚定)和 requiresFreshSearch(需要联网搜索最新事实)。 **第四,根据聊天模式分支路由**: - **OPEN\_CHAT** → ExecutionMode.REACT\_AGENT(开放式 Agent,可调用工具和联网搜索) - **AUTO\_DOCUMENT / DOCUMENT** → 进入文档问答的复杂分支:调 ChatQueryRewriteService 把问题改写成检索友好的表达 → 调 KnowledgeRouteService 路由到目标文档(自动模式下自动选) → 调 DocumentQuestionRouter 判断走 GRAPH\_ONLY、GRAPH\_THEN\_EVIDENCE 还是 RETRIEVAL → 文档范围歧义时走 CLARIFICATION 让用户确认 **编排器有动态修正用户文档选择的能力**——prepareExecutionPlan 外层会检查 `executionPlan.getSelectedDocumentId() != null && !Objects.equals(...selectedDocumentId(), taskInfo.selectedDocumentId())`,如果编排阶段修正了文档范围,要同步刷新归档的会话范围 + RunnableConfig 上下文。**编排器不只是"看用户传了啥跑啥",而是"综合各种信号判断到底应该跑啥"**——这是它叫"大脑"的原因。 ### **buildAgentQuestion:不要让 Agent 裸接用户问题** 编排器产出执行计划之后还要把用户问题"包装"一下——这是 buildAgentQuestion 的职责: ```text 系统时间信息: 当前日期是 2025-05-20,时区为 Asia/Shanghai。 [条件化时效约束] [历史摘要(如有)] 用户问题: [原始问题] ``` 为什么要这么包装?**裸传问题会出问题**——用户问"今天天气怎么样"、不告诉 Agent 今天是几号,它可能用训练数据里的旧日期回答;用户问"还记得我刚才说的吗"、不给历史,Agent 直接懵。 包装的核心是**条件化注入**——requiresCurrentDateAnchoring 时给更严格的日期约束、requiresFreshSearch 时强制要求联网搜索并标注来源日期、有 historySummary 才注入历史。**不同问题有不同的 prompt 强度,精确匹配**。 --- ## **十一、ConversationExecutorRegistry:策略模式 + 自动发现** 执行计划产出后,根据 mode 选执行器: ```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); } } public ConversationExecutor get(ExecutionMode mode) { ... } } ``` 这一段代码非常短但有三个值得讲的设计点: **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 种执行器各管一摊: | 执行模式 | 执行器 | 适用场景 | | --- | --- | --- | | REACT\_AGENT | ReactAgentExecutor | 开放式提问,可调工具和联网搜索 | | RETRIEVAL | RagChatExecutor | 文档知识问答,走完整 RAG 检索 | | GRAPH\_ONLY | GraphOnlyExecutor | 纯图查询,结构化导航类问题 | | GRAPH\_THEN\_EVIDENCE | GraphThenEvidenceExecutor | 先图查询定位再补充证据 | | CLARIFICATION | ClarificationExecutor | 文档范围歧义时让用户确认 | **进阶建议**——可以在构造函数里加启动期校验,确保 ExecutionMode 每个枚举值都有对应执行器,启动直接挂掉比运行时遇到再挂安全。这种"启动时校验闭环"是分布式系统里常用的失败前置手段。 --- ## **十二、emitModelChunk:流式输出三件事** 执行器开始产 chunk,每个 chunk 落到 doOnNext 上、调 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 事件 `{"type":"text","content":"...","timestamp":"...","conversationId":"...","exchangeId":...}`、safeEmit 推到 sink。 **safeEmit 是关键的容错**——它捕获推送异常只记 warn 不抛。原因是**emit 失败抛异常 → doOnError 触发 → finishWithFailure,但实际上模型可能还在正常生成、一次推送失败不应影响整个任务**。answerBuffer 里仍有完整数据、收尾落库不受影响。**把"瞬时推送失败"和"任务真正失败"解耦**。 StreamEventWriter 统一定义了 5 种事件类型: | 类型 | 用途 | 推送时机 | 频次 | | --- | --- | --- | --- | | text | 模型输出正文 | 每个 chunk | 极高 | | thinking | 状态提示 | 阶段切换 | 低 | | error | 错误信息 | 失败时 | 一次 | | reference | 引用来源 | 收尾时 | 一次 | | recommend | 推荐追问 | 收尾时 | 一次 | --- ## **十三、finishSuccessfully:8 步成功收尾** 模型输出完毕、doOnComplete 触发收尾。这一段是整条链路最复杂的方法之一——8 步顺序不能乱: ```text 1. CAS 抢一次性收尾权(finalized.compareAndSet(false, true)) 2. 冻结答案/引用/推荐快照(snapshotReferenceList + deduplicateReferences) 3. 开启 finalize/recommendation 阶段追踪 4. 生成或提取推荐追问(澄清模式复用、常规模式调 RecommendationService) 5. 补发引用和推荐事件到 SSE 6. 关闭 SSE 流(safeComplete) 7. 落库 ChatTurnStatus.COMPLETED(把完整快照 + thinkingSteps + tools 都落) 8. 异步刷会话摘要 + cleanup 释放资源 ``` **为什么这个顺序不能乱?** **CAS 抢权放第一**——保证整个收尾只有一次。停止/失败/成功三种路径都会动 finalized 这个标志,谁先抢到谁执行后面的,后到的直接 return。 **补发事件必须在关闭流之前**——关闭后 sink 已 complete,再 emit 不生效,前端收不到引用和推荐。 **推荐追问的两种来源**——澄清模式(CLARIFICATION)下推荐就是编排阶段已经产出的澄清选项、直接复用避免再调推荐服务生成语义不一致的;常规模式下基于"原问题 + 最终答案 + 最近历史轮次"调 RecommendationService 生成、让推荐贴近本轮回答结果。 **事件顺序:正文 → 引用 → 推荐 → 关闭流**——前端按这个顺序稳定渲染,正文流结束后再看到引用和推荐、避免引用插在正文中间让前端状态机变复杂。**用户先看到答案、再看到来源、最后看到下一步建议**——是符合阅读直觉的流程。 **try / catch / finally 三层嵌套**——这是这段代码最值得讲的工程细节: - **try**:补发引用 + 推荐 → 失败只记 warn 不改判会话失败(**正文已成功,补发失败不能推翻"本轮成功"这一主结果**) - **第一层 finally**:safeComplete 关流 → 即使补发失败也要明确告诉前端"本轮结束" - **第二层 try**:refreshDebugTraceRuntimeStats + completeExchange 落库 → 失败提升日志为 error(**回答成功生成但收尾落库失败、需要重点排查**) - **最内层 finally**:safeRefreshConversationSummary + cleanup → 无论前面所有步骤成功失败都执行、保证资源最终释放 **这种多层嵌套不是为了好看,是为了保证每个失败点都有合适的兜底——正文成功不被补发失败覆盖、补发失败不阻塞关流、关流失败不阻塞落库、落库失败不阻塞清理**。每一层都比下一层"更必须执行"。 --- ## **十四、cleanup:彻底释放资源** 收尾的最后一步: ```java private void cleanup(TaskInfo taskInfo) { Disposable leaseRenewalDisposable = taskInfo.leaseRenewalDisposable(); if (leaseRenewalDisposable != null && !leaseRenewalDisposable.isDisposed()) { leaseRenewalDisposable.dispose(); // 1. 停租约续期定时器 } Disposable disposable = taskInfo.disposable(); if (disposable != null && !disposable.isDisposed()) { disposable.dispose(); // 2. 停业务执行流 } releaseLeaseQuietly(...); // 3. 释放 Redis 租约 chatRuntimeRegistry.remove(...); // 4. 移出本地注册表 } ``` 四个动作分别释放四种资源——续期定时器、业务执行流、Redis 租约、本地注册表。**任何一步失败都用 quietly 版本只记日志不抛异常——cleanup 是终态不应阻塞,租约即使没主动释放 30 秒后也会自动过期**。 cleanup 的本质是**让系统状态回到"这个 conversationId 完全干净、可以接受新请求"的初始态**。 --- ## **十五、异常处理:贯穿全链路** 最后讲一下贯穿整条链路的异常处理思想。链路里每一步都可能失败,统一思路是分层处理: - **致命失败(租约抢占失败、bootstrap 失败)**:立即拒绝、释放已占资源、返回 rejectionFlux - **非致命失败(补发引用、推荐生成)**:降级处理只记 warn,不影响主结果 - **彻底失败(模型异常、执行器异常)**:走 finishWithFailure 完整失败收尾——CAS 抢权 + 失败状态落库 + 关流 + cleanup **任何失败都不影响其他会话**——每个 conversationId 是隔离的、租约+本地注册表保证互斥、失败时所有相关资源都被释放、不会留下悬空状态。 --- ## **十六、整条链路的 12 步串成一句话** 最后用一句话把整条链路收回来: **前端 POST** `/api/chat/stream`**、Controller 返回 Flux 让 WebFlux 以 SSE 持续推送、Flux.defer 把启动逻辑延后到订阅时执行;buildLaunchPlan 把外部参数规范化成不可变的 StreamLaunchPlan(校验问题非空、自动生成 conversationId、按模式严格校验 selectedDocumentId、注入当前日期锚点);****用 Redis 分布式租约(TTL 30 秒 + 10 秒续期 + token 防误释放 + 续期失败立刻停止)保证同会话全局只有一个生成任务****;bootstrapConversation 三步联动——创建轮次归档、组装 TaskInfo 运行时上下文(线程安全集合 + AtomicBoolean finalized + RunnableConfig.context 共享 Map + debugTrace 预占位)、注册到本地运行态做节点内互斥;bindClientChannel 用 doOnSubscribe/doOnCancel 把 sink 转成 Flux,订阅时点火、取消时停任务;activateGeneration 启动租约续期定时器和执行链路;buildConversationExecution 用 Flux.defer + Mono.fromCallable + subscribeOn(boundedElastic) + flatMapMany + publishOn 把"发 thinking → 编排执行计划 → 选执行器 → 消费 chunk → 收尾"串成响应式流;prepareExecutionPlan 调编排器装载会话记忆 + 判时效性 + 按聊天模式路由到 5 种执行模式之一,再用 buildAgentQuestion 把时间锚点 + 时效约束 + 历史摘要 + 原始问题拼成最终 prompt;ConversationExecutorRegistry 用 EnumMap + Spring 自动注入 List 实现策略模式 + 自动发现;emitModelChunk 把每个 chunk 同时做三件事——append 到 answerBuffer、CAS 记首包耗时、safeEmit 推 SSE text 事件;finishSuccessfully 用 CAS 抢一次性收尾权,按"正文 → 引用 → 推荐 → 关闭流 → 落库 COMPLETED → 摘要刷新 → cleanup 释放租约"的严格顺序闭环,try/catch/finally 三层嵌套保证每个失败点都有合适兜底。** 整条链路的核心思想可以归到八个字:**响应式驱动、状态机闭环**。响应式驱动让流式输出和异步处理浑然一体;状态机闭环让任何路径——成功、失败、停止——系统状态都能收敛到干净的初始态。 --- ## **附:一分钟极简版** > 这条链路从前端 POST 开始到答案完整返回结束,分准备/生成/收尾三段。 > > **准备段**:Controller 返回 Flux 走 SSE 协议;Flux.defer 把副作用延后到订阅时;buildLaunchPlan 校验参数、normalize ID、按模式校验文档、注入时间锚点;Redis 租约(TTL 30 秒 + 10 秒续期 + token 防误释放)做跨节点互斥;bootstrap 创建轮次 + 组装 TaskInfo + 注册本地运行态做节点内互斥(双重防御)。 > > **生成段**:bindClientChannel 用 doOnSubscribe/doOnCancel 点火和取消;buildConversationExecution 用 Mono.fromCallable + subscribeOn(boundedElastic) + flatMapMany + publishOn 串响应式流;编排器装载会话记忆、判时效、按模式路由到 5 种执行模式之一;buildAgentQuestion 把时间锚点 + 时效约束 + 历史摘要 + 原问题拼成最终 prompt;ExecutorRegistry 用 EnumMap + Spring 自动注入实现策略模式;emitModelChunk 每个 chunk 做三件事——append answerBuffer、CAS 记首包耗时、safeEmit 推 SSE。 > > **收尾段**:finishSuccessfully 用 CAS 抢一次性收尾权,按"正文 → 引用 → 推荐 → 关闭流 → 落库 COMPLETED → 摘要刷新 → cleanup"严格顺序闭环;try/catch/finally 三层嵌套保证每个失败点都有合适兜底;cleanup 释放续期定时器、执行流、Redis 租约、本地注册表。 > > 整条链路核心思想:**响应式驱动 + 状态机闭环**——任何路径系统状态都能收敛到干净初始态。 --- **企业级项目导航**:⬅️ [[03-提问到返回|03-提问到返回]] | 04-白话讲解 | ➡️ [[00-总览|00-总览]]