--- title: "03-提问到返回" created: 2026-05-21 aliases: - 提问到返回 tags: - 项目 --- # 提问到返回 我们接着上一篇“异步索引构建:落库、向量化与收尾”往下走。 文档处理流水线讲完后,整个 RAG 系统的“**入料阶段**”就闭环了: ```text 上传 → 解析 → 切块 → 向量化 → 索引落库 文档变成可被检索的向量 + 关键词索引 存在 PGVector 和 Elasticsearch 里 ``` 但用户感知不到这些——用户只关心**"我问一个问题,系统怎么给我答案"**。 这一篇就是把这个"使用阶段"彻底拆开来看。学完后你会理解: ```text 1. Controller 为什么返回 Flux 而不是普通 JSON 2. Flux.defer 这一层包装的真实意图 3. 启动计划如何把外部参数变成内部可信赖的数据结构 4. Redis 分布式租约的抢占 + 续期 + 失败自动停止机制 5. TaskInfo 作为运行时"万能上下文"的设计 6. 编排器如何分析意图选择执行模式 7. 执行器注册表的策略模式 + 自动发现 8. 流式输出的事件推送和首包计时 9. 成功收尾的 8 个步骤和事件顺序 10. 异常处理如何贯穿全链路保证资源释放 ``` 下一篇会进入编排器内部——历史摘要怎么生成、意图怎么分析、执行模式怎么决策。 --- ### **一、整体认知:这一节在做什么** 可以一句话概括: > 一次聊天请求从前端 POST 开始,Controller 直接返回 Flux 让 Spring WebFlux 以 SSE 协议持续推送事件,Flux.defer 把启动逻辑延后到前端真正订阅那一刻才执行;然后 buildLaunchPlan 把外部参数规范化成 StreamLaunchPlan(校验问题非空、自动生成 conversationId、按模式校验 selectedDocumentId、注入当前日期锚点);接着用 Redis 分布式租约(TTL 30 秒 + 10 秒续期 + 续期失败自动停止)保证同会话全局只有一个生成任务在跑;bootstrapConversation 创建轮次归档记录、组装 TaskInfo 运行时上下文(单播 sink、线程安全集合、AtomicBoolean finalized、RunnableConfig.context 共享 Map)、注册到运行态;bindClientChannel 用 doOnSubscribe/doOnCancel 把 sink 转成可订阅的 Flux,订阅时点火、取消时停任务;buildConversationExecution 用 Flux.defer + Mono.fromCallable + subscribeOn(boundedElastic) 把"分析中提示 → 编排执行计划 → 选执行器 → 消费 chunk → 收尾"串成响应式流;编排器 prepareExecutionPlan 通过 summarizeHistory 装载会话记忆、判断时效性、按聊天模式路由到 5 种执行模式之一,再通过 buildAgentQuestion 把时间锚点 + 时效约束 + 历史摘要 + 原始问题拼成最终 prompt;ConversationExecutorRegistry 用 EnumMap + Spring 自动注入实现策略模式路由;emitModelChunk 把每个 chunk 同时做三件事——append 到 answerBuffer、记首包耗时、用 streamEventWriter.text 推给前端;finishSuccessfully 用 CAS 抢一次性收尾权,然后按"正文 → 引用 → 推荐 → 关闭流 → 落库 → 摘要刷新 → cleanup 释放租约"顺序闭环。 #### **整体流程图** ```mermaid flowchart TD A[前端 POST /api/chat/stream] --> B[Controller.stream
返回 Flux String] B --> C[Flux.defer 延迟启动] C --> D{前端订阅?} D -->|是| E[buildLaunchPlan
规范化参数] E --> F[claimConversationLease
抢 Redis 租约] F -->|失败| G[rejectionFlux 拒绝流] F -->|成功| H[bootstrapConversation] H --> H1[startExchange 创建轮次] H1 --> H2[createTaskInfo 组装上下文] H2 --> H3[chatRuntimeRegistry.register] H3 --> I[bindClientChannel
doOnSubscribe/doOnCancel] I --> J[activateGeneration] J --> J1[startLeaseRenewal 续期] J1 --> J2[buildConversationExecution] J2 --> K[发送 thinking 事件] K --> L[prepareExecutionPlan
编排器] L --> M[ExecutorRegistry.get
策略路由] M --> N[executor.execute
模型生成] N --> O[emitModelChunk
逐块推送] O --> P{完成?} P -->|正常| Q[finishSuccessfully] P -->|异常| R[finishWithFailure] Q --> Q1[CAS 抢收尾权] Q1 --> Q2[补发引用 + 推荐] Q2 --> Q3[关闭 SSE 流] Q3 --> Q4[落库 COMPLETED] Q4 --> Q5[cleanup 释放资源] ``` 整条链路分成三大块: ```text 准备阶段:Controller → defer → 启动计划 → 抢租约 → bootstrap → 绑通道 生成阶段:激活 → 编排 → 选执行器 → 模型输出 → 推送 收尾阶段:CAS → 补发事件 → 关闭流 → 落库 → 清理 ``` --- ### **二、入口:Controller 返回 Flux** ```java @PostMapping(value = "/stream", produces = "text/event-stream;charset=UTF-8") public Flux stream(@Valid @RequestBody ChatRequestDto dto) { return businessChatService.openConversationStream(dto); } ``` #### **1. Flux 而非 ResponseEntity 的设计** 传统 Controller 长这样: ```java public ResponseEntity> chat(@RequestBody ChatRequestDto dto) { String answer = service.chat(dto); return ResponseEntity.ok(ApiResponse.success(answer)); } ``` 但聊天场景下这种写法**根本走不通**: ```text 模型生成一个完整回答可能要 30 秒 同步等 30 秒前端用户会以为系统死了 即使加 loading 动画,用户体验也很差 ``` 流式输出的本质是**把"等结果"变成"持续推送"**: ```text 模型每生成一个字就推一次 前端用户看到"边打字边显示" 即使总耗时还是 30 秒,感知延迟降到 1-2 秒 ``` #### **2. text/event-stream 内容类型** ```text produces = "text/event-stream;charset=UTF-8" ``` 这是 **SSE(Server-Sent Events)** 协议的标准 MIME 类型: ```text HTTP 响应不会断开,持续保持连接 后端 push 一段数据,前端的 EventSource 立刻触发回调 比 WebSocket 简单(单向通信就够,不需要双向) 比 long polling 高效(不需要反复重连) ``` #### **3. @Valid 在入口层挡掉非法请求** ```text public Flux stream(@Valid @RequestBody ChatRequestDto dto) ``` `@Valid` 会触发 Bean Validation: ```text ChatRequestDto: @NotBlank question @NotBlank chatMode ``` 如果 question 为空,连 Controller 方法体都进不去,直接 400 拒绝。这是**最外层防御**: ```text 非法请求不消耗任何业务资源(不创建轮次、不抢租约、不调模型) 对系统是最大的保护 ``` #### **4. 不包装 ApiResponse 的选择** ```text // 注释里有这句话: // 这里不再额外包装 ApiResponse,而是直接把服务层生成的 SSE 事件流返回给前端逐段消费 ``` 正常 REST 接口都包装成 `{"code": 0, "data": ..., "msg": "..."}`,这里不包装。原因: ```text SSE 不是单次响应,是事件流 事件结构由 StreamEventWriter 统一规定(type + content + timestamp + metadata) 再包一层 ApiResponse 反而打破事件流的标准结构 ``` 错误处理通过 **error 类型的 SSE 事件** 推给前端,而不是 HTTP 状态码或包装结构。这是**流式接口的设计哲学**:协议规范由事件流自己定义,不沿用 REST 包装。 --- ### **三、Flux.defer:延迟启动的真实意图** ```java public Flux openConversationStream(ChatRequestDto request) { return Flux.defer(() -> openDeferredConversationStream(request)); } ``` #### **1. 没有 defer 会发生什么?** 设想去掉 defer: ```java public Flux openConversationStream(ChatRequestDto request) { return openDeferredConversationStream(request); // 直接返回 } ``` `openDeferredConversationStream` 内部会做: ```text 1. buildLaunchPlan(创建对象) 2. claimConversationLease(占 Redis 锁) 3. bootstrapConversation(创建数据库轮次) 4. bindClientChannel(创建 Flux) ``` 问题是:**这些副作用在 Controller return 之前就执行了**: ```text Controller 把 Flux 对象 return 给 Spring WebFlux WebFlux 框架内部还要走一系列处理:写 header、绑定 sink 等 这期间客户端可能因为各种原因(网络抖、客户端崩、超时)无法订阅 但租约已经占了、轮次已经创建了 → 资源泄漏 ``` #### **2. defer 的语义** ```text 没 defer:Flux 对象创建时就执行内部逻辑(eager) 有 defer:Flux 对象创建时只是占位,订阅时才执行内部逻辑(lazy) ``` 类比 JavaScript 的 Promise: ```text new Promise(executor) → executor 立刻执行 Flux.defer(supplier) → supplier 在订阅时才执行 ``` #### **3. 为什么这种延迟特别重要?** 普通业务 API 没有这个问题: ```sql SELECT * FROM user WHERE id = 1 即使客户端没收到,这个 SELECT 也没有副作用 重新调一次也无害 ``` 但聊天 API 有**强副作用**: ```text 抢 Redis 锁:占用了就别人抢不到 创建数据库轮次:留下记录 注册到运行态:占内存 ``` 如果客户端没真正订阅,这些副作用都**浪费**了。defer 让所有副作用**绑定到订阅事件**: ```text 订阅成功 → 执行副作用 → 后续业务推进 订阅失败 → 副作用根本没发生 → 系统状态干净 ``` 这是**响应式编程的核心思想**:把副作用推迟到最晚的合适时机,让系统状态可控。 --- ### **四、buildLaunchPlan:外部参数 → 内部数据结构** ```java private StreamLaunchPlan buildLaunchPlan(ChatRequestDto request) { String question = normalizeQuestion(request.getQuestion()); String conversationId = normalizeConversationId(request.getConversationId()); ChatQueryMode chatMode = parseRequiredChatMode(request.getChatMode()); KnowledgeDocumentDescriptor selectedDocument = resolveSelectedDocument(...); LocalDate currentDate = LocalDate.now(CHAT_ZONE_ID); ... return new StreamLaunchPlan(question, conversationId, chatMode, ...); } ``` 这一步是个**反腐层**——把外部不可控参数转成内部可信赖结构。 #### **1. normalizeConversationId:可选参数的优雅处理** ```java private String normalizeConversationId(String conversationId) { if (StrUtil.isNotBlank(conversationId)) { return conversationId.trim(); } return UUID.randomUUID().toString().replace("-", ""); } ``` ##### **传与不传的语义** ```text 传了 conversationId:继续已有会话 不传 conversationId:开启新会话 ``` ##### **trim 的细节** ```java return conversationId.trim(); ``` 为什么要 trim? ```text 前端通过 URL/cookie/localStorage 传 ID 时可能带空格 不 trim:" abc123 " 和 "abc123" 被当成两个会话 trim:统一去掉边界空格 ``` ##### **UUID.replace("-", "") 的考量** ```text UUID.randomUUID().toString().replace("-", "") ``` ```text 原始 UUID:"550e8400-e29b-41d4-a716-446655440000"(36 字符) 去掉短横线:"550e8400e29b41d4a716446655440000"(32 字符) ``` 为什么去横线? ```text 更短(节省存储) URL 安全(没有特殊字符) 作为 Redis key 不需要转义 "-" 在日志里容易和别的"-"混淆 ``` 是种**轻量惯例**——细节虽小但贯穿全系统会让代码更整齐。 #### **2. parseRequiredChatMode:字符串 → 枚举的强转** ```java private ChatQueryMode parseRequiredChatMode(String value) { ChatQueryMode chatMode = parseOptionalChatMode(value); if (chatMode == null) { throw new IllegalArgumentException("chatMode 不能为空"); } return chatMode; } private ChatQueryMode parseOptionalChatMode(String value) { if (StrUtil.isBlank(value) || "ALL".equalsIgnoreCase(value.trim())) { return null; } try { return ChatQueryMode.valueOf(value.trim().toUpperCase()); } catch (IllegalArgumentException exception) { throw new IllegalArgumentException("chatMode 非法: " + value, exception); } } ``` ##### **Optional 和 Required 拆成两个方法** ```text parseOptional:用于列表查询(可不传) parseRequired:用于聊天请求(必须传) ``` 把"是否必填"的语义放到方法名而不是参数里,调用方一眼能看出意图。 ##### **"ALL" 的特殊处理** ```text if (StrUtil.isBlank(value) || "ALL".equalsIgnoreCase(value.trim())) { return null; } ``` 为什么 ALL 等同于空? ```text 列表查询场景:"查所有聊天模式" 用 ?chatMode=ALL 比 ?chatMode= 更直观 内部统一用 null 表示"不过滤" ALL 是用户友好的别名,内部归一化 ``` 这是**API 设计**和**内部模型**的解耦。 ##### **toUpperCase 大小写不敏感** ```text ChatQueryMode.valueOf(value.trim().toUpperCase()) ``` 前端传 `"open_chat"` / `"Open_Chat"` / `"OPEN_CHAT"` 都能识别。这是**前端容错**——前端开发可能记不准枚举大小写。 ##### **异常包装** ```java catch (IllegalArgumentException exception) { throw new IllegalArgumentException("chatMode 非法: " + value, exception); } ``` `valueOf` 抛的异常信息是 `"No enum constant ChatQueryMode.xxx"`,不够友好。包装成 `"chatMode 非法: xxx"` 让前端能直接展示给用户。 #### **3. resolveSelectedDocument:模式相关的复杂校验** ```java if (chatMode == ChatQueryMode.OPEN_CHAT) { if (normalizedDocumentId != null) { throw new IllegalArgumentException("开放式提问模式下不能传 selectedDocumentId"); } return null; } ... if (chatMode == ChatQueryMode.DOCUMENT) { if (normalizedDocumentId == null) { throw new IllegalArgumentException("当前文档问答模式下必须选择一个文档"); } return documentKnowledgeService.listRetrievableDocuments().stream() .filter(item -> Objects.equals(item.getDocumentId(), resolvedDocumentId)) .findFirst() .orElseThrow(() -> new IllegalArgumentException("所选文档当前不可检索")); } ``` ##### **三种模式的强校验** | 聊天模式 | selectedDocumentId 规则 | 校验逻辑 | | --- | --- | --- | | OPEN\_CHAT | 不允许传 | 传了报错 | | AUTO\_DOCUMENT | 不允许传 | 传了报错 | | DOCUMENT | 必须传 + 必须可检索 | 不传报错,文档不存在/已下线也报错 | ##### **OPEN\_CHAT 和 AUTO\_DOCUMENT 都不允许传** 为什么这两个明明不同的模式,对 selectedDocumentId 的规则一样? ```text OPEN_CHAT:不查任何文档,纯 Agent 模式 AUTO_DOCUMENT:由系统自动判断查哪个文档(可能多个,可能不查) 两者共同点:文档由"系统决定"而非"用户指定" 传 selectedDocumentId 反而是干扰 ``` 如果用户传了,说明用户对模式语义理解错了。**严格拒绝**比"忽略不处理"更好: ```text 忽略:用户以为传了有效,实际系统当没看见,造成困惑 拒绝:明确告知"这个模式下不要传",用户立刻知道用错了 ``` ##### **"当前可检索"的二次校验** ```java documentKnowledgeService.listRetrievableDocuments().stream() .filter(item -> Objects.equals(item.getDocumentId(), resolvedDocumentId)) .findFirst() .orElseThrow(...) ``` 为什么不直接 `documentMapper.selectById(documentId)`? ```text selectById:看文档是否存在 listRetrievableDocuments:看文档是否"现在可被检索" 可能场景: 文档存在(数据库有记录) → 但 indexStatus = BUILD_FAILED → 不可检索 文档存在 → 但被用户下线 → 不可检索 文档存在 → 但索引在重建中 → 暂时不可检索 ``` 用 listRetrievableDocuments 把"业务可用"语义封装在一处,避免每次都重新写过滤条件。 #### **4. 时间锚点的注入** ```java LocalDate currentDate = LocalDate.now(CHAT_ZONE_ID); String currentDateText = formatCurrentDate(currentDate); ``` 这两个字段为什么要在启动计划里就准备好? ```text 启动计划是不可变快照 后续 buildAgentQuestion / 编排器 / 工具调用都用同一个时间 避免不同环节算出不同的"now"(虽然差几毫秒但语义上要一致) ``` ##### **CHAT\_ZONE\_ID 显式指定时区** ```text LocalDate.now(CHAT_ZONE_ID) // 通常是 Asia/Shanghai ``` 为什么不用 `LocalDate.now()`? ```text LocalDate.now():使用 JVM 默认时区 JVM 时区受部署环境影响(Docker 容器、操作系统设置) 服务部署到不同节点可能"今天"不一样 ``` 显式指定 `Asia/Shanghai` 让"今天"在所有节点一致。这是**全局一致性**的小细节。 --- ### **五、Redis 分布式租约** #### **1. 抢占** ```java private boolean claimConversationLease(StreamLaunchPlan launchPlan) { return redisLeaseManager.acquire( launchPlan.getLeaseKey(), launchPlan.getLeaseOwnerToken(), CHAT_RUNNING_LEASE_TTL ); } ``` ##### **三个参数的设计** ```text leaseKey:"chat:running:{conversationId}" 按 conversationId 加锁 不同会话彼此不影响 leaseOwnerToken:UUID 本次请求的"身份证" 续期 / 释放时要带上验证 防止误释放别人的锁 TTL:30 秒 防止节点宕机后锁永久残留 需要持续续期保持有效 ``` ##### **Token 验证防误释放** 考虑场景: ```text 节点 A 抢到锁(token=T1) 节点 A 卡住,任务超过 30 秒 锁过期,节点 B 抢到锁(token=T2) 节点 A 恢复,执行 release(leaseKey) 无 token 验证:释放了 B 的锁! 有 token 验证:发现 token 不匹配,拒绝释放 ``` 加 token 验证防止**误释放**,这是 Redis 分布式锁的标准模式(Redlock 思想)。 #### **2. 续期** ```java private Disposable startLeaseRenewal(TaskInfo taskInfo) { return Flux.interval(CHAT_RUNNING_LEASE_RENEW_INTERVAL, CHAT_RUNNING_LEASE_RENEW_INTERVAL) .subscribe( ignored -> renewLeaseOrStop(taskInfo), error -> log.warn("租约续期任务出现异常...", error) ); } ``` ##### **TTL 30 秒 + 10 秒续期的取舍** ```text 为什么 TTL 不直接设 1 小时,免去续期? 节点宕机 → 锁要等 1 小时才释放 → 用户 1 小时都没法重发 为什么续期间隔不是 25 秒,接近 TTL? 任何延迟(GC、网络抖)都可能错过续期 25 秒间隔 + 30 秒 TTL → 5 秒安全余量太小 10 秒间隔 + 30 秒 TTL → 20 秒余量更稳 为什么续期不是 1 秒一次? 续期太频繁 → Redis 压力大 10 秒已经足够频繁 ``` 这是 **TTL = 3 × 续期间隔** 的经验法则: ```text 允许最多 2 次续期失败仍不至于让锁过期 给 Redis 抖动、网络延迟留余量 ``` ##### **续期失败立刻停止** ```java if (!renewed) { log.warn("会话租约续期失败,准备停止当前会话..."); leaseRenewalDisposable.dispose(); stopTask(taskInfo, "会话租约已失效,已停止生成"); } ``` 为什么不让任务继续跑? ```text 租约失效有两种原因: 1. Redis 故障 → 没法判断别人是否抢了 2. 真有别人抢了 → 必须立刻停止避免双写 保守做法:不管哪种原因,都停止当前任务 代价:用户可能需要重发(但只是这一次) 收益:避免数据混乱(可能影响多个会话) ``` 这是"宁可错杀不可放过"的安全策略——租约的本质是互斥,互斥一旦不确定就要回到安全态。 ##### **续期定时器本身要先停** ```java if (leaseRenewalDisposable != null && !leaseRenewalDisposable.isDisposed()) { leaseRenewalDisposable.dispose(); } stopTask(taskInfo, "..."); ``` 为什么先 dispose 定时器再 stopTask? ```text 定时器还在运行 → 10 秒后又触发 renewLeaseOrStop → 又检测续期失败 → 又调 stopTask 循环触发,日志大量重复 ``` 先停定时器**切断重复触发的源头**,再执行真正的停止逻辑。 #### **3. 租约的整个生命周期** ```mermaid sequenceDiagram participant C as 客户端 participant S as Service participant R as Redis participant T as 续期定时器 C->>S: POST /api/chat/stream S->>R: SET chat:running:{cid} token EX 30 R-->>S: OK S->>T: 启动续期任务(每10秒) S->>C: 流式输出开始 loop 每 10 秒 T->>R: EXPIRE chat:running:{cid} 30 R-->>T: OK end Note over S,C: ... 流式输出中 ... S->>R: DEL chat:running:{cid}(成功收尾) S->>T: dispose() ``` --- ### **六、bootstrapConversation:三步联动** ```java exchangeView = conversationArchiveStore.startExchange(...); TaskInfo taskInfo = createTaskInfo(launchPlan, exchangeView); if (!chatRuntimeRegistry.register(taskInfo)) { failBootstrappedExchange(...); releaseLeaseQuietly(...); return BootstrapResult.rejected(...); } return BootstrapResult.ready(bindClientChannel(taskInfo)); ``` #### **1. 三步的强关联** ```text 1. startExchange:数据库创建轮次记录 → 拿到 exchangeId 2. createTaskInfo:内存组装运行时上下文 → 用到 exchangeId 3. register:把 taskInfo 注册到运行态 → 用 conversationId 做 key ``` 每一步都依赖前面的结果,必须**顺序执行**且**全部成功**才算 bootstrap 完成。 #### **2. 双重防御:为什么抢到租约还可能注册失败?** ```java if (!chatRuntimeRegistry.register(taskInfo)) { failBootstrappedExchange(...); releaseLeaseQuietly(...); return BootstrapResult.rejected("该会话当前正在执行中,请稍后再试"); } ``` 注释里有一段话:"极端情况下,即使抢到租约,也可能在运行态注册时发现已有同会话任务占用,必须补偿性收尾"。 设想这个场景: ```text 节点 A 服务请求 1(conversationId=X) 抢到 Redis 租约 创建轮次,注册 taskInfo 任务执行中 节点 B 服务请求 2(conversationId=X) Redis 租约已被 A 占,理论上抢不到 但 Redis 主从延迟时... A 占的锁还没同步到从库 B 访问从库 → 看到锁不存在 → 也"抢到" 两个节点都以为自己拿到了锁 A 的 chatRuntimeRegistry 已经注册 B 的 chatRuntimeRegistry 在注册时发现同 cid 已存在 → 拒绝 ``` `chatRuntimeRegistry` 是**节点本地的内存注册表**,对**同节点上的并发**做最终防御。这是**双重防御**: ```text Redis 租约:跨节点互斥(应对正常情况) 本地注册表:节点内互斥(应对 Redis 短暂不一致) ``` #### **3. 失败时的补偿性收尾** ```java failBootstrappedExchange(...); // 把已创建的轮次标记为失败 releaseLeaseQuietly(...); // 释放租约 ``` 注意顺序: ```text 1. 先标记轮次失败(数据库状态推进) 2. 再释放租约(允许后续请求重发) ``` 为什么这个顺序? ```text 如果先释放租约: 锁刚释放,后续请求进来,看到失败的轮次也开始建新轮次 两个轮次叠加(虽然不会数据冲突,但污染数据) 先标记失败: 确保数据库状态收敛 再释放锁,下一次进来面对的是干净状态 ``` #### **4. try-catch 兜底** ```java catch (RuntimeException exception) { releaseLeaseQuietly(launchPlan.getLeaseKey(), launchPlan.getLeaseOwnerToken()); if (exchangeView != null) { failBootstrappedExchange(...); } return BootstrapResult.rejected(buildErrorMessage(exception)); } ``` bootstrap 三步任意一步抛异常都进 catch: ```text 1. 释放租约(已经抢到的话) 2. 如果 exchange 已创建,标记为失败 3. 返回拒绝 ``` **不让任何资源残留**。这种 catch 写法的关键: ```text exchangeView 用外部声明的局部变量(不在 try 里声明) catch 块能访问到它来判断是否要清理 如果声明在 try 内部,catch 拿不到 ``` 这是 Java 中**try-catch 跨作用域**的常见技巧。 --- ### **七、TaskInfo:运行时"万能上下文"** #### **1. 字段分类** 按用途分组: ```text 身份标识(不可变): conversationId, exchangeId, traceId, leaseKey, leaseOwnerToken 请求参数(不可变): question, chatMode, selectedDocumentId, currentDate 状态机: executionPlan(后期填充) debugTrace(后期更新) finalized(原子标志) firstResponseTimeMs(原子计时) 运行态集合(线程安全): sink, answerBuffer, thinkingSteps, references, usedTools 基础设施: runnableConfig(Agent 配置) traceRecorder(追踪) eventMetadata(SSE 元数据) 资源句柄: disposable(执行流订阅) leaseRenewalDisposable(续期任务) ``` #### **2. 线程安全的精确选择** 不同字段用不同的线程安全机制: ```text StringBuffer answerBuffer: 需要 append 操作,内部 synchronized 线程安全但单线程比 StringBuilder 慢 Collections.synchronizedList(new ArrayList<>()): 需要 add 操作,简单同步包装 遍历时要外部加锁 ConcurrentHashMap.newKeySet(): 需要 add/contains,高并发场景 内部用分段锁,比 synchronizedSet 高效 AtomicBoolean finalized: 单一 boolean,需要 CAS 比 synchronized 块轻量 AtomicLong firstResponseTimeMs: 单一 long,需要 compareAndSet 保证只设第一次的值 Sinks.Many: Reactor 提供的响应式通道 内部已经处理并发 ``` ##### **为什么不全用一种?** ```text StringBuffer 用于答案累积:append 频繁,数据量大 synchronizedList 用于步骤记录:add 不频繁 ConcurrentHashMap 用于工具集合:add 可能并发 AtomicBoolean/Long 用于简单状态:轻量 CAS 每种场景选最匹配的工具,而不是一刀切 ``` #### **3. RunnableConfig.context() 共享 Map** ```java runnableConfig.context().put(ChatContextKeys.EVENT_SINK, sink); runnableConfig.context().put(ChatContextKeys.QUESTION, launchPlan.getQuestion()); ... ``` ##### **设计意图** 这是个**线程上下文容器**,让深层组件能拿到运行态: ```text 检索器需要 traceId 来记录召回过程 工具调用需要 sink 来推 thinking 事件 模型层需要 currentDate 来回答时间问题 ``` 如果不用 context,方法签名会爆炸: ```text 没 context: retrieve(question, traceId, sink, references, currentDate, ...) 每加一个共享参数,所有方法签名都要改 用 context: retrieve(question, config) 需要什么 config.context().get(KEY) 拿 ``` ##### **类比** ```text RunnableConfig.context ≈ Spring 的 SecurityContext ≈ Node.js 的 async_hooks ≈ Go 的 context.Context ``` 都是同一种模式——**调用链路上的隐式上下文**。 ##### **代价** ```text 优势:解耦,签名简洁,新增字段不破坏 API 代价:隐式依赖,代码 grep 不到谁在用 缓解:用 ChatContextKeys 集中定义所有 key,有据可查 ``` #### **4. initializeDebugTrace 的预占位** ```java ChatDebugTrace debugTrace = initializeDebugTrace(null); runnableConfig.context().put(ChatContextKeys.DEBUG_TRACE, debugTrace); ``` 注意这里**传 null** 给 initializeDebugTrace。注释说:"在真正生成 executionPlan 之前,先放一个空白调试轨迹,保证链路中随时都能读取到 debugTrace"。 ##### **为什么需要空白占位?** ```text prepareExecutionPlan 后会创建真正的 debugTrace 但 prepareExecutionPlan 本身可能调用某些组件 那些组件可能去 context 取 debugTrace 如果 context 还没放,取出来是 null,NPE ``` ##### **解决方案** ```text 先放一个 empty/blank debugTrace(空容器) 组件取出来不会 NPE,只是写入空容器 prepareExecutionPlan 完成后再替换为真实的 ``` 这是**预创建空对象避免 NPE**的常见模式(类似 Null Object Pattern)。 --- ### **八、bindClientChannel:订阅即点火** ```java private Flux bindClientChannel(TaskInfo taskInfo) { return taskInfo.sink().asFlux() .doOnSubscribe(ignored -> activateGeneration(taskInfo)) .doOnCancel(() -> stopTask(taskInfo, "客户端已取消请求")); } ``` #### **1. doOnSubscribe vs doOnNext** ```text doOnSubscribe:订阅时触发一次 doOnNext:每个元素触发 doOnComplete:正常完成时触发 doOnError:异常时触发 doOnCancel:订阅被取消时触发 ``` `doOnSubscribe` 是**整条链路的"点火开关"**——前端建立 SSE 连接的那一刻触发,触发后才执行: ```text startLeaseRenewal:启动续期定时器 buildConversationExecution:启动执行链路 ``` #### **2. doOnCancel:前端断开时的兜底** ```text 用户场景: 用户问完问题,模型正在生成 用户突然关掉浏览器/切到别的页面 SSE 连接断开 没有 doOnCancel: 模型继续生成完整答案 数据库照常落库 但前端没人接收,所有 chunk 被丢弃 浪费算力和模型 token 有 doOnCancel: 检测到客户端断开 → stopTask 主动取消模型生成 清理资源,释放租约 ``` 这是**断线感知**的重要机制。 #### **3. 为什么不在 Service 入口就启动?** 设想直接在 openConversationStream 里启动: ```java public Flux openConversationStream(ChatRequestDto request) { StreamLaunchPlan plan = buildLaunchPlan(request); claimLease(plan); bootstrap(plan); activateGeneration(...); // 直接启动 return Flux<...>(plan.sink); } ``` 问题: ```text 1. 没有 defer:Flux 对象一创建就启动,前端没订阅就跑 2. 没有 doOnCancel:前端断开后无感知 3. 启动逻辑和 Flux 对象耦合 ``` 用 `Flux.defer + doOnSubscribe + doOnCancel` 的组合: ```text defer:订阅时才进入业务逻辑 doOnSubscribe:订阅时点火生成 doOnCancel:取消时兜底清理 ``` 三者配合形成**生命周期完整的响应式管道**。 --- ### **九、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)); ``` #### **1. 嵌套 Flux.defer 的意图** 外层 `bindClientChannel` 返回的 Flux 已经在 defer 里了,为什么这里又来一个 defer? ```text 外层 defer:订阅时才点火 内层 defer:每次重新订阅都重新执行 thinking 事件 + prepareExecutionPlan ``` 理论上单次会话不会重订阅,但 defer 把这种**重入安全性**当成基础保证,比依赖"不会重订阅"更稳。 #### **2. Mono.fromCallable + subscribeOn(boundedElastic)** ```text Mono.fromCallable(() -> prepareExecutionPlan(taskInfo)) .subscribeOn(Schedulers.boundedElastic()) ``` ##### **Mono.fromCallable 的作用** ```text prepareExecutionPlan 是一个同步阻塞方法(包含数据库查询、检索、压缩等) 直接调用会阻塞当前线程 Mono.fromCallable 把它包装成响应式异步任务 ``` ##### **subscribeOn(boundedElastic) 的作用** ```text boundedElastic:Reactor 提供的"弹性边界"线程池 适合 IO 密集型任务(数据库、网络调用) 有上限(避免无限扩张) 空闲时回收线程 ``` ##### **为什么不用 parallel 或 single?** ```text Schedulers.parallel:CPU 密集型(线程数 = CPU 核心数) prepareExecutionPlan 主要是 IO,不适合 Schedulers.single:单线程 所有任务串行,并发量上不去 不适合多用户同时聊天 Schedulers.boundedElastic:IO 优化 完美匹配"调数据库 + 检索 + 调外部"的场景 ``` #### **3. flatMapMany:Mono → Flux 的转换** ```java .flatMapMany(plan -> { ConversationExecutor executor = conversationExecutorRegistry.get(plan.getMode()); return executor.execute(taskInfo); // 返回 Flux }); ``` ##### **为什么是 flatMapMany 而不是 map?** ```text map:T → R 一对一变换 flatMap:T → Mono 异步一对一 flatMapMany:T → Flux 异步一对多 ``` prepareExecutionPlan 产出一个 plan,executor.execute 产出**多个 chunk**: ```text plan → executor → [chunk1, chunk2, chunk3, ...] 一对多关系,只能用 flatMapMany ``` ##### **为什么不能直接拿到 executor 后再外层 .map?** ```text .map(plan -> executor.execute(taskInfo)) 返回 Flux>(嵌套) 需要再 .flatMap(flux -> flux) 解嵌套 不如 .flatMapMany 直接搞定 ``` #### **4. publishOn vs subscribeOn** ```text .publishOn(Schedulers.boundedElastic()) ``` ##### **两者区别** ```text subscribeOn:订阅链路向上传播,影响整个流的上游执行线程 publishOn:从这里往下,切换到指定线程池 ``` ##### **这里 publishOn 的意图** ```text 模型每产生一个 chunk,doOnNext 会执行 emitModelChunk emitModelChunk 内部: append 到 answerBuffer safeEmit 到 sink 如果 doOnNext 在模型推送的线程(可能是 LLM SDK 内部线程): 长时间占用 LLM 线程 可能阻塞下一个 chunk 推送 publishOn 切到 boundedElastic: 每个 chunk 在新线程上处理 LLM 线程立刻释放,可以推下一个 chunk ``` 这是 Reactor 中**生产者和消费者解耦**的标准做法。 #### **5. doOnNext / doOnError / doOnComplete 的三态闭环** ```java .doOnNext(chunk -> emitModelChunk(...)) // 每个 chunk .doOnError(error -> finishWithFailure(...)) // 异常 .doOnComplete(() -> finishSuccessfully(...)); // 正常完成 ``` 这是响应式流的**三种终态**: ```text 持续中:doOnNext 不断触发 正常完成:doOnComplete 触发一次,流终止 异常终止:doOnError 触发一次,流终止 ``` doOnError 和 doOnComplete **互斥**——一个流要么以 complete 结束要么以 error 结束,不会两者都触发。 #### **6. 完整数据流图** ```mermaid sequenceDiagram participant F as 前端 participant C as Controller participant S as Service participant O as Orchestrator participant E as Executor participant M as Model F->>C: POST /api/chat/stream C->>S: openConversationStream S-->>C: Flux(deferred) C-->>F: 200 OK, text/event-stream F->>C: 订阅 SSE C->>S: doOnSubscribe 触发 S->>S: activateGeneration S->>F: SSE event: thinking S->>O: prepareExecutionPlan (boundedElastic) O-->>S: ConversationExecutionPlan S->>E: executor.execute E->>M: 调模型 API loop 模型生成 M-->>E: chunk E-->>S: Flux push S->>S: emitModelChunk (publishOn) S->>F: SSE event: text end M-->>E: 完成 E-->>S: doOnComplete S->>S: finishSuccessfully S->>F: SSE event: reference S->>F: SSE event: recommend S->>F: 关闭流 ``` --- ### **十、编排器和 buildAgentQuestion** #### **1. 编排器是大脑** ```java private ConversationExecutionPlan prepareExecutionPlan(TaskInfo taskInfo) { ConversationExecutionPlan executionPlan = chatPreparationOrchestrator.prepare(taskInfo); executionPlan.setAgentQuestion(buildAgentQuestion(executionPlan)); if (executionPlan.getSelectedDocumentId() != null && !Objects.equals(executionPlan.getSelectedDocumentId(), taskInfo.selectedDocumentId())) { conversationArchiveStore.refreshSessionScope(...); ... } taskInfo.setExecutionPlan(executionPlan); ... return executionPlan; } ``` ##### **编排器的核心职责** ```text 1. 装载会话记忆(summarizeHistory) 2. 判断时效性(requiresCurrentDateAnchoring / requiresFreshSearch) 3. 路由执行模式(REACT_AGENT / RETRIEVAL / GRAPH_ONLY / ...) 4. 必要时改写文档范围 ``` ##### **文档范围动态修正** ```java if (executionPlan.getSelectedDocumentId() != null && !Objects.equals(executionPlan.getSelectedDocumentId(), taskInfo.selectedDocumentId())) { conversationArchiveStore.refreshSessionScope(...); putContextIfNotNull(taskInfo.runnableConfig(), ChatContextKeys.SELECTED_DOCUMENT_ID, ...); ... } ``` 这是个有趣的设计——**编排器可以修正用户的文档选择**。 什么场景下会修正? ```text AUTO_DOCUMENT 模式:用户没选文档 编排器分析问题 → 决定该查 doc_5 把 selectedDocumentId 设为 5 需要同步更新归档和上下文 DOCUMENT 模式但选错了: 编排器发现问题和当前文档完全不匹配 可能切到 CLARIFICATION 模式让用户确认 ``` 修正后要做三件事: ```text 1. refreshSessionScope:更新数据库归档的会话范围 2. 更新 runnableConfig.context:让后续工具看到新的文档 ID 3. 内存对象 taskInfo 在 setExecutionPlan 时一起更新 ``` #### **2. buildAgentQuestion:为什么不能裸传问题** ##### **裸问题的问题** ```text 用户问:"今天天气怎么样?" 直接传给模型: "今天" 对模型来说没有锚点 模型可能用训练数据里的某个日期当"今天" 或者拒绝回答说"我不知道今天日期" ``` ##### **包装后的 prompt** ```text 系统时间信息: 当前日期是 2025年5月20日(星期二),时区为 Asia/Shanghai。 当前问题包含相对时间或强时效语义。当用户提到"今天、明天、昨天..."时, 必须以这个日期为准,不要把搜索结果里的旧日期误当成今天。 当前问题需要核实最新外部事实,回答前必须优先调用联网搜索工具。 如果搜索结果里的日期与当前日期不一致,必须明确说明来源日期。 相关会话背景: 用户之前问过北京天气,本轮可能延续这一话题。 用户问题: 今天天气怎么样? ``` ##### **三段式 prompt 结构** ```text 1. 系统时间锚点:统一时间基准 2. 上下文摘要:让 Agent 知道之前聊了什么 3. 原始问题:用户的真实诉求 ``` 这种"约束 + 上下文 + 问题"的 prompt 结构是 LLM 应用的常见模式。 ##### **条件化 prompt** ```java if (executionPlan.isRequiresCurrentDateAnchoring()) { builder.append("当前问题包含相对时间或强时效语义..."); } else { builder.append("当用户提到"今天、明天、昨天..."时,必须以这个日期为准。\n"); } if (executionPlan.isRequiresFreshSearch()) { builder.append("当前问题需要核实最新外部事实,回答前必须优先调用联网搜索工具。\n"); ... } ``` 不同问题给不同强度的约束: ```text 普通问题(不需要时效性): 只提一句"如果用户提到今天,以这个日期为准" 强时效问题(需要锚点): 详细说明"不要把搜索结果里的旧日期误当成今天" 强调日期一致性校验 强时效 + 需要联网: 强制调用搜索工具 要求标注来源日期 无法找到匹配结果时要明确说明不确定性,不要编造 ``` ##### **为什么不一刀切都给最强约束?** ```text Prompt 越长 → token 消耗越多 → 成本越高 → 推理越慢 弱时效问题(比如"什么是机器学习")给联网搜索约束: 模型可能莫名其妙调搜索工具 浪费 token,降低响应速度 强时效问题(比如"今天股市")不给联网约束: 模型用训练数据回答 → 给出过时信息 → 用户被误导 ``` **精确匹配 prompt 强度和问题特征**,是 RAG 系统调优的关键之一。 --- ### **十一、ConversationExecutorRegistry:策略模式 + 自动发现** ```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) { ConversationExecutor executor = executorMap.get(mode); if (executor == null) { throw new IllegalStateException("未找到执行模式对应的执行器: " + mode); } return executor; } } ``` #### **1. Spring 自动注入 List** ```text public ConversationExecutorRegistry(List executors) { ``` 这是 Spring 的"魔法": ```java Spring 扫描所有实现了 ConversationExecutor 接口的 Bean 自动收集成 List 注入构造函数 新增一个 Executor 类(加 @Component) → 自动加入 List → 自动注册 ``` ##### **对比传统注册方式** ```java 传统: @Bean public ExecutorRegistry registry() { ExecutorRegistry r = new ExecutorRegistry(); r.add(new ReactAgentExecutor()); r.add(new RagChatExecutor()); r.add(new GraphOnlyExecutor()); ...每加一个执行器都要改这里 return r; } 自动发现: 新增 GraphThenEvidenceExecutor → 加 @Component → 自动注册 完全不用动 ExecutorRegistry ``` #### **2. EnumMap 的选择** ```text new EnumMap<>(ExecutionMode.class) ``` ##### **为什么不用 HashMap?** ```text EnumMap:专门为枚举 key 优化,内部用数组实现 O(1) 查找,内存紧凑 比 HashMap 性能更好 HashMap:通用实现,要算 hashCode,处理冲突 对枚举 key 来说浪费 ``` 虽然差异微小,但**正确选择数据结构**是工程素养。 #### **3. mode() 方法的设计** ```typescript public interface ConversationExecutor { ExecutionMode mode(); // 执行器自己声明处理哪种模式 Flux execute(TaskInfo taskInfo); } ``` 这种"自我声明"模式比外部注解或配置更灵活: ```text Bean 自己最清楚自己处理什么 注册表只负责"问 Bean → 收集映射" 新增模式不需要改注册表代码 ``` #### **4. 启动期校验** ```java public ConversationExecutor get(ExecutionMode mode) { ConversationExecutor executor = executorMap.get(mode); if (executor == null) { throw new IllegalStateException("未找到执行模式对应的执行器: " + mode); } return executor; } ``` ##### **抛 IllegalStateException 而非业务异常** ```text IllegalStateException 表示"程序状态错误" 不是用户输入错误,是开发者忘了实现某个 Executor 应该让程序崩溃 + 报警 → 立刻修复 ``` ##### **更好的做法:启动期校验** 可以在 ExecutorRegistry 构造函数里加: ```java public ConversationExecutorRegistry(List executors) { for (ConversationExecutor executor : executors) { executorMap.put(executor.mode(), executor); } // 启动期校验:确保所有 ExecutionMode 都有对应执行器 for (ExecutionMode mode : ExecutionMode.values()) { if (!executorMap.containsKey(mode)) { throw new IllegalStateException("ExecutionMode " + mode + " 没有对应执行器"); } } } ``` 启动直接挂掉 → 比运行时遇到再挂安全得多。 #### **5. 整体架构** ```mermaid flowchart LR A[Orchestrator 决定模式] --> B[Registry.get mode] B --> C{EnumMap 查找} C -->|REACT_AGENT| D[ReactAgentExecutor] C -->|RETRIEVAL| E[RagChatExecutor] C -->|GRAPH_ONLY| F[GraphOnlyExecutor] C -->|GRAPH_THEN_EVIDENCE| G[GraphThenEvidenceExecutor] C -->|CLARIFICATION| H[ClarificationExecutor] D --> I[execute returns Flux String] E --> I F --> I G --> I H --> I ``` 新增执行模式只需要: ```java 1. 在 ExecutionMode 枚举加新值 2. 实现 ConversationExecutor 接口 3. 加 @Component 注解 4. 重启服务 完全不用改 Registry 和 Orchestrator(除了 Orchestrator 的路由逻辑) ``` 这是**开闭原则**(Open-Closed Principle)的标准实践——对扩展开放,对修改关闭。 --- ### **十二、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())); } ``` #### **1. 三件事的顺序** ```text 1. append 到 answerBuffer:累积完整答案,收尾时落库用 2. 记录首包耗时:监控指标,衡量响应速度 3. 推送到前端:用户即时看到 ``` ##### **为什么是这个顺序?** ```text 先 append 再推送: 确保数据已经累积到 buffer 即使后面推送失败,buffer 里也有完整数据 收尾落库不受推送失败影响 先记首包再推送: 时间戳更精确(贴近实际接收时间) ``` #### **2. CAS 记录首包耗时** ```text if (taskInfo.firstResponseTimeMs().get() == 0L) { taskInfo.firstResponseTimeMs() .compareAndSet(0L, System.currentTimeMillis() - taskInfo.startTime()); } ``` ##### **为什么先 get 判 0 再 CAS?** ```text 直接 CAS:每次 chunk 都尝试 CAS,无效操作多 chunk 可能上千个,每次都 CAS 浪费 先 get:大多数情况快速跳过(O(1) 读) 只有第一次 get == 0 时才进入 CAS 显著降低开销 ``` ##### **CAS 而非简单赋值** ```text 为什么不直接 firstResponseTimeMs.set(...)? 考虑边界: 多个线程同时进入 emitModelChunk(理论上不会,但保险起见) 第一个线程进入 if,准备 set 第二个线程也进入 if,也准备 set 两个 set 都执行 → 后者覆盖前者 用 CAS: 只有第一个 compareAndSet 成功 第二个失败,值不变 保证记录的是真正的"首包"时间 ``` ##### **首包耗时的业务价值** ```text TTFB(Time To First Byte) = 首包耗时 是聊天 API 的核心 SLA 指标 监控告警:首包耗时 > 5 秒 → 报警 A/B 测试:对比不同模型/提示词的首包速度 用户体验:首包越快,用户感知越好 ``` #### **3. safeEmit 的兜底** ```java safeEmit(taskInfo.sink(), streamEventWriter.text(chunk, ...)); ``` ##### **safeEmit 的实现 (推测)** ```java private void safeEmit(Sinks.Many sink, String event) { try { sink.tryEmitNext(event); } catch (Exception e) { log.warn("emit 失败", e); } } ``` ##### **为什么 emit 可能失败?** ```text sink 已经被 dispose(任务被取消) sink 已经 complete(任务正常结束) 背压:消费速度跟不上生产速度 ``` ##### **为什么不让失败抛出?** ```text emit 失败抛异常: doOnError 触发 → finishWithFailure 实际上模型可能还在正常生成 一次推送失败不应影响整个任务 safeEmit 只记 warn: 继续累积 chunk 到 buffer 最终落库还是完整的 ``` 这是**容错设计**——把"瞬时推送失败"和"任务真正失败"解耦。 #### **4. 流式输出的体验** 模型每输出一小段就立刻推: ```text 模型生成"根据您提供的文档,Spring WebFlux 是..." 分成 chunk: "根据" "您" "提供的" "文档," "Spring " "WebFlux " "是..." 每个 chunk 立刻推送 → 用户看到"打字机效果" ``` ##### **chunk 粒度的取舍** ```text chunk 太大(整段才推): 用户等很久才看到东西 "边生成边显示"效果消失 chunk 太小(每个 token 一推): 网络开销大(每个 SSE 事件有 header 开销) 前端渲染压力大 通常折中:几个 token 一个 chunk 模型 SDK 会根据网络条件自动调整 ``` --- ### **十三、StreamEventWriter:SSE 事件标准化** ```java public String text(String content, StreamEventMetadata metadata) { return write(event("text", content, metadata)); } private Map event(String type, Object content, StreamEventMetadata metadata) { Map payload = new LinkedHashMap<>(); payload.put("type", type); payload.put("content", content); payload.put("timestamp", Instant.now().toString()); if (metadata != null) { if (metadata.conversationId() != null) { payload.put("conversationId", metadata.conversationId()); } if (metadata.exchangeId() != null && metadata.exchangeId() > 0) { payload.put("exchangeId", metadata.exchangeId()); } } return payload; } ``` #### **1. 五种事件类型** | 类型 | 用途 | 推送时机 | 频次 | | --- | --- | --- | --- | | text | 模型输出正文 | 每个 chunk | 极高 | | thinking | 状态提示 | 阶段切换 | 低 | | error | 错误信息 | 失败时 | 一次 | | reference | 引用来源 | 收尾时 | 一次 | | recommend | 推荐追问 | 收尾时 | 一次 | ##### **设计哲学** ```text type 字段表明事件类别 → 前端按类型分发处理 content 字段灵活承载 → 字符串、List、Map 都能放 timestamp + metadata → 调试和追踪信息 ``` #### **2. LinkedHashMap 的选择** ```java Map payload = new LinkedHashMap<>(); ``` ##### **为什么不用 HashMap?** ```text LinkedHashMap:保持插入顺序 序列化成 JSON 时字段顺序固定 "type" 始终在第一个位置 HashMap:不保证顺序 每次序列化字段顺序可能不同 日志和调试时不一致 ``` 虽然 JSON 顺序在协议层不重要,但**人眼阅读和日志比对**时一致顺序很有帮助。 #### **3. 条件性放入字段** ```java if (metadata.conversationId() != null) { payload.put("conversationId", metadata.conversationId()); } if (metadata.exchangeId() != null && metadata.exchangeId() > 0) { payload.put("exchangeId", metadata.exchangeId()); } ``` ##### **为什么不直接 put(value 可能为 null)?** ```text 直接 put null: JSON 序列化出 "conversationId":null 前端要判断 null,代码繁琐 增加传输大小 条件 put: null 时字段根本不出现 前端代码可以直接用 hasOwnProperty 或 ?. 安全访问 ``` 这是**Java 后端 → JSON 前端**的常见处理细节。 #### **4. 标准事件结构** ```json { "type": "text", "content": "根据文档内容,", "timestamp": "2025-05-20T10:30:00.123Z", "conversationId": "abc123", "exchangeId": 42 } ``` ##### **前端处理代码** ```java const eventSource = new EventSource('/api/chat/stream', {...}); eventSource.onmessage = (event) => { const payload = JSON.parse(event.data); switch (payload.type) { case 'text': appendToAnswer(payload.content); break; case 'thinking': showThinkingState(payload.content); break; case 'reference': renderReferences(payload.content); break; case 'recommend': renderRecommendations(payload.content); break; case 'error': showError(payload.content); break; } }; ``` 统一的事件结构让**前端代码极其简洁**。 --- ### **十四、finishSuccessfully:8 步成功收尾** ```java private void finishSuccessfully(TaskInfo taskInfo) { if (!taskInfo.finalized().compareAndSet(false, true)) { return; } String answer = taskInfo.answerBuffer().toString(); List uniqueReferences = deduplicateReferences(...); ConversationTraceRecorder.StageHandle finalizeStage = ...; ConversationTraceRecorder.StageHandle recommendationStage = ...; List recommendations; if (...CLARIFICATION...) { recommendations = ...clarificationOptions; } else { recommendations = recommendationService.generateRecommendations(...); } try { if (!uniqueReferences.isEmpty()) { safeEmit(sink, streamEventWriter.references(...)); } if (!recommendations.isEmpty()) { safeEmit(sink, streamEventWriter.recommendations(...)); } } catch (RuntimeException exception) { log.warn("补发引用或推荐事件失败", exception); } finally { try { safeComplete(taskInfo.sink()); } catch (...) {} try { refreshDebugTraceRuntimeStats(taskInfo); conversationArchiveStore.completeExchange(...); ... } catch (...) {} finally { safeRefreshConversationSummary(...); cleanup(taskInfo); } } } ``` #### **1. CAS 抢"唯一收尾权"** ```text if (!taskInfo.finalized().compareAndSet(false, true)) { return; } ``` ##### **为什么需要 CAS?** ```text 成功收尾 / 失败收尾 / 主动停止 都要执行收尾逻辑 都可能在不同线程被触发: doOnComplete:模型完成线程 doOnError:任意失败线程 stopTask(用户取消):前端线程 stopTask(续期失败):续期定时器线程 没 CAS: 可能两个收尾同时执行 重复落库 → 数据库唯一约束冲突 重复释放租约 → 释放别人的锁 CAS: 只有第一个 compareAndSet 成功 其他线程直接 return,不重复执行 ``` ##### **CAS 的语义** ```text compareAndSet(expected, newValue): 如果当前值 == expected: 设置为 newValue,返回 true 否则: 不修改,返回 false 第一次调用:finalized 是 false → CAS 成功 → 设为 true 第二次调用:finalized 已经是 true → CAS 失败 → 返回 false ``` #### **2. 数据快照** ```java String answer = taskInfo.answerBuffer().toString(); List uniqueReferences = deduplicateReferences(snapshotReferenceList(taskInfo.references())); ``` ##### **为什么要快照?** ```text references 是线程安全 List,但仍然可能被异步线程写入 比如有个晚到的工具结果 收尾期间继续 add → 落库的数据和推送给前端的不一致 snapshotReferenceList:在某个时刻拷贝出来 后续不再受外部修改影响 保证一致性 ``` ##### **deduplicateReferences:去重** ```text 不同检索路径可能命中同一文档 向量检索召回 doc_5 关键词检索也召回 doc_5 references 里有两份 doc_5 去重:同 documentId + chunkId 算同一引用 用户看到的是干净的引用列表 落库的也是去重后的 ``` #### **3. CLARIFICATION 模式的特殊处理** ```java if (taskInfo.executionPlan() != null && taskInfo.executionPlan().getMode() == ExecutionMode.CLARIFICATION) { recommendations = ...getClarificationOptions(); } else { recommendations = recommendationService.generateRecommendations(...); } ``` ##### **为什么 CLARIFICATION 复用澄清选项?** ```yaml 正常模式: 问:今天天气? 答:晴朗 25 度 推荐:明天天气? / 这周降水? → 推荐内容是"答案的延伸" CLARIFICATION 模式: 问:介绍一下产品 答:你想了解哪个产品? A/B/C 澄清选项:[A 产品介绍, B 产品介绍, C 产品介绍] 如果再调推荐服务: 基于"介绍一下产品"+ "你想了解哪个产品"生成 可能产生:产品有什么? / 怎么购买? 之类 但用户当前阶段需要的是 A/B/C 选择 额外推荐反而干扰 ``` 直接复用澄清选项作为推荐——**"问什么就推什么"**,不浪费一次推荐生成。 #### **4. 补发事件 + 关闭流的顺序** ```text 1. 补发 references 2. 补发 recommendations 3. 关闭流(safeComplete) ``` ##### **为什么这个顺序?** ```text 关闭流后再 emit:无效(流已关闭) 所以必须先 emit 再关闭 事件之间也有顺序: references 在 recommendations 前: 前端先看到"答案的依据" 再看到"接下来可以问什么" 逻辑上是"先回顾再展望" ``` ##### **safeComplete 的作用** ```java sink.tryEmitComplete(); ``` ```text 通知前端:本次流结束 前端 EventSource.close 触发 连接断开,资源释放 ``` #### **5. try-catch-finally 的精妙嵌套** ```java try { // 补发事件 } catch (RuntimeException exception) { log.warn(...); // 不改判失败 } finally { try { safeComplete(...); // 关闭流 } catch (...) {} try { refreshDebugTraceRuntimeStats(...); conversationArchiveStore.completeExchange(...); // 落库 } catch (...) { log.error(...); // 落库失败是 error 级别 } finally { safeRefreshConversationSummary(...); // 异步刷摘要 cleanup(...); // 资源清理 } } ``` ##### **三层 try-catch 的职责** ```text 最外层 try-catch: 捕获补发事件的异常 warn 级别,不影响主流程 中间 try-catch: 分别保护 safeComplete / 落库 每个独立 try,失败不影响下一个 最内层 finally: 无论前面所有步骤成功失败,都要执行 刷摘要 + cleanup 保证资源最终释放 ``` ##### **日志级别的精确选择** ```text 补发事件失败:warn 答案已生成成功,只是补充信息没推 用户体验不完美但不致命 关闭流失败:warn 流可能已自动关闭,影响小 落库失败:error 用户看到了答案,但数据库没记录 需要紧急排查 摘要刷新失败:silent(在 safeRefresh 内部处理) 最不重要的步骤,不该影响主流程 ``` ##### **为什么补发失败不改判会话?** 注释明确说:"正文已经成功完成,补发事件失败不应推翻'本轮回答成功生成'这一主结果"。 ```yaml 用户视角: 答案已经在屏幕上看到了 ✓ 引用没看到 → 体验差但不致命 推荐没看到 → 体验差但不致命 技术视角: 模型 token 已经消耗 ✓ 数据库轮次还要落库 ✓ 如果改判失败: 数据库状态:FAILED 用户实际体验:成功 数据和体验割裂,运维排查困难 ``` #### **6. 事件推送顺序总览** ```mermaid sequenceDiagram participant E as Executor participant S as Sink participant F as 前端 Note over E,F: 模型生成阶段 loop 每个 chunk E->>S: emit text event S->>F: SSE: text end Note over E,F: 收尾阶段(finishSuccessfully) E->>S: emit reference event S->>F: SSE: reference E->>S: emit recommend event S->>F: SSE: recommend E->>S: complete signal S->>F: 关闭连接 Note over E: 落库 + 清理(用户已断开,异步进行) ``` #### **7. cleanup 的彻底清理** ```java private void cleanup(TaskInfo taskInfo) { Disposable leaseRenewalDisposable = taskInfo.leaseRenewalDisposable(); if (leaseRenewalDisposable != null && !leaseRenewalDisposable.isDisposed()) { leaseRenewalDisposable.dispose(); } Disposable disposable = taskInfo.disposable(); if (disposable != null && !disposable.isDisposed()) { disposable.dispose(); } releaseLeaseQuietly(taskInfo.leaseKey(), taskInfo.leaseOwnerToken()); chatRuntimeRegistry.remove(taskInfo.conversationId(), taskInfo); } ``` ##### **清理顺序** ```text 1. 停续期定时器(否则继续触发 Redis 操作) 2. 停业务执行流(确保 Reactor 资源释放) 3. 释放 Redis 租约(允许后续请求) 4. 移除运行态注册(允许同节点接受新请求) ``` ##### **每步的 isDisposed 检查** ```text if (... != null && !disposable.isDisposed()) ``` ```text 为什么要检查 isDisposed? 可能这个 disposable 已经被某个分支主动 dispose 过 重复 dispose 一般无害,但有些实现可能抛异常 防御性编程,确保不出异常 ``` ##### **releaseLeaseQuietly:静默释放** ```text release 失败的可能原因: 租约已过期(自动释放了) Redis 故障 失败也不报错: cleanup 是终态,不应阻塞 租约即使没主动释放,30 秒后也会自动过期 ``` --- ### **十五、总结:整条链路的设计哲学** #### **1. 12 步流程回顾** ```text 1. Controller 接收请求 → 参数校验,转交 Service 2. Flux.defer 延迟启动 → 等前端真正订阅 3. buildLaunchPlan → 规范化参数,构建启动计划 4. claimConversationLease → 抢占 Redis 分布式租约 5. bootstrapConversation → 创建轮次记录、构建 TaskInfo、注册运行态 6. bindClientChannel → 绑定 SSE 输出通道 7. activateGeneration → 启动租约续期 + 执行链路 8. prepareExecutionPlan → 编排器分析意图,生成执行计划 9. ConversationExecutorRegistry.get → 根据模式选择执行器 10. executor.execute → 执行器工作,模型开始生成 11. emitModelChunk → 逐块推送文本到前端 12. finishSuccessfully → 补发引用和推荐、落库归档、释放资源 ``` #### **2. 贯穿全链路的设计原则** ##### **响应式驱动** ```text 全链路无阻塞 → Flux/Mono 串联 publishOn/subscribeOn 精确控制线程 背压由 Reactor 自动处理 ``` ##### **延迟启动 + 副作用收敛** ```text Flux.defer 延迟到订阅时 副作用绑定到订阅事件 未订阅 → 零资源消耗 ``` ##### **多层防御** ```text 入口 @Valid:格式校验 buildLaunchPlan:语义校验 claimLease:Redis 跨节点互斥 chatRuntimeRegistry:节点内互斥 finalized CAS:收尾互斥 ``` ##### **资源管理闭环** ```text 任何资源(租约/订阅/注册)都有"获取 + 释放"对 异常时通过 catch 清理已获取的资源 finally 兜底确保清理 不留悬空资源 ``` ##### **状态可观测** ```text ConversationTraceRecorder 记录每个阶段 StreamEventMetadata 注入每个事件 debugTrace 累积运行态统计 firstResponseTimeMs 监控指标 ``` ##### **失败容错的层级** ```text 致命失败(如租约抢占):立即拒绝 非致命失败(如补发事件):降级处理 彻底失败(如模型异常):走 finishWithFailure 任何失败都不影响其他会话 ``` #### **3. 系统能力图谱** ```mermaid mindmap root((聊天链路)) 协议层 SSE 流式 Flux WebFlux 事件类型 5种 并发控制 Redis 租约 TTL 30s 10s 续期 CAS 收尾 本地注册表 上下文管理 TaskInfo RunnableConfig 会话记忆 时间锚点 路由编排 意图分析 执行模式 5种 策略模式 自动注册 异常恢复 统一失败收尾 补发降级 资源清理 轮次状态机 可观测性 ConversationTrace firstResponseTime debugTrace 事件 metadata ``` --- **企业级项目导航**:⬅️ [[02-前后端模块划分与调用关系|02-前后端模块划分与调用关系]] | 03-提问到返回 | ➡️ [[04-白话讲解|04-白话讲解]]