从用户提问到答案返回的总流程

这篇文档会把 Super Agent 聊天系统后端的完整链路拆开来讲,从用户点击"发送"的那一刻开始,一直到答案流式输出到前端、最后落库收尾为止。每个关键步骤都会贴出对应的源码,加上注释说明它在整条链路里的作用。

看完这篇,你会对"一个问题是怎么从 Controller 一路走到模型输出再回到前端"有一个完整的认知。

总流程概览

先看一张全局流程图,对整条链路有个直观印象:

Fv6c6Yms0QdCgnIwQwnETI8VPGQJ-2202ad22

接下来我们按照这张图的顺序,逐步拆解每个阶段的源码。

入口: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 事件都通过它推给前端
  • answerBufferStringBuffer(线程安全),因为模型输出和收尾落库可能在不同线程
  • finalizedAtomicBoolean,通过 CAS 保证停止/成功/失败的收尾逻辑只执行一次
  • thinkingStepsreferencesusedTools 都用线程安全容器,因为执行过程中多个组件会并发写入

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

这段代码的执行顺序可以用下面这张图来理解:

FqWVKz1q2uwuYAn89geDA_JgIhwL-e3af166c

编排器: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;
    }

    // 文档问答模式 → 需要问题改写 + 知识路由 + 执行模式判定
    // ...(后续文档会详细展开)
}

编排器的路由决策可以用这张图来概括:

FjrgQP3AMZoVuMwQffj9JXO1dDhE-789354b4

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-前后端模块划分与调用关系