--- title: "01-从用户提问到答案返回的总流程" created: 2026-05-21 aliases: - 从用户提问到答案返回的总流程 tags: - 项目 --- # 从用户提问到答案返回的总流程 这篇文档会把 Super Agent 聊天系统后端的完整链路拆开来讲,从用户点击"发送"的那一刻开始,一直到答案流式输出到前端、最后落库收尾为止。每个关键步骤都会贴出对应的源码,加上注释说明它在整条链路里的作用。 看完这篇,你会对"一个问题是怎么从 Controller 一路走到模型输出再回到前端"有一个完整的认知。 ## 总流程概览 先看一张全局流程图,对整条链路有个直观印象: ![[Fv6c6Yms0QdCgnIwQwnETI8VPGQJ-2202ad22.png]] 接下来我们按照这张图的顺序,逐步拆解每个阶段的源码。 ## 入口:Controller 接收请求 一切从前端的 POST 请求开始。前端把用户的问题、会话 ID、聊天模式等信息打包成 `ChatRequestDto`,发到 `/api/chat/stream` 接口。 先看请求参数长什么样: ```java public class ChatRequestDto { @NotBlank(message = "question 不能为空") private String question; // 用户输入的问题 private String conversationId; // 会话 ID,不传则自动生成新会话 @NotBlank(message = "chatMode 不能为空") private String chatMode; // 聊天模式:OPEN_CHAT / AUTO_DOCUMENT / DOCUMENT private String selectedDocumentId; // 当前文档问答模式下,用户选择的文档 ID } ``` Controller 就做两件事:接参数、转交给 Service: ```java @AllArgsConstructor @RestController @RequestMapping("/api/chat") public class BusinessChatController { private final BusinessChatService businessChatService; /** * 打开一个流式会话。 *

* 该接口返回的是 SSE 文本流,前端可以持续接收“思考中”、“正文增量”、“引用”、“推荐追问”等事件。 *

* * @param dto 前端提交的聊天请求,包含问题、会话 ID、聊天模式、选中文档等信息 * @return SSE 字符串流,内容由服务层按事件格式持续输出 */ @PostMapping(value = "/stream", produces = "text/event-stream;charset=UTF-8") public Flux stream(@Valid @RequestBody ChatRequestDto dto) { // 这里不再额外包装 ApiResponse,而是直接把服务层生成的 SSE 事件流返回给前端逐段消费。 return businessChatService.openConversationStream(dto); } } ``` > 为什么返回 Flux 而不是普通 JSON? > > 因为聊天回答是流式生成的,模型每产出一小段文字就立刻推给前端,用户能看到"边想边写"的效果。这里用的是 Spring WebFlux 的 `Flux`,配合 `text/event-stream` 内容类型,实现了 SSE(Server-Sent Events)协议。 ## 延迟启动:Flux.defer 的设计意图 Controller 调用的 `openConversationStream()` 并不会立刻开始干活,而是用 `Flux.defer` 包了一层: ```java public Flux openConversationStream(ChatRequestDto request) { // defer 的作用:把真正的启动逻辑延后到"客户端真正订阅流"的那一刻 // 避免只是创建 Flux 对象时就提前占用租约、创建轮次 return Flux.defer(() -> openDeferredConversationStream(request)); } ``` 这个设计很关键——如果不用 `defer`,Flux 对象一创建就会执行内部逻辑,但这时候前端可能还没准备好接收数据。用了 `defer` 之后,只有前端真正建立 SSE 连接(订阅 Flux)时,后端才会开始抢租约、创建轮次这些操作。 ## 构建启动计划:buildLaunchPlan 真正的启动逻辑在 `openDeferredConversationStream()` 里,第一步就是把前端传来的参数规范化,转成内部使用的 `StreamLaunchPlan` 对象: ```java /** * 把外部请求转换成内部启动计划。 *

* 这一步会完成问题与会话 ID 规范化、聊天模式解析、所选文档校验,以及时间锚点准备。 *

*/ private StreamLaunchPlan buildLaunchPlan(ChatRequestDto request) { // 先校验并规整用户问题,确保下游不会处理空白问题。 String question = normalizeQuestion(request.getQuestion()); // conversationId 允许前端不传;如果不传则为新会话自动生成一个稳定 ID。 String conversationId = normalizeConversationId(request.getConversationId()); ChatQueryMode chatMode = parseRequiredChatMode(request.getChatMode()); // 在当前文档问答模式下,这里会校验 selectedDocumentId 是否合法、是否可检索。 KnowledgeDocumentDescriptor selectedDocument = resolveSelectedDocument(chatMode, request.getSelectedDocumentId()); // 当前日期会被写入 prompt 和上下文,作为处理“今天/最新/本周”等相对时效语义的统一基准。 LocalDate currentDate = LocalDate.now(CHAT_ZONE_ID); String currentDateText = formatCurrentDate(currentDate); return new StreamLaunchPlan( question, conversationId, chatMode, selectedDocument == null ? null : selectedDocument.getDocumentId(), selectedDocument == null ? "" : selectedDocument.getDocumentName(), selectedDocument == null ? null : selectedDocument.getLastIndexTaskId(), // 每个会话共用一个运行租约键,用来防止并发生成。 buildChatLeaseKey(conversationId), // ownerToken 代表本次请求对租约的“所有权”,续期和释放时都靠它校验。 UUID.randomUUID().toString(), currentDate, currentDateText ); } ``` 这一步做的事情不复杂,但很重要——它把外部不可控的前端参数,转换成了内部稳定、可信赖的数据结构。后续所有环节都基于这个 `StreamLaunchPlan` 来工作。 我们展开看看里面几个关键的子方法。 ### normalizeConversationId:会话 ID 规范化 ```java // BusinessChatService.java —— 规范化 conversationId private String normalizeConversationId(String conversationId) { // 前端传了就直接用(去掉首尾空格) if (StrUtil.isNotBlank(conversationId)) { return conversationId.trim(); } // 没传就自动生成一个 UUID 作为新会话的 ID return UUID.randomUUID().toString().replace("-", ""); } ``` 这个设计让前端可以灵活控制:传了 `conversationId` 就是继续已有会话,不传就是开启新会话。 ### parseRequiredChatMode:聊天模式解析 ```java // BusinessChatService.java —— 解析聊天模式 private ChatQueryMode parseRequiredChatMode(String value) { ChatQueryMode chatMode = parseOptionalChatMode(value); if (chatMode == null) { throw new IllegalArgumentException("chatMode 不能为空"); } return chatMode; } private ChatQueryMode parseOptionalChatMode(String value) { // 空值或 "ALL" 表示不过滤(用于列表查询场景) 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); } } ``` 前端传的是字符串(比如 `"OPEN_CHAT"`),这里负责转成枚举。如果传了个不认识的值,直接抛异常拒绝,不会让非法模式流入后续链路。 ### resolveSelectedDocument:文档校验 这个方法根据聊天模式来校验 `selectedDocumentId` 是否合法,不同模式有不同的规则: ```java // BusinessChatService.java —— 校验所选文档 private KnowledgeDocumentDescriptor resolveSelectedDocument(ChatQueryMode chatMode, String selectedDocumentId) { String normalizedDocumentId = StrUtil.trimToNull(selectedDocumentId); if (chatMode == ChatQueryMode.OPEN_CHAT) { // 开放问答模式不绑定文档,传了 selectedDocumentId 就报错 if (normalizedDocumentId != null) { throw new IllegalArgumentException("开放式提问模式下不能传 selectedDocumentId"); } return null; } if (chatMode == ChatQueryMode.AUTO_DOCUMENT) { // 自动知识问答模式也不允许手动指定文档 if (normalizedDocumentId != null) { throw new IllegalArgumentException("自动知识问答模式下不能传 selectedDocumentId"); } return null; } // 当前文档问答模式(DOCUMENT):必须传,而且必须是当前可检索的文档 if (normalizedDocumentId == null) { throw new IllegalArgumentException("当前文档问答模式下必须选择一个文档"); } final Long resolvedDocumentId = parseRequiredLong(normalizedDocumentId, "selectedDocumentId"); // 只允许命中"当前可检索"的文档,避免引用已下线或不可用的数据源 return documentKnowledgeService.listRetrievableDocuments().stream() .filter(item -> Objects.equals(item.getDocumentId(), resolvedDocumentId)) .findFirst() .orElseThrow(() -> new IllegalArgumentException("所选文档当前不可检索: " + normalizedDocumentId)); } ``` 这里的校验逻辑可以总结成一张表: | 聊天模式 | selectedDocumentId 规则 | | --- | --- | | `OPEN_CHAT` | 不允许传,传了就报错 | | `AUTO_DOCUMENT` | 不允许传,文档由系统自动路由 | | `DOCUMENT` | 必须传,且文档必须当前可检索 | 这种"在入口处就把非法参数拦住"的做法,让后续的编排器和执行器可以放心地使用这些参数,不用再做重复校验。 ## 抢占分布式租约 启动计划构建好之后,紧接着就是抢占 Redis 分布式租约。这是为了保证**同一个会话在任意时刻只有一个生成任务在运行**: ```java // BusinessChatService.java —— 抢占租约 private boolean claimConversationLease(StreamLaunchPlan launchPlan) { // 用 Redis 实现分布式锁,TTL 30 秒,后续会定期续期 return redisLeaseManager.acquire( launchPlan.getLeaseKey(), // 键:chat:running:{conversationId} launchPlan.getLeaseOwnerToken(), // 值:本次请求的唯一 token CHAT_RUNNING_LEASE_TTL // 过期时间:30 秒 ); } ``` 如果抢占失败,说明这个会话已经有一个任务在跑了,直接返回拒绝流: ```text // 租约抢占失败,返回错误提示 if (!leaseClaimed) { return rejectionFlux("该会话当前正在执行中,请稍后再试", launchPlan.getConversationId(), null); } ``` > 为什么需要分布式租约? > > 在集群部署场景下,用户可能快速连续点击发送,或者前端重试请求。如果没有租约机制,同一个会话可能在多个节点上同时生成回答,导致数据混乱。Redis 租约保证了全局唯一性。 ## Bootstrap:创建轮次、构建 TaskInfo、注册运行态 拿到租约之后,进入 `bootstrapConversation()`,这一步要做三件事: - 在数据库里创建一条新的轮次(exchange)记录 - 构建 `TaskInfo` 运行时上下文对象 - 把任务注册到内存运行态注册表 ```java // BusinessChatService.java —— 会话 bootstrap /** * 对会话做启动前置处理。 *

* 包括创建一条新的 exchange 归档记录、构建运行时任务对象、注册到运行时注册表,并把 SSE 通道与任务绑定。 *

* * @param launchPlan 已规范化后的启动计划 * @return bootstrap 结果;可能是可执行的流,也可能是一个拒绝原因 */ private BootstrapResult bootstrapConversation(StreamLaunchPlan launchPlan) { // exchangeView 表示本次问答轮次的归档记录,后续无论成功还是失败都依赖它进行收尾落库。 ConversationExchangeView exchangeView = null; try { // 一旦启动流程开始,就先在归档层生成一条“新轮次”,这样后续异常也能被定位到具体 exchange。 exchangeView = conversationArchiveStore.startExchange( launchPlan.getConversationId(), launchPlan.getQuestion(), launchPlan.getChatMode(), launchPlan.getSelectedDocumentId(), launchPlan.getSelectedDocumentName() ); // TaskInfo 聚合了本次会话运行所需的所有状态:SSE sink、trace、引用、上下文等。 TaskInfo taskInfo = createTaskInfo(launchPlan, exchangeView); if (!chatRuntimeRegistry.register(taskInfo)) { // 极端情况下,即使抢到租约,也可能在运行态注册时发现已有同会话任务占用,必须补偿性收尾。 failBootstrappedExchange(launchPlan.getConversationId(), exchangeView.getExchangeId(), "该会话当前正在执行中,请稍后再试"); releaseLeaseQuietly(launchPlan.getLeaseKey(), launchPlan.getLeaseOwnerToken()); return BootstrapResult.rejected("该会话当前正在执行中,请稍后再试"); } // 只有在归档、运行态、SSE 通道都准备好之后,才把流返回给上层。 return BootstrapResult.ready(bindClientChannel(taskInfo)); } catch (RuntimeException exception) { // bootstrap 过程中只要失败,就先释放租约,再把已经创建的轮次标记为失败,避免悬空数据。 releaseLeaseQuietly(launchPlan.getLeaseKey(), launchPlan.getLeaseOwnerToken()); if (exchangeView != null) { failBootstrappedExchange(launchPlan.getConversationId(), exchangeView.getExchangeId(), buildErrorMessage(exception)); } return BootstrapResult.rejected(buildErrorMessage(exception)); } } ``` ### TaskInfo:运行时的"万能上下文" `TaskInfo` 是整条执行链路的核心数据载体,几乎所有组件都要从它身上拿东西。看看它都装了什么: ```java // TaskInfo.java —— 运行时任务上下文 public class TaskInfo { private final String conversationId; // 会话 ID private final long exchangeId; // 轮次 ID private final String question; // 用户问题 private final ChatQueryMode chatMode; // 聊天模式 private volatile ConversationExecutionPlan executionPlan; // 执行计划(后续填充) private final RunnableConfig runnableConfig; // Agent 运行配置 private final ConversationTraceRecorder traceRecorder; // 执行追踪记录器 private final Sinks.Many sink; // SSE 事件推送通道 private final StringBuffer answerBuffer; // 答案累积缓冲区 private final List thinkingSteps; // 思考步骤 private final List references; // 引用来源 private final Set usedTools; // 使用过的工具 private final String leaseKey; // 租约键 private final String leaseOwnerToken; // 租约所有权 token private final long startTime; // 任务开始时间戳 private final AtomicLong firstResponseTimeMs; // 首包耗时 private final AtomicBoolean finalized; // 是否已结束(保证收尾只执行一次) } ``` 这里有几个设计值得注意: - `sink` 是 Reactor 的单播通道,所有 SSE 事件都通过它推给前端 - `answerBuffer` 用 `StringBuffer`(线程安全),因为模型输出和收尾落库可能在不同线程 - `finalized` 用 `AtomicBoolean`,通过 CAS 保证停止/成功/失败的收尾逻辑只执行一次 - `thinkingSteps`、`references`、`usedTools` 都用线程安全容器,因为执行过程中多个组件会并发写入 ### createTaskInfo:TaskInfo 是怎么组装出来的 `TaskInfo` 不是简单 new 出来的,它的构建过程涉及 SSE 通道初始化、Agent 运行配置、上下文注入等一系列操作: ```java // BusinessChatService.java —— 构建运行时任务对象 /** * 构建一个运行中的任务快照对象。 *

* 这里会初始化 SSE sink、RunnableConfig、调试追踪对象、引用集合、工具集合等上下文, * 后续执行链路中的各个组件都会围绕这个 {@link TaskInfo} 协作。 *

*/ private TaskInfo createTaskInfo(StreamLaunchPlan launchPlan, ConversationExchangeView exchangeView) { // 每个会话只有一个单播 sink,确保一条流只服务当前订阅的前端连接。 Sinks.Many sink = Sinks.many().unicast().onBackpressureBuffer(); // RunnableConfig 是底层 ReactAgent 和 checkpoint 体系识别当前线程上下文的关键对象。 RunnableConfig runnableConfig = buildSessionConfig(launchPlan.getConversationId()); // 这些集合会在执行过程中不断追加内容,因此使用线程安全容器保存运行态快照。 List thinkingSteps = Collections.synchronizedList(new ArrayList<>()); List references = Collections.synchronizedList(new ArrayList<>()); Set usedTools = ConcurrentHashMap.newKeySet(); String traceId = UUID.randomUUID().toString().replace("-", ""); // traceRecorder 负责把执行过程切成多个阶段,并记录每一阶段的状态、耗时和附加信息。 ConversationTraceRecorder traceRecorder = new ConversationTraceRecorder( conversationTraceStageStore, retrievalObserveStore, launchPlan.getConversationId(), exchangeView.getExchangeId(), traceId ); // eventMetadata 会被放进每条 SSE 事件中,方便前端知道事件归属哪个会话、哪个轮次。 StreamEventMetadata eventMetadata = new StreamEventMetadata( launchPlan.getConversationId(), exchangeView.getExchangeId() ); // 下方这些 context 值会被执行链路中的工具、检索器、追踪器共同读取。 runnableConfig.context().put(ChatContextKeys.EVENT_SINK, sink); runnableConfig.context().put(ChatContextKeys.EVENT_METADATA, eventMetadata); runnableConfig.context().put(ChatContextKeys.THINKING_STEPS, thinkingSteps); runnableConfig.context().put(ChatContextKeys.REFERENCES, references); runnableConfig.context().put(ChatContextKeys.USED_TOOLS, usedTools); runnableConfig.context().put(ChatContextKeys.TRACE_ID, traceId); runnableConfig.context().put(ChatContextKeys.QUESTION, launchPlan.getQuestion()); runnableConfig.context().put(ChatContextKeys.CHAT_MODE, launchPlan.getChatMode().name()); // 把“当前日期”和它的文本形式显式注入上下文,是为了让相对时间问题有统一的锚点。 runnableConfig.context().put(ChatContextKeys.CURRENT_DATE, launchPlan.getCurrentDate().toString()); runnableConfig.context().put(ChatContextKeys.CURRENT_DATE_TEXT, launchPlan.getCurrentDateText()); // 文档问答模式下,所选文档和索引任务信息也会随上下文透传。 putContextIfNotNull(runnableConfig, ChatContextKeys.SELECTED_DOCUMENT_ID, launchPlan.getSelectedDocumentId()); putContextIfNotBlank(runnableConfig, ChatContextKeys.SELECTED_DOCUMENT_NAME, launchPlan.getSelectedDocumentName()); putContextIfNotNull(runnableConfig, ChatContextKeys.SELECTED_TASK_ID, launchPlan.getSelectedTaskId()); // 在真正生成 executionPlan 之前,先放一个空白调试轨迹,保证链路中随时都能读取到 debugTrace。 ChatDebugTrace debugTrace = initializeDebugTrace(null); runnableConfig.context().put(ChatContextKeys.DEBUG_TRACE, debugTrace); return new TaskInfo( launchPlan.getConversationId(), exchangeView.getExchangeId(), launchPlan.getQuestion(), launchPlan.getChatMode(), traceId, launchPlan.getSelectedDocumentId(), launchPlan.getSelectedDocumentName(), launchPlan.getSelectedTaskId(), launchPlan.getCurrentDate(), launchPlan.getCurrentDateText(), null, debugTrace, runnableConfig, traceRecorder, sink, eventMetadata, launchPlan.getLeaseKey(), launchPlan.getLeaseOwnerToken(), thinkingSteps, references, usedTools, System.currentTimeMillis() ); } ``` 这里有个很重要的设计:`RunnableConfig.context()` 是一个共享的 Map,执行链路中的各个组件(工具、检索器、追踪器)都通过它来读写运行时状态。这样就不需要在每个方法签名里传一堆参数,所有组件都能通过 context 拿到自己需要的东西。 ## 绑定 SSE 通道与激活生成 Bootstrap 完成后,`bindClientChannel()` 把 TaskInfo 的内部 sink 转成前端可消费的 Flux: ```java // BusinessChatService.java —— 绑定客户端通道 /** * 把任务的内部 sink 绑定成可返回给前端的 Flux。 *

* 只有前端真正订阅时,才会触发生成任务启动;如果前端断开订阅,则主动停止当前任务。 *

*/ private Flux bindClientChannel(TaskInfo taskInfo) { return taskInfo.sink().asFlux() // 前端真正订阅时,才异步启动生成任务 .doOnSubscribe(ignored -> activateGeneration(taskInfo)) // 前端断开连接时,主动中断后台生成,避免资源空转 .doOnCancel(() -> stopTask(taskInfo, "客户端已取消请求")); } ``` 这里的 `doOnSubscribe` 是整条链路的"点火开关"——前端建立 SSE 连接的那一刻,后端才真正开始干活。 `activateGeneration()` 做两件事:启动租约续期、启动执行链路: ```java // BusinessChatService.java —— 激活生成 private void activateGeneration(TaskInfo taskInfo) { try { if (taskInfo.finalized().get()) { return; // 任务已被其他线程结束,不再重复启动 } // 长时间执行的会话需要定期续租,否则别的节点会认为任务已失效 Disposable leaseRenewalDisposable = startLeaseRenewal(taskInfo); taskInfo.setLeaseRenewalDisposable(leaseRenewalDisposable); // 真正的生成执行链路在这里启动订阅 Disposable disposable = buildConversationExecution(taskInfo).subscribe(); taskInfo.setDisposable(disposable); // 如果任务在刚启动后立即被标记为结束,主动 dispose if (taskInfo.finalized().get() && !disposable.isDisposed()) { disposable.dispose(); } } catch (RuntimeException exception) { finishWithFailure(taskInfo, exception); } } ``` > 租约续期机制 > > 租约 TTL 是 30 秒,每 10 秒续期一次。如果续期失败(比如 Redis 连接断了),会自动停止当前会话,防止无租约状态下继续执行。这个设计保证了即使节点宕机,租约也会在 30 秒后自动释放,不会永久阻塞后续请求。 展开看看租约续期的具体实现: ### startLeaseRenewal:定时续期 ```java // BusinessChatService.java —— 启动租约续期任务 private Disposable startLeaseRenewal(TaskInfo taskInfo) { // 每隔 10 秒触发一次续期 return Flux.interval(CHAT_RUNNING_LEASE_RENEW_INTERVAL, CHAT_RUNNING_LEASE_RENEW_INTERVAL) .subscribe( ignored -> renewLeaseOrStop(taskInfo), error -> log.warn("租约续期任务出现异常, conversationId={}", taskInfo.conversationId(), error) ); } ``` ### renewLeaseOrStop:续期失败自动停止 ```java // BusinessChatService.java —— 续期或停止 private void renewLeaseOrStop(TaskInfo taskInfo) { // 尝试续期 boolean renewed = redisLeaseManager.renew( taskInfo.leaseKey(), taskInfo.leaseOwnerToken(), CHAT_RUNNING_LEASE_TTL ); if (renewed) { return; // 续期成功,继续执行 } // 续期失败,说明租约已经被别人抢走或者 Redis 出了问题 log.warn("会话租约续期失败,准备停止当前会话, conversationId={}", taskInfo.conversationId()); // 先停掉续期定时器自身,避免重复触发 Disposable leaseRenewalDisposable = taskInfo.leaseRenewalDisposable(); if (leaseRenewalDisposable != null && !leaseRenewalDisposable.isDisposed()) { leaseRenewalDisposable.dispose(); } // 主动停止当前会话 stopTask(taskInfo, "会话租约已失效,已停止生成"); } ``` 这个"续期失败就自动停止"的机制很关键——它保证了系统不会出现"租约已经过期但任务还在跑"的幽灵状态。 ## 核心:组装执行流 buildConversationExecution 这是整条链路最关键的方法,把"准备执行计划 → 选择执行器 → 消费模型输出 → 收尾"串成一条完整的响应式流: ```java // BusinessChatService.java —— 组装完整的对话执行流 /** * 组装完整的对话执行流。 *

* 链路大致分为:发送“正在分析”提示 -> 准备执行计划 -> 按计划选择执行器 -> 消费模型输出 -> * 正常完成时收尾,异常时失败收尾。 *

*/ private Flux buildConversationExecution(TaskInfo taskInfo) { return Flux.defer(() -> { // 先给前端发一个"分析中"的状态事件,让用户知道系统在工作 safeEmit(taskInfo.sink(), streamEventWriter.thinking("正在分析问题上下文。", taskInfo.eventMetadata())); return Mono.fromCallable(() -> prepareExecutionPlan(taskInfo)) // 计划准备包含检索、压缩历史、读取摘要等 IO 操作,放到弹性线程池 .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)); } ``` 这段代码的执行顺序可以用下面这张图来理解: ![[FqWVKz1q2uwuYAn89geDA_JgIhwL-e3af166c.png]] ## 编排器:prepareExecutionPlan 执行计划的生成由 `ChatPreparationOrchestrator` 负责,它是整条链路的"大脑",决定了这次问答到底走哪条路。 ```java // BusinessChatService.java —— 准备执行计划 /** * 准备本轮会话真正执行所需的编排计划。 *

* 这一层会调用编排器分析历史上下文、决定执行模式、构造 agentQuestion,并在必要时刷新会话绑定的文档范围。 *

*/ private ConversationExecutionPlan prepareExecutionPlan(TaskInfo taskInfo) { // 编排器会综合问题、历史、摘要、文档选择等信息生成本轮执行计划。 ConversationExecutionPlan executionPlan = chatPreparationOrchestrator.prepare(taskInfo); // agentQuestion 是最终喂给 Agent 的问题文本,会补充时间锚点和上下文摘要。 executionPlan.setAgentQuestion(buildAgentQuestion(executionPlan)); if (executionPlan.getSelectedDocumentId() != null && !Objects.equals(executionPlan.getSelectedDocumentId(), taskInfo.selectedDocumentId())) { // 如果编排阶段修正了文档范围,需要同步刷新归档中的会话范围与运行上下文。 conversationArchiveStore.refreshSessionScope( taskInfo.conversationId(), executionPlan.getChatMode(), executionPlan.getSelectedDocumentId(), executionPlan.getSelectedDocumentName() ); putContextIfNotNull(taskInfo.runnableConfig(), ChatContextKeys.SELECTED_DOCUMENT_ID, executionPlan.getSelectedDocumentId()); putContextIfNotBlank(taskInfo.runnableConfig(), ChatContextKeys.SELECTED_DOCUMENT_NAME, executionPlan.getSelectedDocumentName()); putContextIfNotNull(taskInfo.runnableConfig(), ChatContextKeys.SELECTED_TASK_ID, executionPlan.getSelectedTaskId()); } // 把最终执行计划和对应的调试轨迹回写到任务对象,供执行链路和收尾阶段复用。 taskInfo.setExecutionPlan(executionPlan); taskInfo.setDebugTrace(initializeDebugTrace(executionPlan)); taskInfo.runnableConfig().context().put(ChatContextKeys.DEBUG_TRACE, taskInfo.debugTrace()); return executionPlan; } ``` 编排器内部的 `prepare()` 方法做了很多事,我们拆开来看核心逻辑: ```java // ChatPreparationOrchestrator.java —— 编排器核心逻辑 public ConversationExecutionPlan prepare(TaskInfo taskInfo) { String conversationId = taskInfo.conversationId(); String question = taskInfo.question(); ChatQueryMode chatMode = taskInfo.chatMode(); Long selectedDocumentId = taskInfo.selectedDocumentId(); String selectedDocumentName = taskInfo.selectedDocumentName(); Long selectedTaskId = taskInfo.selectedTaskId(); LocalDate currentDate = taskInfo.currentDate(); String currentDateText = taskInfo.currentDateText(); ConversationTraceRecorder traceRecorder = taskInfo.traceRecorder(); ConversationTraceRecorder.StageHandle memoryStage = traceRecorder == null ? null : traceRecorder.startStage(ConversationTraceStageCode.MEMORY, chatMode == null ? "" : chatMode.name(), "正在装载会话记忆与最近窗口。", null); ConversationMemoryContext memoryContext; try { // 第一步:加载会话记忆(长期摘要 + 最近对话窗口) memoryContext = summarizeHistory(conversationId, traceRecorder); if (traceRecorder != null) { traceRecorder.completeStage(memoryStage, "会话记忆装载完成。", java.util.Map.of( "compressionApplied", memoryContext != null && memoryContext.isCompressionApplied(), "coveredExchangeId", memoryContext == null ? 0L : memoryContext.getCoveredExchangeId(), "coveredExchangeCount", memoryContext == null ? 0 : memoryContext.getCoveredExchangeCount(), "compressionCount", memoryContext == null ? 0 : memoryContext.getCompressionCount(), "longTermSummary", memoryContext == null ? "" : safeText(memoryContext.getLongTermSummary()), "recentTranscript", memoryContext == null ? "" : safeText(memoryContext.getRecentTranscript()), "answerRecentTranscript", memoryContext == null ? "" : safeText(memoryContext.getAnswerRecentTranscript()) )); } } catch (RuntimeException exception) { if (traceRecorder != null) { traceRecorder.failStage(memoryStage, "会话记忆装载失败。", exception.getMessage(), null); } throw exception; } HistoryPlanningContext historyPlanningContext = buildHistoryPlanningContext(memoryContext); // 第二步:构建历史上下文,供后续问题改写和回答使用 String historySummary = buildPlanningHistory(memoryContext, historyPlanningContext); AnswerHistoryContext answerHistoryContext = buildAnswerHistoryContext( question, memoryContext == null ? "" : memoryContext.getAnswerRecentTranscript() ); // 第三步:判断时效性——用户问的是不是"今天""最新"这类需要实时信息的问题 boolean requiresCurrentDateAnchoring = TimeSensitiveQueryHelper.requiresCurrentDateAnchoring(question); boolean requiresFreshSearch = TimeSensitiveQueryHelper.requiresFreshSearch(question); if (chatMode == null) { throw new IllegalArgumentException("chatMode 不能为空"); } // 第四步:根据聊天模式走不同分支 if (chatMode == ChatQueryMode.OPEN_CHAT) { ConversationExecutionPlan plan = basePlan(question, chatMode, memoryContext, historyPlanningContext, historySummary, answerHistoryContext, currentDate, currentDateText, requiresCurrentDateAnchoring, requiresFreshSearch) .mode(ExecutionMode.REACT_AGENT) .build(); if (traceRecorder != null) { ConversationTraceRecorder.StageHandle routeStage = traceRecorder.startStage(ConversationTraceStageCode.ROUTE, ExecutionMode.REACT_AGENT.name(), "路由到开放式 Agent。", null); traceRecorder.completeStage(routeStage, "已判定走开放式 Agent 路径。", java.util.Map.of( "chatMode", chatMode.name(), "executionMode", ExecutionMode.REACT_AGENT.name(), "requiresFreshSearch", requiresFreshSearch, "requiresCurrentDateAnchoring", requiresCurrentDateAnchoring )); } return plan; } // 文档问答模式 → 需要问题改写 + 知识路由 + 执行模式判定 // ...(后续文档会详细展开) } ``` 编排器的路由决策可以用这张图来概括: ![[FjrgQP3AMZoVuMwQffj9JXO1dDhE-789354b4.png]] ### buildAgentQuestion:构造最终喂给 Agent 的问题 编排器生成执行计划之后,还需要把用户的原始问题"包装"一下,加上时间锚点和历史摘要,让 Agent 有足够的上下文来回答: ```java // BusinessChatService.java —— 构造 Agent 问题 /** * 构造最终发给 Agent 的问题文本。 *

* 这里会把时间锚点、时效性约束、历史摘要和原始问题拼接成统一 prompt, * 让 Agent 在处理“今天/最新/本周”等表达时有明确基准。 *

*/ private String buildAgentQuestion(ConversationExecutionPlan executionPlan) { StringBuilder builder = new StringBuilder(); // 注入系统时间信息,作为处理相对时间的统一基准 builder.append("系统时间信息:\n"); builder.append("当前日期是 ").append(executionPlan.getCurrentDateText()) .append(",时区为 Asia/Shanghai。\n"); if (executionPlan.isRequiresCurrentDateAnchoring()) { // 对强时效问题补充更严格的日期约束 builder.append("当前问题包含相对时间或强时效语义。"); builder.append("当用户提到"今天、明天、昨天、现在、当前、最新、本周、本月、今年"等表达时,"); builder.append("必须以这个日期为准,不要把搜索结果里的旧日期误当成今天。\n"); } else { builder.append("当用户提到"今天、明天、昨天、现在、当前、最新"等相对时间时,必须以这个日期为准。\n"); } if (executionPlan.isRequiresFreshSearch()) { // 强制联网核实最新事实 builder.append("当前问题需要核实最新外部事实,回答前必须优先调用联网搜索工具。\n"); builder.append("如果搜索结果里的日期与当前日期不一致,必须明确说明来源日期。\n"); builder.append("如果无法找到与当前日期匹配的可靠结果,要明确说明不确定性,不要编造最新信息。\n"); } // 如果有历史摘要,也一并注入,让 Agent 知道之前聊了什么 if (StrUtil.isNotBlank(executionPlan.getHistorySummary())) { builder.append("\n相关会话背景:\n"); builder.append(executionPlan.getHistorySummary()).append("\n"); } // 最后才是用户的原始问题 builder.append("\n用户问题:\n"); builder.append(executionPlan.getOriginalQuestion()); return builder.toString(); } ``` 这个方法的设计思路是:**不要让 Agent 裸接用户问题**。用户问"今天天气怎么样",如果不告诉 Agent 今天是几号,它可能会用训练数据里的旧日期来回答。通过在问题前面注入时间锚点和历史摘要,Agent 就有了足够的上下文来给出准确的回答。 ## 执行器注册表:策略模式路由 编排器确定了执行模式之后,`ConversationExecutorRegistry` 负责找到对应的执行器。这里用的是经典的**策略模式 + 注册表**: ```java // ConversationExecutorRegistry.java —— 执行器注册表 @Component public class ConversationExecutorRegistry { // 用 EnumMap 存储执行模式到执行器的映射,查找效率 O(1) private final Map executorMap = new EnumMap<>(ExecutionMode.class); // Spring 会自动注入所有实现了 ConversationExecutor 接口的 Bean 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; } } ``` 所有执行器都实现同一个接口: ```typescript // ConversationExecutor.java —— 统一执行器接口 public interface ConversationExecutor { ExecutionMode mode(); // 声明自己处理哪种模式 Flux execute(TaskInfo taskInfo); // 执行并返回文本流 } ``` 目前系统支持的执行模式有: | 执行模式 | 对应执行器 | 适用场景 | | --- | --- | --- | | `REACT_AGENT` | ReactAgentExecutor | 开放式提问,支持工具调用和联网搜索 | | `RETRIEVAL` | RagChatExecutor | 文档知识问答,走 RAG 检索链路 | | `GRAPH_ONLY` | GraphOnlyExecutor | 纯图查询,适合结构化导航类问题 | | `GRAPH_THEN_EVIDENCE` | GraphThenEvidenceExecutor | 先图查询再补充证据 | | `CLARIFICATION` | ClarificationExecutor | 文档范围歧义时,向用户确认 | 这个设计的好处是:新增一种聊天模式,只需要实现 `ConversationExecutor` 接口并注册为 Spring Bean,注册表会自动发现它,完全不需要改动已有代码。 ## 流式输出:emitModelChunk 执行器工作过程中,模型每产出一小段文字,就会通过 Flux 推出来,然后被 `emitModelChunk()` 处理: ```java // BusinessChatService.java —— 处理模型输出的单个增量片段 /** * 处理模型输出的单个增量片段。 *

* 每收到一个 chunk,就同时做三件事:追加答案缓冲区、记录首包耗时、向前端发送文本事件。 *

*/ private void emitModelChunk(TaskInfo taskInfo, String chunk) { // answerBuffer 持续累积,最终落库时用它拿到完整答案 taskInfo.answerBuffer().append(chunk); // 首包耗时只记录第一次收到正文输出的时刻 if (taskInfo.firstResponseTimeMs().get() == 0L) { taskInfo.firstResponseTimeMs() .compareAndSet(0L, System.currentTimeMillis() - taskInfo.startTime()); } // 每个 chunk 都即时推给前端,形成"边生成边展示"的效果 safeEmit(taskInfo.sink(), streamEventWriter.text(chunk, taskInfo.eventMetadata())); } ``` `StreamEventWriter` 负责把内容包装成标准的 JSON 事件格式: ```java // StreamEventWriter.java —— SSE 事件格式化 @Component public class StreamEventWriter { private final ObjectMapper objectMapper; // 文本事件:模型输出的正文增量 public String text(String content, StreamEventMetadata metadata) { return write(event("text", content, metadata)); } // 思考事件:分析中的状态提示 public String thinking(String content, StreamEventMetadata metadata) { return write(event("thinking", content, metadata)); } // 错误事件:执行失败时的错误信息 public String error(String content, StreamEventMetadata metadata) { return write(event("error", content, metadata)); } // 引用事件:检索命中的来源文档 public String references(List references, StreamEventMetadata metadata) { Map payload = event("reference", references, metadata); payload.put("count", references != null ? references.size() : 0); return write(payload); } // 推荐事件:生成的推荐追问 public String recommendations(List recommendations, StreamEventMetadata metadata) { Map payload = event("recommend", recommendations, metadata); payload.put("count", recommendations != null ? recommendations.size() : 0); return write(payload); } // 统一事件结构:type + content + timestamp + 会话元信息 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; } } ``` 前端收到的每条 SSE 事件大概长这样: ```json { "type": "text", "content": "根据文档内容,", "timestamp": "2025-05-20T10:30:00.123Z", "conversationId": "abc123", "exchangeId": 42 } ``` ## 成功收尾:finishSuccessfully 模型输出完毕后,`doOnComplete` 触发成功收尾逻辑。这一步要做的事情不少: ```java // BusinessChatService.java —— 成功收尾 /** * 处理会话成功完成时的统一收尾。 *

* 这个方法位于“模型正文已经正常输出完成”之后,是一次成功对话真正结束前的最后一道总收口。 * 它承担的不是单一动作,而是一整套按顺序执行的完成态闭环: * 1. 通过 CAS 把任务标记为 finalized,确保成功收尾只执行一次; * 2. 从运行态缓冲区中冻结最终答案、引用和推荐追问所需的数据快照; * 3. 在追踪体系中开启 finalize/recommendation 阶段,便于调试和耗时分析; * 4. 生成或提取推荐追问,并把推荐阶段标记为完成; * 5. 向前端补发“引用”和“推荐追问”事件; * 6. 关闭 SSE 流,告诉前端这一轮输出已经彻底结束; * 7. 以 {@link ChatTurnStatus#COMPLETED} 状态把最终结果完整落库; * 8. 异步刷新会话摘要,并清理租约、订阅、运行态注册表等临时资源。 *

*

* 这里的顺序不能随意打乱。特别是: * “补发事件”必须发生在“关闭 SSE 流”之前,否则前端会收不到引用和推荐; * “落库和清理”要放在 finally 中,保证即使补发事件失败,也不会让会话停留在未收尾状态。 *

*/ private void finishSuccessfully(TaskInfo taskInfo) { // finalized 从 false 置为 true 说明当前线程拿到了“唯一一次成功收尾权”; // 如果这里失败,表示别的线程已经做过停止/失败/成功收尾,本次直接退出避免重复落库。 if (!taskInfo.finalized().compareAndSet(false, true)) { return; } // answer 是最终要持久化的完整回答文本; // uniqueReferences 先对运行态引用做快照再去重,避免后续落库和前端展示出现重复证据。 String answer = taskInfo.answerBuffer().toString(); List uniqueReferences = deduplicateReferences(snapshotReferenceList(taskInfo.references())); // finalizeStage 用来观测“成功收尾”本身的耗时和状态,而不是正文生成耗时。 ConversationTraceRecorder.StageHandle finalizeStage = taskInfo.traceRecorder() == null ? null : taskInfo.traceRecorder().startStage( org.javaup.ai.chatagent.model.trace.ConversationTraceStageCode.FINALIZE, taskInfo.executionPlan() == null || taskInfo.executionPlan().getMode() == null ? "" : taskInfo.executionPlan().getMode().name(), "正在收尾已完成会话。", null ); // recommendationStage 单独拆出来,是因为推荐追问生成可能本身就是一个有成本、可失败、可观测的子阶段。 ConversationTraceRecorder.StageHandle recommendationStage = taskInfo.traceRecorder() == null ? null : taskInfo.traceRecorder().startStage( org.javaup.ai.chatagent.model.trace.ConversationTraceStageCode.RECOMMENDATION, taskInfo.executionPlan() == null || taskInfo.executionPlan().getMode() == null ? "" : taskInfo.executionPlan().getMode().name(), "正在生成推荐追问。", null ); List recommendations; if (taskInfo.executionPlan() != null && taskInfo.executionPlan().getMode() == org.javaup.ai.chatagent.rag.model.ExecutionMode.CLARIFICATION) { // 澄清模式下,“推荐追问”本质上就是编排阶段已经产出的澄清选项; // 这里直接复用,避免再调用推荐服务生成一批与澄清问题语义不一致的新建议。 recommendations = taskInfo.executionPlan().getClarificationOptions() == null ? List.of() : new ArrayList<>(taskInfo.executionPlan().getClarificationOptions()); } else { // 常规问答模式下,再基于“原问题 + 最终答案 + 最近历史轮次”生成后续建议, // 这样推荐内容能更贴近本轮回答结果,而不是只基于用户问题孤立生成。 recommendations = recommendationService.generateRecommendations( taskInfo.question(), answer, historicalRecentExchanges(taskInfo), taskInfo.traceRecorder() ); } if (taskInfo.traceRecorder() != null) { // 推荐阶段到这里就算结束,无论这些推荐稍后能否成功发到前端, // 至少“推荐内容本身已经生成出来”这一事实要被追踪记录下来。 taskInfo.traceRecorder().completeStage(recommendationStage, "推荐追问生成完成。", Map.of( "recommendationCount", recommendations.size(), "recommendations", recommendations )); } try { // 正文流结束后再补发引用,前端可以按“正文 -> 引用”的顺序稳定渲染, // 避免引用事件插入正文中间导致页面状态机更复杂。 if (!uniqueReferences.isEmpty()) { safeEmit(taskInfo.sink(), streamEventWriter.references(uniqueReferences, taskInfo.eventMetadata())); } // 推荐追问放在引用之后发送,表示“本轮回答及其证据已经给齐,接下来给你下一步建议”。 if (!recommendations.isEmpty()) { safeEmit(taskInfo.sink(), streamEventWriter.recommendations(recommendations, taskInfo.eventMetadata())); } } catch (RuntimeException exception) { // 这里故意只记 warn,不把会话整体改判为失败; // 因为正文已经成功完成,补发事件失败不应推翻“本轮回答成功生成”这一主结果。 log.warn("补发引用或推荐事件失败, conversationId={}, exchangeId={}", taskInfo.conversationId(), taskInfo.exchangeId(), exception); } finally { try { // 无论补发事件成功与否,都要主动关闭 SSE 流,明确告诉前端“这一轮已经结束”。 safeComplete(taskInfo.sink()); } catch (RuntimeException exception) { log.warn("关闭成功完成的 SSE 流失败, conversationId={}, exchangeId={}", taskInfo.conversationId(), taskInfo.exchangeId(), exception); } try { // 先把运行期间累积的模型调用、工具调用等统计补写进 debugTrace, // 再统一把本轮结果以 COMPLETED 状态落库,保证数据库里存的是最终完整快照。 refreshDebugTraceRuntimeStats(taskInfo); conversationArchiveStore.completeExchange( taskInfo.conversationId(), taskInfo.exchangeId(), answer, // thinkingSteps / references / tools 都在收尾时取快照,避免并发修改影响持久化结果。 snapshotStringList(taskInfo.thinkingSteps()), uniqueReferences, recommendations, snapshotUsedTools(taskInfo.usedTools()), taskInfo.debugTrace(), ChatTurnStatus.COMPLETED, "", toNullable(taskInfo.firstResponseTimeMs().get()), System.currentTimeMillis() - taskInfo.startTime() ); if (taskInfo.traceRecorder() != null) { // finalize 阶段在“落库成功”后才标记完成,这样 trace 中的完成态才真正代表整轮闭环完成。 taskInfo.traceRecorder().completeStage(finalizeStage, "会话已按完成状态收尾。", Map.of( "finalStatus", ChatTurnStatus.COMPLETED.name(), "referenceCount", uniqueReferences.size(), "recommendationCount", recommendations.size(), "answerLength", answer.length() )); } } catch (RuntimeException exception) { // 到这里说明“回答是成功生成的”,但“成功态收尾”失败了,通常是落库或调试追踪更新失败; // 因此日志级别提升为 error,方便排查为什么用户看到了回答但后台没有完整持久化。 log.error("成功会话收尾落库失败, conversationId={}, exchangeId={}", taskInfo.conversationId(), taskInfo.exchangeId(), exception); if (taskInfo.traceRecorder() != null) { taskInfo.traceRecorder().failStage(finalizeStage, "完成态收尾失败。", exception.getMessage(), null); } } finally { // 这里无论前面的“补发事件”“关闭流”“落库”是否成功,都必须做最终清理: // 1. 尝试刷新长期摘要,让后续轮次能拿到最新上下文; // 2. 释放租约、订阅和运行态注册,避免同一个 conversationId 被永久占用。 safeRefreshConversationSummary(taskInfo.conversationId()); cleanup(taskInfo); } } } ``` 收尾阶段的事件推送顺序是:**正文 → 引用 → 推荐追问 → 关闭流**。前端按这个顺序依次渲染,用户先看到答案,再看到引用来源和推荐的下一步问题。 ## 资源清理:cleanup 最后一步是释放所有运行态资源: ```java // BusinessChatService.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(); } // 释放 Redis 租约,其他节点或后续请求才能再次启动同一个会话 releaseLeaseQuietly(taskInfo.leaseKey(), taskInfo.leaseOwnerToken()); // 从运行时注册表移除当前任务 chatRuntimeRegistry.remove(taskInfo.conversationId(), taskInfo); } ``` ## 总结 回顾一下整条链路,一个用户问题从发出到收到答案,后端一共经历了这些阶段: - **Controller 接收请求** → 参数校验,转交 Service - **Flux.defer 延迟启动** → 等前端真正订阅才开始 - **buildLaunchPlan** → 规范化参数,构建启动计划 - **claimConversationLease** → 抢占 Redis 分布式租约 - **bootstrapConversation** → 创建轮次记录、构建 TaskInfo、注册运行态 - **bindClientChannel** → 绑定 SSE 输出通道 - **activateGeneration** → 启动租约续期 + 执行链路 - **prepareExecutionPlan** → 编排器分析意图,生成执行计划 - **ConversationExecutorRegistry.get** → 根据模式选择执行器 - **executor.execute** → 执行器工作,模型开始生成 - **emitModelChunk** → 逐块推送文本到前端 - **finishSuccessfully** → 补发引用和推荐、落库归档、释放资源 每一步都有明确的职责边界,异常处理也贯穿始终——任何一步失败都会走统一的失败收尾逻辑,保证不会留下悬空的租约或半开的会话状态。 --- **企业级项目导航**:⬅️ [[09-落库向量化收尾|09-落库向量化收尾]] | 01-从用户提问到答案返回的总流程 | ➡️ [[02-前后端模块划分与调用关系|02-前后端模块划分与调用关系]]