提问到返回

我们接着上一篇“异步索引构建:落库、向量化与收尾”往下走。

文档处理流水线讲完后,整个 RAG 系统的“入料阶段”就闭环了:

上传 → 解析 → 切块 → 向量化 → 索引落库
文档变成可被检索的向量 + 关键词索引
存在 PGVector 和 Elasticsearch 里

但用户感知不到这些——用户只关心**"我问一个问题,系统怎么给我答案"**。

这一篇就是把这个"使用阶段"彻底拆开来看。学完后你会理解:

1. Controller 为什么返回 Flux<String> 而不是普通 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 释放租约"顺序闭环。

整体流程图

flowchart TD
    A[前端 POST /api/chat/stream] --> B[Controller.stream<br/>返回 Flux String]
    B --> C[Flux.defer 延迟启动]
    C --> D{前端订阅?}
    D -->|是| E[buildLaunchPlan<br/>规范化参数]
    E --> F[claimConversationLease<br/>抢 Redis 租约]
    F -->|失败| G[rejectionFlux 拒绝流]
    F -->|成功| H[bootstrapConversation]
    H --> H1[startExchange 创建轮次]
    H1 --> H2[createTaskInfo 组装上下文]
    H2 --> H3[chatRuntimeRegistry.register]
    H3 --> I[bindClientChannel<br/>doOnSubscribe/doOnCancel]
    I --> J[activateGeneration]
    J --> J1[startLeaseRenewal 续期]
    J1 --> J2[buildConversationExecution]
    J2 --> K[发送 thinking 事件]
    K --> L[prepareExecutionPlan<br/>编排器]
    L --> M[ExecutorRegistry.get<br/>策略路由]
    M --> N[executor.execute<br/>模型生成]
    N --> O[emitModelChunk<br/>逐块推送]
    O --> P{完成?}
    P -->|正常| Q[finishSuccessfully]
    P -->|异常| R[finishWithFailure]
    Q --> Q1[CAS 抢收尾权]
    Q1 --> Q2[补发引用 + 推荐]
    Q2 --> Q3[关闭 SSE 流]
    Q3 --> Q4[落库 COMPLETED]
    Q4 --> Q5[cleanup 释放资源]

整条链路分成三大块:

准备阶段:Controller → defer → 启动计划 → 抢租约 → bootstrap → 绑通道
生成阶段:激活 → 编排 → 选执行器 → 模型输出 → 推送
收尾阶段:CAS → 补发事件 → 关闭流 → 落库 → 清理

二、入口:Controller 返回 Flux

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

1. Flux 而非 ResponseEntity 的设计

传统 Controller 长这样:

public ResponseEntity<ApiResponse<String>> chat(@RequestBody ChatRequestDto dto) {
    String answer = service.chat(dto);
    return ResponseEntity.ok(ApiResponse.success(answer));
}

但聊天场景下这种写法根本走不通

模型生成一个完整回答可能要 30 秒
同步等 30 秒前端用户会以为系统死了
即使加 loading 动画,用户体验也很差

流式输出的本质是把"等结果"变成"持续推送"

模型每生成一个字就推一次
前端用户看到"边打字边显示"
即使总耗时还是 30 秒,感知延迟降到 1-2 秒

2. text/event-stream 内容类型

produces = "text/event-stream;charset=UTF-8"

这是 SSE(Server-Sent Events) 协议的标准 MIME 类型:

HTTP 响应不会断开,持续保持连接
后端 push 一段数据,前端的 EventSource 立刻触发回调
比 WebSocket 简单(单向通信就够,不需要双向)
比 long polling 高效(不需要反复重连)

3. @Valid 在入口层挡掉非法请求

public Flux<String> stream(@Valid @RequestBody ChatRequestDto dto)

@Valid 会触发 Bean Validation:

ChatRequestDto:
    @NotBlank question
    @NotBlank chatMode

如果 question 为空,连 Controller 方法体都进不去,直接 400 拒绝。这是最外层防御

非法请求不消耗任何业务资源(不创建轮次、不抢租约、不调模型)
对系统是最大的保护

4. 不包装 ApiResponse 的选择

// 注释里有这句话:
// 这里不再额外包装 ApiResponse,而是直接把服务层生成的 SSE 事件流返回给前端逐段消费

正常 REST 接口都包装成 {"code": 0, "data": ..., "msg": "..."},这里不包装。原因:

SSE 不是单次响应,是事件流
事件结构由 StreamEventWriter 统一规定(type + content + timestamp + metadata)
再包一层 ApiResponse 反而打破事件流的标准结构

错误处理通过 error 类型的 SSE 事件 推给前端,而不是 HTTP 状态码或包装结构。这是流式接口的设计哲学:协议规范由事件流自己定义,不沿用 REST 包装。


三、Flux.defer:延迟启动的真实意图

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

1. 没有 defer 会发生什么?

设想去掉 defer:

public Flux<String> openConversationStream(ChatRequestDto request) {
    return openDeferredConversationStream(request);  // 直接返回
}

openDeferredConversationStream 内部会做:

1. buildLaunchPlan(创建对象)
2. claimConversationLease(占 Redis 锁)
3. bootstrapConversation(创建数据库轮次)
4. bindClientChannel(创建 Flux)

问题是:这些副作用在 Controller return 之前就执行了

Controller 把 Flux 对象 return 给 Spring WebFlux
WebFlux 框架内部还要走一系列处理:写 header、绑定 sink 等
这期间客户端可能因为各种原因(网络抖、客户端崩、超时)无法订阅
但租约已经占了、轮次已经创建了 → 资源泄漏

2. defer 的语义

没 defer:Flux 对象创建时就执行内部逻辑(eager)
有 defer:Flux 对象创建时只是占位,订阅时才执行内部逻辑(lazy)

类比 JavaScript 的 Promise:

new Promise(executor) → executor 立刻执行
Flux.defer(supplier) → supplier 在订阅时才执行

3. 为什么这种延迟特别重要?

普通业务 API 没有这个问题:

SELECT * FROM user WHERE id = 1
即使客户端没收到,这个 SELECT 也没有副作用
重新调一次也无害

但聊天 API 有强副作用

抢 Redis 锁:占用了就别人抢不到
创建数据库轮次:留下记录
注册到运行态:占内存

如果客户端没真正订阅,这些副作用都浪费了。defer 让所有副作用绑定到订阅事件

订阅成功 → 执行副作用 → 后续业务推进
订阅失败 → 副作用根本没发生 → 系统状态干净

这是响应式编程的核心思想:把副作用推迟到最晚的合适时机,让系统状态可控。


四、buildLaunchPlan:外部参数 → 内部数据结构

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:可选参数的优雅处理

private String normalizeConversationId(String conversationId) {
    if (StrUtil.isNotBlank(conversationId)) {
        return conversationId.trim();
    }
    return UUID.randomUUID().toString().replace("-", "");
}
传与不传的语义
传了 conversationId:继续已有会话
不传 conversationId:开启新会话
trim 的细节
return conversationId.trim();

为什么要 trim?

前端通过 URL/cookie/localStorage 传 ID 时可能带空格
不 trim:" abc123 " 和 "abc123" 被当成两个会话
trim:统一去掉边界空格
UUID.replace("-", "") 的考量
UUID.randomUUID().toString().replace("-", "")
原始 UUID:"550e8400-e29b-41d4-a716-446655440000"(36 字符)
去掉短横线:"550e8400e29b41d4a716446655440000"(32 字符)

为什么去横线?

更短(节省存储)
URL 安全(没有特殊字符)
作为 Redis key 不需要转义
"-" 在日志里容易和别的"-"混淆

是种轻量惯例——细节虽小但贯穿全系统会让代码更整齐。

2. parseRequiredChatMode:字符串 → 枚举的强转

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 拆成两个方法
parseOptional:用于列表查询(可不传)
parseRequired:用于聊天请求(必须传)

把"是否必填"的语义放到方法名而不是参数里,调用方一眼能看出意图。

"ALL" 的特殊处理
if (StrUtil.isBlank(value) || "ALL".equalsIgnoreCase(value.trim())) {
    return null;
}

为什么 ALL 等同于空?

列表查询场景:"查所有聊天模式" 用 ?chatMode=ALL 比 ?chatMode= 更直观
内部统一用 null 表示"不过滤"
ALL 是用户友好的别名,内部归一化

这是API 设计内部模型的解耦。

toUpperCase 大小写不敏感
ChatQueryMode.valueOf(value.trim().toUpperCase())

前端传 "open_chat" / "Open_Chat" / "OPEN_CHAT" 都能识别。这是前端容错——前端开发可能记不准枚举大小写。

异常包装
catch (IllegalArgumentException exception) {
    throw new IllegalArgumentException("chatMode 非法: " + value, exception);
}

valueOf 抛的异常信息是 "No enum constant ChatQueryMode.xxx",不够友好。包装成 "chatMode 非法: xxx" 让前端能直接展示给用户。

3. resolveSelectedDocument:模式相关的复杂校验

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 的规则一样?

OPEN_CHAT:不查任何文档,纯 Agent 模式
AUTO_DOCUMENT:由系统自动判断查哪个文档(可能多个,可能不查)

两者共同点:文档由"系统决定"而非"用户指定"
传 selectedDocumentId 反而是干扰

如果用户传了,说明用户对模式语义理解错了。严格拒绝比"忽略不处理"更好:

忽略:用户以为传了有效,实际系统当没看见,造成困惑
拒绝:明确告知"这个模式下不要传",用户立刻知道用错了
"当前可检索"的二次校验
documentKnowledgeService.listRetrievableDocuments().stream()
    .filter(item -> Objects.equals(item.getDocumentId(), resolvedDocumentId))
    .findFirst()
    .orElseThrow(...)

为什么不直接 documentMapper.selectById(documentId)

selectById:看文档是否存在
listRetrievableDocuments:看文档是否"现在可被检索"

可能场景:
    文档存在(数据库有记录) → 但 indexStatus = BUILD_FAILED → 不可检索
    文档存在 → 但被用户下线 → 不可检索
    文档存在 → 但索引在重建中 → 暂时不可检索

用 listRetrievableDocuments 把"业务可用"语义封装在一处,避免每次都重新写过滤条件。

4. 时间锚点的注入

LocalDate currentDate = LocalDate.now(CHAT_ZONE_ID);
String currentDateText = formatCurrentDate(currentDate);

这两个字段为什么要在启动计划里就准备好?

启动计划是不可变快照
后续 buildAgentQuestion / 编排器 / 工具调用都用同一个时间
避免不同环节算出不同的"now"(虽然差几毫秒但语义上要一致)
CHAT_ZONE_ID 显式指定时区
LocalDate.now(CHAT_ZONE_ID)  // 通常是 Asia/Shanghai

为什么不用 LocalDate.now()

LocalDate.now():使用 JVM 默认时区
JVM 时区受部署环境影响(Docker 容器、操作系统设置)
服务部署到不同节点可能"今天"不一样

显式指定 Asia/Shanghai 让"今天"在所有节点一致。这是全局一致性的小细节。


五、Redis 分布式租约

1. 抢占

private boolean claimConversationLease(StreamLaunchPlan launchPlan) {
    return redisLeaseManager.acquire(
        launchPlan.getLeaseKey(),
        launchPlan.getLeaseOwnerToken(),
        CHAT_RUNNING_LEASE_TTL
    );
}
三个参数的设计
leaseKey:"chat:running:{conversationId}"
    按 conversationId 加锁
    不同会话彼此不影响

leaseOwnerToken:UUID
    本次请求的"身份证"
    续期 / 释放时要带上验证
    防止误释放别人的锁

TTL:30 秒
    防止节点宕机后锁永久残留
    需要持续续期保持有效
Token 验证防误释放

考虑场景:

节点 A 抢到锁(token=T1)
节点 A 卡住,任务超过 30 秒
锁过期,节点 B 抢到锁(token=T2)
节点 A 恢复,执行 release(leaseKey)
    无 token 验证:释放了 B 的锁!
    有 token 验证:发现 token 不匹配,拒绝释放

加 token 验证防止误释放,这是 Redis 分布式锁的标准模式(Redlock 思想)。

2. 续期

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 秒续期的取舍
为什么 TTL 不直接设 1 小时,免去续期?
    节点宕机 → 锁要等 1 小时才释放 → 用户 1 小时都没法重发

为什么续期间隔不是 25 秒,接近 TTL?
    任何延迟(GC、网络抖)都可能错过续期
    25 秒间隔 + 30 秒 TTL → 5 秒安全余量太小
    10 秒间隔 + 30 秒 TTL → 20 秒余量更稳

为什么续期不是 1 秒一次?
    续期太频繁 → Redis 压力大
    10 秒已经足够频繁

这是 TTL = 3 × 续期间隔 的经验法则:

允许最多 2 次续期失败仍不至于让锁过期
给 Redis 抖动、网络延迟留余量
续期失败立刻停止
if (!renewed) {
    log.warn("会话租约续期失败,准备停止当前会话...");
    leaseRenewalDisposable.dispose();
    stopTask(taskInfo, "会话租约已失效,已停止生成");
}

为什么不让任务继续跑?

租约失效有两种原因:
    1. Redis 故障 → 没法判断别人是否抢了
    2. 真有别人抢了 → 必须立刻停止避免双写

保守做法:不管哪种原因,都停止当前任务
代价:用户可能需要重发(但只是这一次)
收益:避免数据混乱(可能影响多个会话)

这是"宁可错杀不可放过"的安全策略——租约的本质是互斥,互斥一旦不确定就要回到安全态。

续期定时器本身要先停
if (leaseRenewalDisposable != null && !leaseRenewalDisposable.isDisposed()) {
    leaseRenewalDisposable.dispose();
}
stopTask(taskInfo, "...");

为什么先 dispose 定时器再 stopTask?

定时器还在运行 → 10 秒后又触发 renewLeaseOrStop → 又检测续期失败 → 又调 stopTask
循环触发,日志大量重复

先停定时器切断重复触发的源头,再执行真正的停止逻辑。

3. 租约的整个生命周期

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:三步联动

exchangeView = conversationArchiveStore.startExchange(...);
TaskInfo taskInfo = createTaskInfo(launchPlan, exchangeView);
if (!chatRuntimeRegistry.register(taskInfo)) {
    failBootstrappedExchange(...);
    releaseLeaseQuietly(...);
    return BootstrapResult.rejected(...);
}
return BootstrapResult.ready(bindClientChannel(taskInfo));

1. 三步的强关联

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

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

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

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

注释里有一段话:"极端情况下,即使抢到租约,也可能在运行态注册时发现已有同会话任务占用,必须补偿性收尾"。

设想这个场景:

节点 A 服务请求 1(conversationId=X)
    抢到 Redis 租约
    创建轮次,注册 taskInfo
    任务执行中

节点 B 服务请求 2(conversationId=X)
    Redis 租约已被 A 占,理论上抢不到

但 Redis 主从延迟时...
    A 占的锁还没同步到从库
    B 访问从库 → 看到锁不存在 → 也"抢到"
    两个节点都以为自己拿到了锁

A 的 chatRuntimeRegistry 已经注册
B 的 chatRuntimeRegistry 在注册时发现同 cid 已存在 → 拒绝

chatRuntimeRegistry节点本地的内存注册表,对同节点上的并发做最终防御。这是双重防御

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

3. 失败时的补偿性收尾

failBootstrappedExchange(...);  // 把已创建的轮次标记为失败
releaseLeaseQuietly(...);        // 释放租约

注意顺序:

1. 先标记轮次失败(数据库状态推进)
2. 再释放租约(允许后续请求重发)

为什么这个顺序?

如果先释放租约:
    锁刚释放,后续请求进来,看到失败的轮次也开始建新轮次
    两个轮次叠加(虽然不会数据冲突,但污染数据)

先标记失败:
    确保数据库状态收敛
    再释放锁,下一次进来面对的是干净状态

4. try-catch 兜底

catch (RuntimeException exception) {
    releaseLeaseQuietly(launchPlan.getLeaseKey(), launchPlan.getLeaseOwnerToken());
    if (exchangeView != null) {
        failBootstrappedExchange(...);
    }
    return BootstrapResult.rejected(buildErrorMessage(exception));
}

bootstrap 三步任意一步抛异常都进 catch:

1. 释放租约(已经抢到的话)
2. 如果 exchange 已创建,标记为失败
3. 返回拒绝

不让任何资源残留。这种 catch 写法的关键:

exchangeView 用外部声明的局部变量(不在 try 里声明)
catch 块能访问到它来判断是否要清理
如果声明在 try 内部,catch 拿不到

这是 Java 中try-catch 跨作用域的常见技巧。


七、TaskInfo:运行时"万能上下文"

1. 字段分类

按用途分组:

身份标识(不可变):
    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. 线程安全的精确选择

不同字段用不同的线程安全机制:

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<String>:
    Reactor 提供的响应式通道
    内部已经处理并发
为什么不全用一种?
StringBuffer 用于答案累积:append 频繁,数据量大
synchronizedList 用于步骤记录:add 不频繁
ConcurrentHashMap 用于工具集合:add 可能并发
AtomicBoolean/Long 用于简单状态:轻量 CAS

每种场景选最匹配的工具,而不是一刀切

3. RunnableConfig.context() 共享 Map

runnableConfig.context().put(ChatContextKeys.EVENT_SINK, sink);
runnableConfig.context().put(ChatContextKeys.QUESTION, launchPlan.getQuestion());
...
设计意图

这是个线程上下文容器,让深层组件能拿到运行态:

检索器需要 traceId 来记录召回过程
工具调用需要 sink 来推 thinking 事件
模型层需要 currentDate 来回答时间问题

如果不用 context,方法签名会爆炸:

没 context:
    retrieve(question, traceId, sink, references, currentDate, ...)
    每加一个共享参数,所有方法签名都要改

用 context:
    retrieve(question, config)
    需要什么 config.context().get(KEY) 拿
类比
RunnableConfig.context ≈ Spring 的 SecurityContext
            ≈ Node.js 的 async_hooks
            ≈ Go 的 context.Context

都是同一种模式——调用链路上的隐式上下文

代价
优势:解耦,签名简洁,新增字段不破坏 API
代价:隐式依赖,代码 grep 不到谁在用
缓解:用 ChatContextKeys 集中定义所有 key,有据可查

4. initializeDebugTrace 的预占位

ChatDebugTrace debugTrace = initializeDebugTrace(null);
runnableConfig.context().put(ChatContextKeys.DEBUG_TRACE, debugTrace);

注意这里传 null 给 initializeDebugTrace。注释说:"在真正生成 executionPlan 之前,先放一个空白调试轨迹,保证链路中随时都能读取到 debugTrace"。

为什么需要空白占位?
prepareExecutionPlan 后会创建真正的 debugTrace
但 prepareExecutionPlan 本身可能调用某些组件
那些组件可能去 context 取 debugTrace
如果 context 还没放,取出来是 null,NPE
解决方案
先放一个 empty/blank debugTrace(空容器)
组件取出来不会 NPE,只是写入空容器
prepareExecutionPlan 完成后再替换为真实的

这是预创建空对象避免 NPE的常见模式(类似 Null Object Pattern)。


八、bindClientChannel:订阅即点火

private Flux<String> bindClientChannel(TaskInfo taskInfo) {
    return taskInfo.sink().asFlux()
        .doOnSubscribe(ignored -> activateGeneration(taskInfo))
        .doOnCancel(() -> stopTask(taskInfo, "客户端已取消请求"));
}

1. doOnSubscribe vs doOnNext

doOnSubscribe:订阅时触发一次
doOnNext:每个元素触发
doOnComplete:正常完成时触发
doOnError:异常时触发
doOnCancel:订阅被取消时触发

doOnSubscribe整条链路的"点火开关"——前端建立 SSE 连接的那一刻触发,触发后才执行:

startLeaseRenewal:启动续期定时器
buildConversationExecution:启动执行链路

2. doOnCancel:前端断开时的兜底

用户场景:
    用户问完问题,模型正在生成
    用户突然关掉浏览器/切到别的页面
    SSE 连接断开

没有 doOnCancel:
    模型继续生成完整答案
    数据库照常落库
    但前端没人接收,所有 chunk 被丢弃
    浪费算力和模型 token

有 doOnCancel:
    检测到客户端断开 → stopTask
    主动取消模型生成
    清理资源,释放租约

这是断线感知的重要机制。

3. 为什么不在 Service 入口就启动?

设想直接在 openConversationStream 里启动:

public Flux<String> openConversationStream(ChatRequestDto request) {
    StreamLaunchPlan plan = buildLaunchPlan(request);
    claimLease(plan);
    bootstrap(plan);
    activateGeneration(...);  // 直接启动
    return Flux<...>(plan.sink);
}

问题:

1. 没有 defer:Flux 对象一创建就启动,前端没订阅就跑
2. 没有 doOnCancel:前端断开后无感知
3. 启动逻辑和 Flux 对象耦合

Flux.defer + doOnSubscribe + doOnCancel 的组合:

defer:订阅时才进入业务逻辑
doOnSubscribe:订阅时点火生成
doOnCancel:取消时兜底清理

三者配合形成生命周期完整的响应式管道


九、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));

1. 嵌套 Flux.defer 的意图

外层 bindClientChannel 返回的 Flux 已经在 defer 里了,为什么这里又来一个 defer?

外层 defer:订阅时才点火
内层 defer:每次重新订阅都重新执行 thinking 事件 + prepareExecutionPlan

理论上单次会话不会重订阅,但 defer 把这种重入安全性当成基础保证,比依赖"不会重订阅"更稳。

2. Mono.fromCallable + subscribeOn(boundedElastic)

Mono.fromCallable(() -> prepareExecutionPlan(taskInfo))
    .subscribeOn(Schedulers.boundedElastic())
Mono.fromCallable 的作用
prepareExecutionPlan 是一个同步阻塞方法(包含数据库查询、检索、压缩等)
直接调用会阻塞当前线程
Mono.fromCallable 把它包装成响应式异步任务
subscribeOn(boundedElastic) 的作用
boundedElastic:Reactor 提供的"弹性边界"线程池
适合 IO 密集型任务(数据库、网络调用)
有上限(避免无限扩张)
空闲时回收线程
为什么不用 parallel 或 single?
Schedulers.parallel:CPU 密集型(线程数 = CPU 核心数)
    prepareExecutionPlan 主要是 IO,不适合

Schedulers.single:单线程
    所有任务串行,并发量上不去
    不适合多用户同时聊天

Schedulers.boundedElastic:IO 优化
    完美匹配"调数据库 + 检索 + 调外部"的场景

3. flatMapMany:Mono → Flux 的转换

.flatMapMany(plan -> {
    ConversationExecutor executor = conversationExecutorRegistry.get(plan.getMode());
    return executor.execute(taskInfo);  // 返回 Flux<String>
});
为什么是 flatMapMany 而不是 map?
map:T → R 一对一变换
flatMap:T → Mono<R> 异步一对一
flatMapMany:T → Flux<R> 异步一对多

prepareExecutionPlan 产出一个 plan,executor.execute 产出多个 chunk

plan → executor → [chunk1, chunk2, chunk3, ...]
一对多关系,只能用 flatMapMany
为什么不能直接拿到 executor 后再外层 .map?
.map(plan -> executor.execute(taskInfo))
返回 Flux<Flux<String>>(嵌套)
需要再 .flatMap(flux -> flux) 解嵌套
不如 .flatMapMany 直接搞定

4. publishOn vs subscribeOn

.publishOn(Schedulers.boundedElastic())
两者区别
subscribeOn:订阅链路向上传播,影响整个流的上游执行线程
publishOn:从这里往下,切换到指定线程池
这里 publishOn 的意图
模型每产生一个 chunk,doOnNext 会执行 emitModelChunk
emitModelChunk 内部:
    append 到 answerBuffer
    safeEmit 到 sink

如果 doOnNext 在模型推送的线程(可能是 LLM SDK 内部线程):
    长时间占用 LLM 线程
    可能阻塞下一个 chunk 推送

publishOn 切到 boundedElastic:
    每个 chunk 在新线程上处理
    LLM 线程立刻释放,可以推下一个 chunk

这是 Reactor 中生产者和消费者解耦的标准做法。

5. doOnNext / doOnError / doOnComplete 的三态闭环

.doOnNext(chunk -> emitModelChunk(...))   // 每个 chunk
.doOnError(error -> finishWithFailure(...)) // 异常
.doOnComplete(() -> finishSuccessfully(...)); // 正常完成

这是响应式流的三种终态

持续中:doOnNext 不断触发
正常完成:doOnComplete 触发一次,流终止
异常终止:doOnError 触发一次,流终止

doOnError 和 doOnComplete 互斥——一个流要么以 complete 结束要么以 error 结束,不会两者都触发。

6. 完整数据流图

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. 编排器是大脑

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;
}
编排器的核心职责
1. 装载会话记忆(summarizeHistory)
2. 判断时效性(requiresCurrentDateAnchoring / requiresFreshSearch)
3. 路由执行模式(REACT_AGENT / RETRIEVAL / GRAPH_ONLY / ...)
4. 必要时改写文档范围
文档范围动态修正
if (executionPlan.getSelectedDocumentId() != null
    && !Objects.equals(executionPlan.getSelectedDocumentId(), taskInfo.selectedDocumentId())) {
    conversationArchiveStore.refreshSessionScope(...);
    putContextIfNotNull(taskInfo.runnableConfig(), ChatContextKeys.SELECTED_DOCUMENT_ID, ...);
    ...
}

这是个有趣的设计——编排器可以修正用户的文档选择

什么场景下会修正?

AUTO_DOCUMENT 模式:用户没选文档
    编排器分析问题 → 决定该查 doc_5
    把 selectedDocumentId 设为 5
    需要同步更新归档和上下文

DOCUMENT 模式但选错了:
    编排器发现问题和当前文档完全不匹配
    可能切到 CLARIFICATION 模式让用户确认

修正后要做三件事:

1. refreshSessionScope:更新数据库归档的会话范围
2. 更新 runnableConfig.context:让后续工具看到新的文档 ID
3. 内存对象 taskInfo 在 setExecutionPlan 时一起更新

2. buildAgentQuestion:为什么不能裸传问题

裸问题的问题
用户问:"今天天气怎么样?"
直接传给模型:
    "今天" 对模型来说没有锚点
    模型可能用训练数据里的某个日期当"今天"
    或者拒绝回答说"我不知道今天日期"
包装后的 prompt
系统时间信息:
当前日期是 2025年5月20日(星期二),时区为 Asia/Shanghai。
当前问题包含相对时间或强时效语义。当用户提到"今天、明天、昨天..."时,
必须以这个日期为准,不要把搜索结果里的旧日期误当成今天。

当前问题需要核实最新外部事实,回答前必须优先调用联网搜索工具。
如果搜索结果里的日期与当前日期不一致,必须明确说明来源日期。

相关会话背景:
用户之前问过北京天气,本轮可能延续这一话题。

用户问题:
今天天气怎么样?
三段式 prompt 结构
1. 系统时间锚点:统一时间基准
2. 上下文摘要:让 Agent 知道之前聊了什么
3. 原始问题:用户的真实诉求

这种"约束 + 上下文 + 问题"的 prompt 结构是 LLM 应用的常见模式。

条件化 prompt
if (executionPlan.isRequiresCurrentDateAnchoring()) {
    builder.append("当前问题包含相对时间或强时效语义...");
} else {
    builder.append("当用户提到"今天、明天、昨天..."时,必须以这个日期为准。\n");
}

if (executionPlan.isRequiresFreshSearch()) {
    builder.append("当前问题需要核实最新外部事实,回答前必须优先调用联网搜索工具。\n");
    ...
}

不同问题给不同强度的约束:

普通问题(不需要时效性):
    只提一句"如果用户提到今天,以这个日期为准"

强时效问题(需要锚点):
    详细说明"不要把搜索结果里的旧日期误当成今天"
    强调日期一致性校验

强时效 + 需要联网:
    强制调用搜索工具
    要求标注来源日期
    无法找到匹配结果时要明确说明不确定性,不要编造
为什么不一刀切都给最强约束?
Prompt 越长 → token 消耗越多 → 成本越高 → 推理越慢
弱时效问题(比如"什么是机器学习")给联网搜索约束:
    模型可能莫名其妙调搜索工具
    浪费 token,降低响应速度

强时效问题(比如"今天股市")不给联网约束:
    模型用训练数据回答 → 给出过时信息 → 用户被误导

精确匹配 prompt 强度和问题特征,是 RAG 系统调优的关键之一。


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

@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) {
        ConversationExecutor executor = executorMap.get(mode);
        if (executor == null) {
            throw new IllegalStateException("未找到执行模式对应的执行器: " + mode);
        }
        return executor;
    }
}

1. Spring 自动注入 List

public ConversationExecutorRegistry(List<ConversationExecutor> executors) {

这是 Spring 的"魔法":

Spring 扫描所有实现了 ConversationExecutor 接口的 Bean
自动收集成 List 注入构造函数
新增一个 Executor 类(加 @Component) → 自动加入 List → 自动注册
对比传统注册方式
传统:
    @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 的选择

new EnumMap<>(ExecutionMode.class)
为什么不用 HashMap?
EnumMap:专门为枚举 key 优化,内部用数组实现
    O(1) 查找,内存紧凑
    比 HashMap 性能更好

HashMap:通用实现,要算 hashCode,处理冲突
    对枚举 key 来说浪费

虽然差异微小,但正确选择数据结构是工程素养。

3. mode() 方法的设计

public interface ConversationExecutor {
    ExecutionMode mode();  // 执行器自己声明处理哪种模式
    Flux<String> execute(TaskInfo taskInfo);
}

这种"自我声明"模式比外部注解或配置更灵活:

Bean 自己最清楚自己处理什么
注册表只负责"问 Bean → 收集映射"
新增模式不需要改注册表代码

4. 启动期校验

public ConversationExecutor get(ExecutionMode mode) {
    ConversationExecutor executor = executorMap.get(mode);
    if (executor == null) {
        throw new IllegalStateException("未找到执行模式对应的执行器: " + mode);
    }
    return executor;
}
抛 IllegalStateException 而非业务异常
IllegalStateException 表示"程序状态错误"
不是用户输入错误,是开发者忘了实现某个 Executor
应该让程序崩溃 + 报警 → 立刻修复
更好的做法:启动期校验

可以在 ExecutorRegistry 构造函数里加:

public ConversationExecutorRegistry(List<ConversationExecutor> 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. 整体架构

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

新增执行模式只需要:

1. 在 ExecutionMode 枚举加新值
2. 实现 ConversationExecutor 接口
3.@Component 注解
4. 重启服务
完全不用改 Registry 和 Orchestrator(除了 Orchestrator 的路由逻辑)

这是开闭原则(Open-Closed Principle)的标准实践——对扩展开放,对修改关闭。


十二、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()));
}

1. 三件事的顺序

1. append 到 answerBuffer:累积完整答案,收尾时落库用
2. 记录首包耗时:监控指标,衡量响应速度
3. 推送到前端:用户即时看到
为什么是这个顺序?
先 append 再推送:
    确保数据已经累积到 buffer
    即使后面推送失败,buffer 里也有完整数据
    收尾落库不受推送失败影响

先记首包再推送:
    时间戳更精确(贴近实际接收时间)

2. CAS 记录首包耗时

if (taskInfo.firstResponseTimeMs().get() == 0L) {
    taskInfo.firstResponseTimeMs()
        .compareAndSet(0L, System.currentTimeMillis() - taskInfo.startTime());
}
为什么先 get 判 0 再 CAS?
直接 CAS:每次 chunk 都尝试 CAS,无效操作多
    chunk 可能上千个,每次都 CAS 浪费

先 get:大多数情况快速跳过(O(1) 读)
    只有第一次 get == 0 时才进入 CAS
    显著降低开销
CAS 而非简单赋值
为什么不直接 firstResponseTimeMs.set(...)?

考虑边界:
    多个线程同时进入 emitModelChunk(理论上不会,但保险起见)
    第一个线程进入 if,准备 set
    第二个线程也进入 if,也准备 set
    两个 set 都执行 → 后者覆盖前者

用 CAS:
    只有第一个 compareAndSet 成功
    第二个失败,值不变
    保证记录的是真正的"首包"时间
首包耗时的业务价值
TTFB(Time To First Byte) = 首包耗时
是聊天 API 的核心 SLA 指标

监控告警:首包耗时 > 5 秒 → 报警
A/B 测试:对比不同模型/提示词的首包速度
用户体验:首包越快,用户感知越好

3. safeEmit 的兜底

safeEmit(taskInfo.sink(), streamEventWriter.text(chunk, ...));
safeEmit 的实现 (推测)
private void safeEmit(Sinks.Many<String> sink, String event) {
    try {
        sink.tryEmitNext(event);
    } catch (Exception e) {
        log.warn("emit 失败", e);
    }
}
为什么 emit 可能失败?
sink 已经被 dispose(任务被取消)
sink 已经 complete(任务正常结束)
背压:消费速度跟不上生产速度
为什么不让失败抛出?
emit 失败抛异常:
    doOnError 触发 → finishWithFailure
    实际上模型可能还在正常生成
    一次推送失败不应影响整个任务

safeEmit 只记 warn:
    继续累积 chunk 到 buffer
    最终落库还是完整的

这是容错设计——把"瞬时推送失败"和"任务真正失败"解耦。

4. 流式输出的体验

模型每输出一小段就立刻推:

模型生成"根据您提供的文档,Spring WebFlux 是..."
分成 chunk:
    "根据"
    "您"
    "提供的"
    "文档,"
    "Spring "
    "WebFlux "
    "是..."

每个 chunk 立刻推送 → 用户看到"打字机效果"
chunk 粒度的取舍
chunk 太大(整段才推):
    用户等很久才看到东西
    "边生成边显示"效果消失

chunk 太小(每个 token 一推):
    网络开销大(每个 SSE 事件有 header 开销)
    前端渲染压力大

通常折中:几个 token 一个 chunk
    模型 SDK 会根据网络条件自动调整

十三、StreamEventWriter:SSE 事件标准化

public String text(String content, StreamEventMetadata metadata) {
    return write(event("text", content, metadata));
}

private Map<String, Object> event(String type, Object content, StreamEventMetadata metadata) {
    Map<String, Object> 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 推荐追问 收尾时 一次
设计哲学
type 字段表明事件类别 → 前端按类型分发处理
content 字段灵活承载 → 字符串、List、Map 都能放
timestamp + metadata → 调试和追踪信息

2. LinkedHashMap 的选择

Map<String, Object> payload = new LinkedHashMap<>();
为什么不用 HashMap?
LinkedHashMap:保持插入顺序
    序列化成 JSON 时字段顺序固定
    "type" 始终在第一个位置

HashMap:不保证顺序
    每次序列化字段顺序可能不同
    日志和调试时不一致

虽然 JSON 顺序在协议层不重要,但人眼阅读和日志比对时一致顺序很有帮助。

3. 条件性放入字段

if (metadata.conversationId() != null) {
    payload.put("conversationId", metadata.conversationId());
}
if (metadata.exchangeId() != null && metadata.exchangeId() > 0) {
    payload.put("exchangeId", metadata.exchangeId());
}
为什么不直接 put(value 可能为 null)?
直接 put null:
    JSON 序列化出 "conversationId":null
    前端要判断 null,代码繁琐
    增加传输大小

条件 put:
    null 时字段根本不出现
    前端代码可以直接用 hasOwnProperty 或 ?. 安全访问

这是Java 后端 → JSON 前端的常见处理细节。

4. 标准事件结构

{
    "type": "text",
    "content": "根据文档内容,",
    "timestamp": "2025-05-20T10:30:00.123Z",
    "conversationId": "abc123",
    "exchangeId": 42
}
前端处理代码
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 步成功收尾

private void finishSuccessfully(TaskInfo taskInfo) {
    if (!taskInfo.finalized().compareAndSet(false, true)) {
        return;
    }

    String answer = taskInfo.answerBuffer().toString();
    List<SearchReference> uniqueReferences = deduplicateReferences(...);

    ConversationTraceRecorder.StageHandle finalizeStage = ...;
    ConversationTraceRecorder.StageHandle recommendationStage = ...;

    List<String> 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 抢"唯一收尾权"

if (!taskInfo.finalized().compareAndSet(false, true)) {
    return;
}
为什么需要 CAS?
成功收尾 / 失败收尾 / 主动停止 都要执行收尾逻辑
都可能在不同线程被触发:
    doOnComplete:模型完成线程
    doOnError:任意失败线程
    stopTask(用户取消):前端线程
    stopTask(续期失败):续期定时器线程

没 CAS:
    可能两个收尾同时执行
    重复落库 → 数据库唯一约束冲突
    重复释放租约 → 释放别人的锁

CAS:
    只有第一个 compareAndSet 成功
    其他线程直接 return,不重复执行
CAS 的语义
compareAndSet(expected, newValue):
    如果当前值 == expected:
        设置为 newValue,返回 true
    否则:
        不修改,返回 false

第一次调用:finalized 是 false → CAS 成功 → 设为 true
第二次调用:finalized 已经是 true → CAS 失败 → 返回 false

2. 数据快照

String answer = taskInfo.answerBuffer().toString();
List<SearchReference> uniqueReferences = deduplicateReferences(snapshotReferenceList(taskInfo.references()));
为什么要快照?
references 是线程安全 List,但仍然可能被异步线程写入
    比如有个晚到的工具结果
    收尾期间继续 add → 落库的数据和推送给前端的不一致

snapshotReferenceList:在某个时刻拷贝出来
    后续不再受外部修改影响
    保证一致性
deduplicateReferences:去重
不同检索路径可能命中同一文档
    向量检索召回 doc_5
    关键词检索也召回 doc_5
    references 里有两份 doc_5

去重:同 documentId + chunkId 算同一引用
    用户看到的是干净的引用列表
    落库的也是去重后的

3. CLARIFICATION 模式的特殊处理

if (taskInfo.executionPlan() != null
    && taskInfo.executionPlan().getMode() == ExecutionMode.CLARIFICATION) {
    recommendations = ...getClarificationOptions();
}
else {
    recommendations = recommendationService.generateRecommendations(...);
}
为什么 CLARIFICATION 复用澄清选项?
正常模式:
    问:今天天气?
    答:晴朗 25 
    推荐:明天天气? / 这周降水?
     推荐内容是"答案的延伸"

CLARIFICATION 模式:
    问:介绍一下产品
    答:你想了解哪个产品? A/B/C
    澄清选项:[A 产品介绍, B 产品介绍, C 产品介绍]

    如果再调推荐服务:
        基于"介绍一下产品"+ "你想了解哪个产品"生成
        可能产生:产品有什么? / 怎么购买? 之类
        但用户当前阶段需要的是 A/B/C 选择
        额外推荐反而干扰

直接复用澄清选项作为推荐——"问什么就推什么",不浪费一次推荐生成。

4. 补发事件 + 关闭流的顺序

1. 补发 references
2. 补发 recommendations
3. 关闭流(safeComplete)
为什么这个顺序?
关闭流后再 emit:无效(流已关闭)
所以必须先 emit 再关闭

事件之间也有顺序:
    references 在 recommendations 前:
        前端先看到"答案的依据"
        再看到"接下来可以问什么"
        逻辑上是"先回顾再展望"
safeComplete 的作用
sink.tryEmitComplete();
通知前端:本次流结束
前端 EventSource.close 触发
连接断开,资源释放

5. try-catch-finally 的精妙嵌套

try {
    // 补发事件
} catch (RuntimeException exception) {
    log.warn(...);  // 不改判失败
}
finally {
    try {
        safeComplete(...);  // 关闭流
    } catch (...) {}

    try {
        refreshDebugTraceRuntimeStats(...);
        conversationArchiveStore.completeExchange(...);  // 落库
    } catch (...) {
        log.error(...);  // 落库失败是 error 级别
    }
    finally {
        safeRefreshConversationSummary(...);  // 异步刷摘要
        cleanup(...);  // 资源清理
    }
}
三层 try-catch 的职责
最外层 try-catch:
    捕获补发事件的异常
    warn 级别,不影响主流程

中间 try-catch:
    分别保护 safeComplete / 落库
    每个独立 try,失败不影响下一个

最内层 finally:
    无论前面所有步骤成功失败,都要执行
    刷摘要 + cleanup
    保证资源最终释放
日志级别的精确选择
补发事件失败:warn
    答案已生成成功,只是补充信息没推
    用户体验不完美但不致命

关闭流失败:warn
    流可能已自动关闭,影响小

落库失败:error
    用户看到了答案,但数据库没记录
    需要紧急排查

摘要刷新失败:silent(在 safeRefresh 内部处理)
    最不重要的步骤,不该影响主流程
为什么补发失败不改判会话?

注释明确说:"正文已经成功完成,补发事件失败不应推翻'本轮回答成功生成'这一主结果"。

用户视角:
    答案已经在屏幕上看到了 
    引用没看到  体验差但不致命
    推荐没看到  体验差但不致命

技术视角:
    模型 token 已经消耗 
    数据库轮次还要落库 

如果改判失败:
    数据库状态:FAILED
    用户实际体验:成功
    数据和体验割裂,运维排查困难

6. 事件推送顺序总览

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 的彻底清理

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);
}
清理顺序
1. 停续期定时器(否则继续触发 Redis 操作)
2. 停业务执行流(确保 Reactor 资源释放)
3. 释放 Redis 租约(允许后续请求)
4. 移除运行态注册(允许同节点接受新请求)
每步的 isDisposed 检查
if (... != null && !disposable.isDisposed())
为什么要检查 isDisposed?
    可能这个 disposable 已经被某个分支主动 dispose 过
    重复 dispose 一般无害,但有些实现可能抛异常
    防御性编程,确保不出异常
releaseLeaseQuietly:静默释放
release 失败的可能原因:
    租约已过期(自动释放了)
    Redis 故障

失败也不报错:
    cleanup 是终态,不应阻塞
    租约即使没主动释放,30 秒后也会自动过期

十五、总结:整条链路的设计哲学

1. 12 步流程回顾

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. 贯穿全链路的设计原则

响应式驱动
全链路无阻塞 → Flux/Mono 串联
publishOn/subscribeOn 精确控制线程
背压由 Reactor 自动处理
延迟启动 + 副作用收敛
Flux.defer 延迟到订阅时
副作用绑定到订阅事件
未订阅 → 零资源消耗
多层防御
入口 @Valid:格式校验
buildLaunchPlan:语义校验
claimLease:Redis 跨节点互斥
chatRuntimeRegistry:节点内互斥
finalized CAS:收尾互斥
资源管理闭环
任何资源(租约/订阅/注册)都有"获取 + 释放"对
异常时通过 catch 清理已获取的资源
finally 兜底确保清理
不留悬空资源
状态可观测
ConversationTraceRecorder 记录每个阶段
StreamEventMetadata 注入每个事件
debugTrace 累积运行态统计
firstResponseTimeMs 监控指标
失败容错的层级
致命失败(如租约抢占):立即拒绝
非致命失败(如补发事件):降级处理
彻底失败(如模型异常):走 finishWithFailure
任何失败都不影响其他会话

3. 系统能力图谱

mindmap
    root((聊天链路))
        协议层
            SSE 流式
            Flux WebFlux
            事件类型 5种
        并发控制
            Redis 租约
            TTL 30s
            10s 续期
            CAS 收尾
            本地注册表
        上下文管理
            TaskInfo
            RunnableConfig
            会话记忆
            时间锚点
        路由编排
            意图分析
            执行模式 5种
            策略模式
            自动注册
        异常恢复
            统一失败收尾
            补发降级
            资源清理
            轮次状态机
        可观测性
            ConversationTrace
            firstResponseTime
            debugTrace
            事件 metadata

企业级项目导航:⬅️ 02-前后端模块划分与调用关系 | 03-提问到返回 | ➡️ 04-白话讲解