聊天系统整体架构讲解

我们从用户按下"发送"那一刻开始,把整个聊天系统的链路完整走一遍。整条链路分 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()

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

Controller 就做两件事:接参数、转交 Service。但这一个方法里藏着四个值得讲的设计点:

返回值是 Flux<String> 而不是普通 JSON。聊天回答是流式生成的——模型每产出一小段文字就立刻推给前端,用户看到"边想边写"的效果,感知延迟从 30 秒降到 1-2 秒。如果写成 ResponseEntity<ApiResponse<String>> 同步等结果,模型生成 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()

public Flux<String> 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 设计的常用原则。

最后还会注入两个时间字段——currentDatecurrentDateText,作为后续处理"今天/最新/本周"等相对时间表达的统一基准。这里特意用 LocalDate.now(CHAT_ZONE_ID) 显式指定 Asia/Shanghai 时区——不用默认时区是因为 JVM 默认时区受部署环境影响(Docker、操作系统设置),不同节点的"今天"可能不一样,显式指定让全局一致

整个 buildLaunchPlan 的产出是一个不可变快照,所有字段都 final——后续所有环节都基于这个快照工作,外部参数从此和内部状态彻底解耦。在入口处就把非法参数全拦住了,后续的编排器和执行器就可以放心使用,不用再做重复校验。

步骤 4:Redis 分布式租约抢占

代码位置BusinessChatService.claimConversationLease()

private boolean claimConversationLease(StreamLaunchPlan launchPlan) {
    return redisLeaseManager.acquire(
        launchPlan.getLeaseKey(),        // chat:running:{conversationId}
        launchPlan.getLeaseOwnerToken(), // 本次请求的唯一 token (UUID)
        CHAT_RUNNING_LEASE_TTL           // 30 秒
    );
}

这一步保证同一个会话在集群任意时刻只有一个生成任务在运行。集群部署下用户快速点两次发送或前端重试,没有租约的话两个生成任务可能在不同节点上同时跑——同一个 conversationId 下两条轮次记录、两份答案文本、两份引用,前端状态机彻底乱套。

三个参数的设计都有讲究

  • leaseKeychat: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 做三件强关联的事,必须顺序执行且全部成功

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 的经典应用

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

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()

private Flux<String> 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

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<R> 异步一对多。prepareExecutionPlan 产出一个 plan,executor.execute 产出多个 chunk,是一对多的关系。如果用 map 会拿到 Flux<Flux<String>> 嵌套 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 模式下用户没选文档,编排器分析问题后决定该查哪个文档。修正发生时还要做善后:

if (executionPlan.getSelectedDocumentId() != null
    && !Objects.equals(executionPlan.getSelectedDocumentId(), taskInfo.selectedDocumentId())) {
    conversationArchiveStore.refreshSessionScope(...);  // 同步归档
    putContextIfNotNull(taskInfo.runnableConfig(), ChatContextKeys.SELECTED_DOCUMENT_ID, ...);  // 同步上下文
    ...
}

不光更新执行计划本身,还要同步刷新归档的会话范围 + RunnableConfig 上下文——三个数据源(执行计划、归档、运行上下文)保持一致,防止后续组件读到旧值。

buildAgentQuestion:不要让 Agent 裸接用户问题

最后 buildAgentQuestion 把时间锚点 + 时效约束 + 历史摘要 + 原始问题拼成最终 prompt,结构大致是:

系统时间信息:
当前日期是 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) 查找找到对应执行器:

@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);
        }
    }
    ...
}

这一段代码非常短但有三个值得讲的设计点:

EnumMap 不是 HashMap——枚举做 key 时 EnumMap 是数组实现、查找 O(1) 且常数极小、内存紧凑,是 Java 集合里的小优化。

Spring 自动注入 List<ConversationExecutor> 是 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 同时做三件事:

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 嵌套

这是这段代码最值得讲的工程细节:

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

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 统一走失败收尾,保证租约释放

整体哲学是"主结果优先"——回答生成成功是核心目标,附加信息(引用、推荐、调试追踪、归档落库)失败都不能推翻这个主结果。

状态转换

轮次状态CREATEDCOMPLETED / FAILED / STOPPED,一旦进入终态不可逆转。

租约状态ACQUIREDRENEWINGRELEASED / EXPIRED。正常流程是 ACQUIRED → RENEWING → RELEASED(cleanup 时主动释放);异常流程是续期失败后 EXPIRED 触发自动停止;极端情况是节点宕机后 30 秒 TTL 自动 EXPIRED。

TaskInfo.finalizedfalsetrue(CAS 单向跳变),保证三种收尾路径只有一种生效。

会话记忆体系

最后提一下支撑整条链路的记忆体系,它分三层,每层应对不同的诉求:

  • 长期摘要:通过 LLM 压缩历史对话生成的摘要,注入 AgentQuestion 前缀。作用是"远记忆"——让 Agent 知道之前聊了什么主题但不占太多 token
  • 最近对话窗口:保留最近几轮的原文。作用是"近记忆"——保证短期上下文准确性,问题改写时参照最近表达
  • Checkpoint:ReAct Agent 的状态快照,支持多轮工具调用间的状态传递。作用是"过程记忆"——单次 Agent 执行内的中间状态

为什么要分三层?长期摘要太粗(细节丢了)、最近窗口太窄(远的看不到)、Checkpoint 只在单次 Agent 内有效(跨轮次没用)——任何单一记忆机制都不够,组合起来才覆盖全场景

每一轮收尾时 safeRefreshConversationSummary 异步刷新长期摘要——把刚结束的轮次纳入


企业级项目导航:⬅️ 05-龙虾OpenClaw深度解析:从对话到执行的跨越 | 00-聊天系统整体架构讲解 | ➡️ 01-文档上传到RAG检索完整链路讲解