从用户提问到答案返回的总流程
这篇文档会把 Super Agent 聊天系统后端的完整链路拆开来讲,从用户点击"发送"的那一刻开始,一直到答案流式输出到前端、最后落库收尾为止。每个关键步骤都会贴出对应的源码,加上注释说明它在整条链路里的作用。
看完这篇,你会对"一个问题是怎么从 Controller 一路走到模型输出再回到前端"有一个完整的认知。
总流程概览
先看一张全局流程图,对整条链路有个直观印象:
接下来我们按照这张图的顺序,逐步拆解每个阶段的源码。
入口:Controller 接收请求
一切从前端的 POST 请求开始。前端把用户的问题、会话 ID、聊天模式等信息打包成 ChatRequestDto,发到 /api/chat/stream 接口。
先看请求参数长什么样:
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:
@AllArgsConstructor
@RestController
@RequestMapping("/api/chat")
public class BusinessChatController {
private final BusinessChatService businessChatService;
/**
* 打开一个流式会话。
* <p>
* 该接口返回的是 SSE 文本流,前端可以持续接收“思考中”、“正文增量”、“引用”、“推荐追问”等事件。
* </p>
*
* @param dto 前端提交的聊天请求,包含问题、会话 ID、聊天模式、选中文档等信息
* @return SSE 字符串流,内容由服务层按事件格式持续输出
*/
@PostMapping(value = "/stream", produces = "text/event-stream;charset=UTF-8")
public Flux<String> stream(@Valid @RequestBody ChatRequestDto dto) {
// 这里不再额外包装 ApiResponse,而是直接把服务层生成的 SSE 事件流返回给前端逐段消费。
return businessChatService.openConversationStream(dto);
}
}
为什么返回 Flux 而不是普通 JSON?
因为聊天回答是流式生成的,模型每产出一小段文字就立刻推给前端,用户能看到"边想边写"的效果。这里用的是 Spring WebFlux 的
Flux<String>,配合text/event-stream内容类型,实现了 SSE(Server-Sent Events)协议。
延迟启动:Flux.defer 的设计意图
Controller 调用的 openConversationStream() 并不会立刻开始干活,而是用 Flux.defer 包了一层:
public Flux<String> openConversationStream(ChatRequestDto request) {
// defer 的作用:把真正的启动逻辑延后到"客户端真正订阅流"的那一刻
// 避免只是创建 Flux 对象时就提前占用租约、创建轮次
return Flux.defer(() -> openDeferredConversationStream(request));
}
这个设计很关键——如果不用 defer,Flux 对象一创建就会执行内部逻辑,但这时候前端可能还没准备好接收数据。用了 defer 之后,只有前端真正建立 SSE 连接(订阅 Flux)时,后端才会开始抢租约、创建轮次这些操作。
构建启动计划:buildLaunchPlan
真正的启动逻辑在 openDeferredConversationStream() 里,第一步就是把前端传来的参数规范化,转成内部使用的 StreamLaunchPlan 对象:
/**
* 把外部请求转换成内部启动计划。
* <p>
* 这一步会完成问题与会话 ID 规范化、聊天模式解析、所选文档校验,以及时间锚点准备。
* </p>
*/
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 规范化
// BusinessChatService.java —— 规范化 conversationId
private String normalizeConversationId(String conversationId) {
// 前端传了就直接用(去掉首尾空格)
if (StrUtil.isNotBlank(conversationId)) {
return conversationId.trim();
}
// 没传就自动生成一个 UUID 作为新会话的 ID
return UUID.randomUUID().toString().replace("-", "");
}
这个设计让前端可以灵活控制:传了 conversationId 就是继续已有会话,不传就是开启新会话。
parseRequiredChatMode:聊天模式解析
// 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 是否合法,不同模式有不同的规则:
// 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 分布式租约。这是为了保证同一个会话在任意时刻只有一个生成任务在运行:
// 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 秒
);
}
如果抢占失败,说明这个会话已经有一个任务在跑了,直接返回拒绝流:
// 租约抢占失败,返回错误提示
if (!leaseClaimed) {
return rejectionFlux("该会话当前正在执行中,请稍后再试",
launchPlan.getConversationId(), null);
}
为什么需要分布式租约?
在集群部署场景下,用户可能快速连续点击发送,或者前端重试请求。如果没有租约机制,同一个会话可能在多个节点上同时生成回答,导致数据混乱。Redis 租约保证了全局唯一性。
Bootstrap:创建轮次、构建 TaskInfo、注册运行态
拿到租约之后,进入 bootstrapConversation(),这一步要做三件事:
- 在数据库里创建一条新的轮次(exchange)记录
- 构建
TaskInfo运行时上下文对象 - 把任务注册到内存运行态注册表
// BusinessChatService.java —— 会话 bootstrap
/**
* 对会话做启动前置处理。
* <p>
* 包括创建一条新的 exchange 归档记录、构建运行时任务对象、注册到运行时注册表,并把 SSE 通道与任务绑定。
* </p>
*
* @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 是整条执行链路的核心数据载体,几乎所有组件都要从它身上拿东西。看看它都装了什么:
// 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<String> sink; // SSE 事件推送通道
private final StringBuffer answerBuffer; // 答案累积缓冲区
private final List<String> thinkingSteps; // 思考步骤
private final List<SearchReference> references; // 引用来源
private final Set<String> 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 运行配置、上下文注入等一系列操作:
// BusinessChatService.java —— 构建运行时任务对象
/**
* 构建一个运行中的任务快照对象。
* <p>
* 这里会初始化 SSE sink、RunnableConfig、调试追踪对象、引用集合、工具集合等上下文,
* 后续执行链路中的各个组件都会围绕这个 {@link TaskInfo} 协作。
* </p>
*/
private TaskInfo createTaskInfo(StreamLaunchPlan launchPlan, ConversationExchangeView exchangeView) {
// 每个会话只有一个单播 sink,确保一条流只服务当前订阅的前端连接。
Sinks.Many<String> sink = Sinks.many().unicast().onBackpressureBuffer();
// RunnableConfig 是底层 ReactAgent 和 checkpoint 体系识别当前线程上下文的关键对象。
RunnableConfig runnableConfig = buildSessionConfig(launchPlan.getConversationId());
// 这些集合会在执行过程中不断追加内容,因此使用线程安全容器保存运行态快照。
List<String> thinkingSteps = Collections.synchronizedList(new ArrayList<>());
List<SearchReference> references = Collections.synchronizedList(new ArrayList<>());
Set<String> 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:
// BusinessChatService.java —— 绑定客户端通道
/**
* 把任务的内部 sink 绑定成可返回给前端的 Flux。
* <p>
* 只有前端真正订阅时,才会触发生成任务启动;如果前端断开订阅,则主动停止当前任务。
* </p>
*/
private Flux<String> bindClientChannel(TaskInfo taskInfo) {
return taskInfo.sink().asFlux()
// 前端真正订阅时,才异步启动生成任务
.doOnSubscribe(ignored -> activateGeneration(taskInfo))
// 前端断开连接时,主动中断后台生成,避免资源空转
.doOnCancel(() -> stopTask(taskInfo, "客户端已取消请求"));
}
这里的 doOnSubscribe 是整条链路的"点火开关"——前端建立 SSE 连接的那一刻,后端才真正开始干活。
activateGeneration() 做两件事:启动租约续期、启动执行链路:
// 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:定时续期
// 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:续期失败自动停止
// 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
这是整条链路最关键的方法,把"准备执行计划 → 选择执行器 → 消费模型输出 → 收尾"串成一条完整的响应式流:
// BusinessChatService.java —— 组装完整的对话执行流
/**
* 组装完整的对话执行流。
* <p>
* 链路大致分为:发送“正在分析”提示 -> 准备执行计划 -> 按计划选择执行器 -> 消费模型输出 ->
* 正常完成时收尾,异常时失败收尾。
* </p>
*/
private Flux<String> 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));
}
这段代码的执行顺序可以用下面这张图来理解:
编排器:prepareExecutionPlan
执行计划的生成由 ChatPreparationOrchestrator 负责,它是整条链路的"大脑",决定了这次问答到底走哪条路。
// BusinessChatService.java —— 准备执行计划
/**
* 准备本轮会话真正执行所需的编排计划。
* <p>
* 这一层会调用编排器分析历史上下文、决定执行模式、构造 agentQuestion,并在必要时刷新会话绑定的文档范围。
* </p>
*/
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() 方法做了很多事,我们拆开来看核心逻辑:
// 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;
}
// 文档问答模式 → 需要问题改写 + 知识路由 + 执行模式判定
// ...(后续文档会详细展开)
}
编排器的路由决策可以用这张图来概括:
buildAgentQuestion:构造最终喂给 Agent 的问题
编排器生成执行计划之后,还需要把用户的原始问题"包装"一下,加上时间锚点和历史摘要,让 Agent 有足够的上下文来回答:
// BusinessChatService.java —— 构造 Agent 问题
/**
* 构造最终发给 Agent 的问题文本。
* <p>
* 这里会把时间锚点、时效性约束、历史摘要和原始问题拼接成统一 prompt,
* 让 Agent 在处理“今天/最新/本周”等表达时有明确基准。
* </p>
*/
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 负责找到对应的执行器。这里用的是经典的策略模式 + 注册表:
// ConversationExecutorRegistry.java —— 执行器注册表
@Component
public class ConversationExecutorRegistry {
// 用 EnumMap 存储执行模式到执行器的映射,查找效率 O(1)
private final Map<ExecutionMode, ConversationExecutor> executorMap =
new EnumMap<>(ExecutionMode.class);
// Spring 会自动注入所有实现了 ConversationExecutor 接口的 Bean
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;
}
}
所有执行器都实现同一个接口:
// ConversationExecutor.java —— 统一执行器接口
public interface ConversationExecutor {
ExecutionMode mode(); // 声明自己处理哪种模式
Flux<String> execute(TaskInfo taskInfo); // 执行并返回文本流
}
目前系统支持的执行模式有:
| 执行模式 | 对应执行器 | 适用场景 |
|---|---|---|
REACT_AGENT |
ReactAgentExecutor | 开放式提问,支持工具调用和联网搜索 |
RETRIEVAL |
RagChatExecutor | 文档知识问答,走 RAG 检索链路 |
GRAPH_ONLY |
GraphOnlyExecutor | 纯图查询,适合结构化导航类问题 |
GRAPH_THEN_EVIDENCE |
GraphThenEvidenceExecutor | 先图查询再补充证据 |
CLARIFICATION |
ClarificationExecutor | 文档范围歧义时,向用户确认 |
这个设计的好处是:新增一种聊天模式,只需要实现 ConversationExecutor 接口并注册为 Spring Bean,注册表会自动发现它,完全不需要改动已有代码。
流式输出:emitModelChunk
执行器工作过程中,模型每产出一小段文字,就会通过 Flux 推出来,然后被 emitModelChunk() 处理:
// BusinessChatService.java —— 处理模型输出的单个增量片段
/**
* 处理模型输出的单个增量片段。
* <p>
* 每收到一个 chunk,就同时做三件事:追加答案缓冲区、记录首包耗时、向前端发送文本事件。
* </p>
*/
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 事件格式:
// 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<SearchReference> references,
StreamEventMetadata metadata) {
Map<String, Object> payload = event("reference", references, metadata);
payload.put("count", references != null ? references.size() : 0);
return write(payload);
}
// 推荐事件:生成的推荐追问
public String recommendations(List<String> recommendations,
StreamEventMetadata metadata) {
Map<String, Object> payload = event("recommend", recommendations, metadata);
payload.put("count", recommendations != null ? recommendations.size() : 0);
return write(payload);
}
// 统一事件结构:type + content + timestamp + 会话元信息
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;
}
}
前端收到的每条 SSE 事件大概长这样:
{
"type": "text",
"content": "根据文档内容,",
"timestamp": "2025-05-20T10:30:00.123Z",
"conversationId": "abc123",
"exchangeId": 42
}
成功收尾:finishSuccessfully
模型输出完毕后,doOnComplete 触发成功收尾逻辑。这一步要做的事情不少:
// BusinessChatService.java —— 成功收尾
/**
* 处理会话成功完成时的统一收尾。
* <p>
* 这个方法位于“模型正文已经正常输出完成”之后,是一次成功对话真正结束前的最后一道总收口。
* 它承担的不是单一动作,而是一整套按顺序执行的完成态闭环:
* 1. 通过 CAS 把任务标记为 finalized,确保成功收尾只执行一次;
* 2. 从运行态缓冲区中冻结最终答案、引用和推荐追问所需的数据快照;
* 3. 在追踪体系中开启 finalize/recommendation 阶段,便于调试和耗时分析;
* 4. 生成或提取推荐追问,并把推荐阶段标记为完成;
* 5. 向前端补发“引用”和“推荐追问”事件;
* 6. 关闭 SSE 流,告诉前端这一轮输出已经彻底结束;
* 7. 以 {@link ChatTurnStatus#COMPLETED} 状态把最终结果完整落库;
* 8. 异步刷新会话摘要,并清理租约、订阅、运行态注册表等临时资源。
* </p>
* <p>
* 这里的顺序不能随意打乱。特别是:
* “补发事件”必须发生在“关闭 SSE 流”之前,否则前端会收不到引用和推荐;
* “落库和清理”要放在 finally 中,保证即使补发事件失败,也不会让会话停留在未收尾状态。
* </p>
*/
private void finishSuccessfully(TaskInfo taskInfo) {
// finalized 从 false 置为 true 说明当前线程拿到了“唯一一次成功收尾权”;
// 如果这里失败,表示别的线程已经做过停止/失败/成功收尾,本次直接退出避免重复落库。
if (!taskInfo.finalized().compareAndSet(false, true)) {
return;
}
// answer 是最终要持久化的完整回答文本;
// uniqueReferences 先对运行态引用做快照再去重,避免后续落库和前端展示出现重复证据。
String answer = taskInfo.answerBuffer().toString();
List<SearchReference> 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<String> 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
最后一步是释放所有运行态资源:
// 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-落库向量化收尾 | 01-从用户提问到答案返回的总流程 | ➡️ 02-前后端模块划分与调用关系
💬 评论