白话讲解

一、把这条链路放在系统里的位置

讲这条链路之前先把它和上一段衔接一下。前面我们花了好几大段讲"入料阶段"——文档怎么从原始文件变成可被检索的结构化资产:上传、解析、切块、向量化、索引落库。入料阶段的终点是 PGVector 和 ES 里躺着一份份准备好被检索的 chunk

这一段是"使用阶段"——用户终于来了,对着输入框敲了一句话点了发送。从那一刻开始,到他屏幕上一个字一个字蹦出答案,再到答案完整收尾、引用展示、推荐追问出现——这中间发生的所有事,是这条链路的全部内容。

讲这条链路之前先立一个核心认知:这是一条响应式 + 流式的链路,不是传统的 REST 请求。所有设计都是围绕"流"展开的——返回值是 Flux 不是 JSON、副作用要绑定到订阅事件、收尾要保证只触发一次。如果按传统 REST 接口的思维去看,会处处觉得别扭——但只要把"流"这个核心抓住,后面所有奇怪的设计都能讲通。

整条链路可以切成三大块:

准备阶段: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():

@PostMapping(value = "/stream", produces = "text/event-stream;charset=UTF-8")
public Flux<String> stream(@Valid @RequestBody ChatRequestDto dto) {
    return businessChatService.openConversationStream(dto);
}

这一个方法就藏了四个值得单独讲的设计点。

第一,返回 Flux 不返回 JSON。传统 REST 接口里这里会写 ResponseEntity<ApiResponse<String>>——同步等结果、统一包装。但聊天场景这么写直接走不通——一个完整回答可能要 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:

public Flux<String> openConversationStream(ChatRequestDto request) {
    return Flux.defer(() -> openDeferredConversationStream(request));
}

这一层 defer 包装非常关键,但它的真实意图新手很容易看不出来。

先讲没有 defer 会怎样。直接 return openDeferredConversationStream(request) 的话,这个方法内部要做的事——buildLaunchPlan、claimConversationLease、bootstrapConversation、bindClientChannel——会在 Controller return 之前 全部执行。问题来了:

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 抢占的三个参数

redisLeaseManager.acquire(leaseKey, leaseOwnerToken, CHAT_RUNNING_LEASE_TTL);
  • leaseKeychat:running:{conversationId} —— 按会话粒度加锁,不同会话彼此不影响。
  • leaseOwnerToken:UUID —— 本次请求对租约的"身份证",续期和释放时要带上做验证。
  • TTL:30 秒 —— 防止节点宕机后锁永久残留。

leaseOwnerToken 防误释放这一点特别要讲。考虑这个场景:

节点 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 秒不够用怎么办?模型生成可能一两分钟。所以要续期:

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 秒已经足够频繁。

续期失败立刻停止任务——这是租约机制最有意思的一个设计:

if (!renewed) {
    leaseRenewalDisposable.dispose();   // 先停定时器自身,避免重复触发
    stopTask(taskInfo, "会话租约已失效,已停止生成");
}

为什么不让任务继续跑?租约失效有两种原因——Redis 故障(没法判断别人是否抢了)、真有别人抢了(必须停止避免双写)。保守做法是不管哪种原因都停——代价是用户可能要重发,收益是避免数据混乱。互斥一旦不确定就要回到安全态——这是分布式系统里的安全策略。

注意这里先 dispose 定时器再 stopTask 的顺序——不先停定时器的话,10 秒后又触发 renewLeaseOrStop、又检测续期失败、又调 stopTask、循环触发、日志大量重复。先切断重复触发的源头再执行真正的停止逻辑

七、bootstrapConversation:三步联动

抢到租约之后进入 bootstrap,要做三件强关联的事:

1. startExchange:数据库创建轮次记录 → 拿到 exchangeId
2. createTaskInfo:内存组装运行时上下文 → 用到 exchangeId
3. register:把 taskInfo 注册到运行态 → 用 conversationId 做 key

每一步都依赖前面的结果,必须顺序执行且全部成功才算 bootstrap 完成。

7.1 双重防御:为什么抢到租约还可能注册失败

if (!chatRuntimeRegistry.register(taskInfo)) {
    failBootstrappedExchange(...);
    releaseLeaseQuietly(...);
    return BootstrapResult.rejected("该会话当前正在执行中,请稍后再试");
}

新手看到这里会问——租约都抢到了,注册怎么还会失败?设想这个场景:

节点 A 服务请求 1(conversationId=X) → 抢到 Redis 锁 → 任务执行中
节点 B 服务请求 2(conversationId=X) → Redis 主从延迟时...
    A 占的锁还没同步到从库
    B 访问从库 → 看到锁不存在 → 也"抢到"
    两个节点都以为自己拿到了锁

这时候本地的 chatRuntimeRegistry 就站出来做最后防御。chatRuntimeRegistry 是节点本地的内存注册表,对同节点上的并发做最终防御。整个互斥体系是双重的:

  • Redis 租约:跨节点互斥(应对正常情况)
  • 本地注册表:节点内互斥(应对 Redis 短暂不一致)

7.2 失败时的补偿性收尾

注册失败时要做两件事,顺序很关键

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:

return taskInfo.sink().asFlux()
    .doOnSubscribe(ignored -> activateGeneration(taskInfo))
    .doOnCancel(() -> stopTask(taskInfo, "客户端已取消请求"));

doOnSubscribe 是整条链路的"点火开关"——前端建立 SSE 连接的那一刻,后端才真正开始干活。这是 Flux.defer 设计意图的最终落点——所有副作用都绑在订阅事件上。

doOnCancel 处理前端断开——用户关掉浏览器、网络断了,前端不再订阅了,后端立刻调 stopTask 中断生成、不让模型空跑浪费 token。

activateGeneration 做两件事:启动租约续期定时器、启动执行链路订阅。两个 Disposable 都存到 TaskInfo 上,停止时统一 dispose。


九、buildConversationExecution:把整条响应式流串起来

这是整条链路最核心的一段代码——把"分析中提示 → 编排执行计划 → 选执行器 → 消费 chunk → 收尾"串成一条响应式流:

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<String>> 嵌套 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 的职责:

系统时间信息:
当前日期是 2025-05-20,时区为 Asia/Shanghai。
[条件化时效约束]

[历史摘要(如有)]

用户问题:
[原始问题]

为什么要这么包装?裸传问题会出问题——用户问"今天天气怎么样"、不告诉 Agent 今天是几号,它可能用训练数据里的旧日期回答;用户问"还记得我刚才说的吗"、不给历史,Agent 直接懵。

包装的核心是条件化注入——requiresCurrentDateAnchoring 时给更严格的日期约束、requiresFreshSearch 时强制要求联网搜索并标注来源日期、有 historySummary 才注入历史。不同问题有不同的 prompt 强度,精确匹配


十一、ConversationExecutorRegistry:策略模式 + 自动发现

执行计划产出后,根据 mode 选执行器:

@Component
public class ConversationExecutorRegistry {
    private final Map<ExecutionMode, ConversationExecutor> executorMap = new EnumMap<>(ExecutionMode.class);

    public ConversationExecutorRegistry(List<ConversationExecutor> 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:

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 步顺序不能乱:

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:彻底释放资源

收尾的最后一步:

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-提问到返回 | 04-白话讲解 | ➡️ 00-总览