1. 本期目标

上一篇文章分析了 AI Code Mother 中的 LangGraph4j 工作流编排机制。我们已经知道,项目把 AI 代码生成拆成了图片收集、提示词增强、智能路由、代码生成、代码质检和项目构建等多个节点。

这一期继续分析一个非常重要的工程细节:

AI 生成内容和工作流执行进度是如何实时返回给前端的?

在 AI Code Mother 中,项目主要使用了两类流式输出方式:

Flux
SSE

其中,Flux 来自 Project Reactor,可以表示一个异步的、连续产生数据的序列。官方文档中也说明,Flux<T> 表示 0 到 N 个异步元素组成的序列,并以完成信号或错误信号结束。(projectreactor.io)

而 SSE,即 Server-Sent Events,可以让服务端通过 text/event-stream 持续向前端推送数据。Spring 官方文档也说明,应用可以返回 Flux<ServerSentEvent> 或类似的多值流式结果。(Home)

本期主要解决几个问题:

1. 为什么 AI 代码生成需要流式返回?
2. Flux 在项目中负责什么?
3. SSE 在项目中负责什么?
4. AppController 如何返回 AI 生成内容?
5. StreamHandlerExecutor 如何处理不同类型的流?
6. 工作流进度如何通过 Flux / SseEmitter 返回?
7. 当前实现有什么值得学习和改进的地方?

2. 为什么 AI 代码生成需要流式返回?

普通接口一般是这样的:

前端发送请求
    ↓
后端执行完整逻辑
    ↓
一次性返回最终结果

这种方式适合查询用户列表、获取文章详情、提交表单等场景。

但是 AI 代码生成不一样。它可能需要较长时间,因为后端需要:

调用大模型
等待模型逐步生成内容
解析模型输出
保存代码文件
记录对话历史
可能还要调用工具或构建项目

如果等所有内容生成完再返回,用户会长时间看不到任何反馈。

所以更好的体验是:

AI 生成一点
    ↓
后端返回一点
    ↓
前端显示一点

也就是边生成、边返回、边展示。

这就是项目中使用 Flux 和 SSE 的核心原因。


3. Flux 和 SSE 的关系

这两个概念容易混在一起,但其实它们负责的层次不同。

可以简单理解为:

Flux:后端内部的数据流表示方式
SSE:后端把数据流推送给前端的 HTTP 输出方式

也就是说:

AI 模型生成 chunk
    ↓
后端用 Flux<String> 表示连续的 chunk
    ↓
Controller 把 chunk 包装成 ServerSentEvent
    ↓
浏览器通过 SSE 持续接收

Flux 更偏向后端编程模型,表示“这里会陆续产生多个数据”;SSE 更偏向网络传输协议,表示“服务端把这些数据持续推送给客户端”。


4. AppController 中的流式接口

项目中最核心的 AI 代码生成接口在:

src/main/java/com/aicode/codemother/controller/AppController.java

对应接口是:

@GetMapping(value = "/chat/gen/code", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
@RateLimit(limitType = RateLimitType.USER, rate = 5, rateInterval = 60, message = "AI 对话请求过于频繁,请稍后再试")
public Flux<ServerSentEvent<String>> chatToGenCode(
        @RequestParam Long appId,
        @RequestParam String message,
        HttpServletRequest request
) {
    ...
}

这里有两个关键信息。

第一,接口返回类型是:

Flux<ServerSentEvent<String>>

这表示后端不是返回一个普通 JSON,而是返回一串 SSE 事件。

第二,接口声明了:

produces = MediaType.TEXT_EVENT_STREAM_VALUE

这表示响应内容类型是 text/event-stream,也就是 SSE 流式响应。源码中 AppController/chat/gen/code 接口正是这样定义的。(GitHub)


5. Controller 层如何包装 SSE 事件?

chatToGenCode 方法中,Controller 先调用 Service 层:

Flux<String> contentFlux = appService.chatToGenCode(appId, message, loginUser);

这里拿到的是普通的 Flux<String>,也就是 AI 生成内容的字符串流。

然后 Controller 会把每个 chunk 包装成 SSE 事件:

return contentFlux
        .map(chunk -> {
            Map<String, String> wrapper = Map.of("d", chunk);
            String jsonData = JSONUtil.toJsonStr(wrapper);
            return ServerSentEvent.<String>builder()
                    .data(jsonData)
                    .build();
        })
        .concatWith(Mono.just(
                ServerSentEvent.<String>builder()
                        .event("done")
                        .data("")
                        .build()
        ));

这段代码说明,后端并不是直接把模型输出的字符串原样发给前端,而是先包装成:

{
  "d": "当前生成内容"
}

然后通过 ServerSentEvent 返回。等整个流结束后,还会追加一个 done 事件,告诉前端生成已经完成。源码中可以看到,项目确实使用 map 包装每个 chunk,并用 concatWith 在最后追加 event("done")。(GitHub)


6. 前端为什么需要 done 事件?

在流式输出中,前端会不断收到内容,例如:

chunk1
chunk2
chunk3
chunk4
...

但是前端还需要知道:

什么时候生成结束?
什么时候可以关闭 loading?
什么时候可以允许用户继续操作?
什么时候可以刷新预览或启用部署按钮?

所以最后的 done 事件很重要。

它的作用不是显示内容,而是给前端一个明确的结束信号。

可以理解为:

普通 data 事件:用于显示 AI 生成内容
done 事件:用于标记本轮生成结束

这种设计比只依赖连接关闭更清晰。因为前端可以直接监听 done 事件,然后执行后续逻辑。


7. Service 层如何产生 Flux?

Controller 只是负责包装 SSE,真正生成 Flux<String> 的地方在:

AppServiceImpl.chatToGenCode()

这个方法会完成参数校验、查询应用、权限校验、读取代码生成类型、保存用户消息和设置监控上下文。之后,它调用:

Flux<String> codeStream = aiCodeGeneratorFacade.generateAndSaveCodeStream(
        message,
        codeGenTypeEnum,
        appId
);

这一步才是真正获得 AI 代码生成流。接着,Service 层不会直接返回 codeStream,而是交给 StreamHandlerExecutor 处理:

return streamHandlerExecutor.doExecute(
        codeStream,
        chatHistoryService,
        appId,
        loginUser,
        codeGenTypeEnum
).doFinally(signalType -> {
    MonitorContextHolder.clearContext();
});

源码中可以看到,Service 层会在调用 AI 前保存用户消息,随后调用 AiCodeGeneratorFacade.generateAndSaveCodeStream() 获得流式结果,再通过 StreamHandlerExecutor 收集 AI 响应并在流结束后保存对话历史,最后通过 doFinally 清理监控上下文。(GitHub)

这个过程可以概括为:

接收用户消息
    ↓
保存用户消息
    ↓
调用 AI 生成流
    ↓
处理流式内容
    ↓
返回给 Controller
    ↓
流结束后保存 AI 回复

8. 为什么不能直接返回原始 codeStream?

从表面看,aiCodeGeneratorFacade.generateAndSaveCodeStream() 已经返回了 Flux<String>,似乎可以直接返回给 Controller。

但项目没有这样做,而是增加了:

StreamHandlerExecutor

原因是:项目不仅要把内容返回给前端,还要在后端保存完整 AI 回复到对话历史中。

流式响应的特点是内容被拆成很多小块:

chunk1 + chunk2 + chunk3 + ...

前端可以边收边显示,但数据库里一般需要保存完整的一条 AI 回复。

所以后端需要一边把 chunk 放行给前端,一边把 chunk 收集起来。等流结束后,再拼接成完整内容保存到 chat_history 表中。

这就是流处理器的作用。


9. StreamHandlerExecutor:根据生成类型选择处理器

StreamHandlerExecutor 位于:

src/main/java/com/aicode/codemother/core/handler/StreamHandlerExecutor.java

它的职责很明确:根据代码生成类型选择不同的流处理器。

源码注释说明:

HTML、MULTI_FILE 使用 SimpleTextStreamHandler
VUE_PROJECT 使用 JsonMessageStreamHandler

对应代码逻辑是:

return switch (codeGenType) {
    case VUE_PROJECT -> jsonMessageStreamHandler.handle(...);
    case HTML, MULTI_FILE -> new SimpleTextStreamHandler().handle(...);
};

这说明项目认为不同生成模式的流式内容格式不同,所以不能用同一个处理方式。StreamHandlerExecutor 明确把 HTML / MULTI_FILE 交给简单文本处理器,把 VUE_PROJECT 交给 JSON 消息处理器。(GitHub)


10. SimpleTextStreamHandler:处理简单文本流

SimpleTextStreamHandler 负责处理 HTMLMULTI_FILE 类型。

它的逻辑比较简单:

StringBuilder aiResponseBuilder = new StringBuilder();

return originFlux
        .map(chunk -> {
            aiResponseBuilder.append(chunk);
            return chunk;
        })
        .doOnComplete(() -> {
            String aiResponse = aiResponseBuilder.toString();
            chatHistoryService.addChatMessage(...);
        })
        .doOnError(error -> {
            String errorMessage = "AI回复失败: " + error.getMessage();
            chatHistoryService.addChatMessage(...);
        });

这段代码有两个关键点。

第一,map 阶段会把每个 chunk 追加到 StringBuilder,同时原样返回给前端。

第二,doOnComplete 阶段会在流结束后,把完整 AI 回复保存到对话历史中。

源码中 SimpleTextStreamHandler 明确说明它用于处理 HTML 和 MULTI_FILE 类型的流式响应,并直接收集完整文本响应。(GitHub)

可以理解为:

AI 生成一个 chunk
    ↓
后端收集这个 chunk
    ↓
同时把这个 chunk 返回给前端
    ↓
流结束后保存完整回复

11. JsonMessageStreamHandler:处理 Vue 项目的复杂流

VUE_PROJECT 的流式输出更复杂,因为它可能包含工具调用信息。

所以项目使用:

JsonMessageStreamHandler

它位于:

src/main/java/com/aicode/codemother/core/handler/JsonMessageStreamHandler.java

源码注释说明,它用于处理 VUE_PROJECT 类型的复杂流式响应,并且包含工具调用信息。它会解析每个 JSON 消息块,然后重组为前端需要的响应格式。(GitHub)

它的核心处理逻辑大致是:

return originFlux
        .map(chunk -> {
            return handleJsonMessageChunk(chunk, chatHistoryStringBuilder, seenToolIds);
        })
        .filter(StrUtil::isNotEmpty)
        .doOnComplete(() -> {
            String aiResponse = chatHistoryStringBuilder.toString();
            chatHistoryService.addChatMessage(...);
        })
        .doOnError(error -> {
            String errorMessage = "AI回复失败: " + error.getMessage();
            chatHistoryService.addChatMessage(...);
        });

相比 SimpleTextStreamHandler,它多了几个动作:

解析 JSON 消息
识别消息类型
处理 AI 普通回复
处理工具调用请求
过滤空字符串
保存整理后的完整 AI 回复

源码中可以看到,它会根据 StreamMessageTypeEnum 区分 AI_RESPONSETOOL_REQUEST。如果是 AI 回复,就拼接响应文本并返回;如果是工具调用请求,就根据工具名称获取工具实例,并返回格式化后的工具调用信息。(GitHub)


12. 为什么 Vue 项目要用 JSON 消息流?

HTML 和 MULTI_FILE 的生成相对简单,模型主要输出代码文本。

但是 Vue 项目生成可能涉及:

创建文件
修改文件
读取目录
调用工具
生成多个工程文件
处理工具调用结果

如果这些内容都混在普通字符串里,前端和后端都很难区分:

这是 AI 正常回复?
这是工具调用?
这是工具调用结果?
这是要展示给用户的内容?
还是只是内部执行信息?

所以 Vue 项目采用 JSON 消息流更合理。

可以理解为:

简单代码生成:文本流就够了
复杂工程生成:需要结构化消息流

这也是为什么 StreamHandlerExecutor 会根据 CodeGenTypeEnum 选择不同处理器。


13. AI 代码生成流式链路总结

到这里,正式 AI 生成接口的流式链路可以总结为:

前端请求 /api/app/chat/gen/code
    ↓
AppController.chatToGenCode()
    ↓
AppServiceImpl.chatToGenCode()
    ↓
AiCodeGeneratorFacade.generateAndSaveCodeStream()
    ↓
得到 Flux<String>
    ↓
StreamHandlerExecutor
    ├─ HTML / MULTI_FILE → SimpleTextStreamHandler
    └─ VUE_PROJECT → JsonMessageStreamHandler
    ↓
返回处理后的 Flux<String>
    ↓
AppController 包装成 Flux<ServerSentEvent<String>>
    ↓
前端通过 SSE 接收内容
    ↓
最后收到 done 事件

这条链路体现了项目的一个核心设计:

AI 输出既要实时返回给前端,也要在后端被收集、整理和保存。


14. 工作流中的 Flux 流式输出

除了正式的 AI 代码生成接口,项目中的工作流模块也提供了 Flux 流式输出。

CodeGenWorkflow 中,有一个方法:

public Flux<String> executeWorkflowWithFlux(String originalPrompt)

它通过 Flux.create 创建一个流,然后使用虚拟线程执行工作流。在工作流执行过程中,它会不断向 sink 推送事件。源码中可以看到,它会先推送 workflow_start,然后在每个步骤完成后推送 step_completed,最后推送 workflow_completed。如果发生异常,则推送 workflow_error。(GitHub)

流程可以理解为:

工作流开始
    ↓
发送 workflow_start
    ↓
执行节点 1
    ↓
发送 step_completed
    ↓
执行节点 2
    ↓
发送 step_completed
    ↓
全部完成
    ↓
发送 workflow_completed

这种流式输出不是为了返回 AI 生成的代码内容,而是为了返回“工作流执行进度”。


15. 工作流事件的格式

CodeGenWorkflow 中有一个辅助方法:

private String formatSseEvent(String eventType, Object data) {
    String jsonData = JSONUtil.toJsonStr(data);
    return "event: " + eventType + "\ndata: " + jsonData + "\n\n";
}

这说明工作流的 Flux 版本虽然返回的是 Flux<String>,但每个字符串已经被格式化成 SSE 事件格式。源码中 formatSseEvent 会把事件名和 JSON 数据拼接成 event: ...\ndata: ...\n\n 的形式。(GitHub)

例如,一个事件可能是:

event: step_completed
data: {"stepNumber":1,"currentStep":"图片收集完成"}

这样前端就可以根据 event 字段区分不同类型的工作流事件。


16. 工作流中的 SseEmitter 输出

除了 Flux 版本,CodeGenWorkflow 还提供了一个基于 SseEmitter 的版本:

public SseEmitter executeWorkflowWithSse(String originalPrompt)

这个方法会创建一个 30 分钟超时的 SseEmitter,然后使用虚拟线程执行工作流。执行过程中,它会通过 sendSseEvent 发送 workflow_startstep_completedworkflow_completedworkflow_error 等事件。(GitHub)

sendSseEvent 的核心逻辑是:

emitter.send(SseEmitter.event()
        .name(eventType)
        .data(data));

源码中也可以看到,如果发送失败,会记录异常日志。(GitHub)

可以简单理解为:

Flux 版本:用响应式数据流推送事件
SseEmitter 版本:用 Spring MVC 的 SSE 发送器推送事件

两者目的相似,都是让前端可以实时看到工作流执行进度。


17. WorkflowSseController:工作流流式接口

项目中还有一个专门的控制器:

WorkflowSseController

它位于:

src/main/java/com/aicode/codemother/controller/WorkflowSseController.java

这个控制器提供了三个接口:

POST /workflow/execute
GET  /workflow/execute-flux
GET  /workflow/execute-sse

其中,/execute 是同步执行工作流;/execute-flux 返回 Flux<String>,并声明 produces = MediaType.TEXT_EVENT_STREAM_VALUE/execute-sse 返回 SseEmitter,同样声明为 text/event-stream。源码注释也说明,这个控制器用于演示 LangGraph4j 工作流的流式输出功能。(GitHub)

可以理解为:

/workflow/execute
    用于普通同步测试

/workflow/execute-flux
    用于 Flux 流式测试

/workflow/execute-sse
    用于 SseEmitter 流式测试

这说明项目不仅对正式代码生成做了流式响应,也专门为工作流执行进度做了流式输出示例。


18. AI 内容流和工作流进度流的区别

这里需要区分两种流。

第一种是 AI 内容流:

用户输入需求
    ↓
模型生成代码
    ↓
后端持续返回代码 chunk

它的核心目的是让用户实时看到 AI 生成的内容。

第二种是工作流进度流:

工作流开始
    ↓
图片收集完成
    ↓
提示词增强完成
    ↓
代码生成完成
    ↓
质检完成
    ↓
构建完成

它的核心目的是让用户实时看到当前执行到哪一步。

所以两者虽然都使用流式输出,但关注点不同:

AI 内容流:返回生成内容
工作流进度流:返回执行状态

在一个更完整的 AI 代码生成平台中,这两类流可以结合起来:

左侧显示工作流进度
右侧显示 AI 生成内容
底部显示工具调用和构建日志

这样用户体验会更好。


19. 这个设计有什么优点?

我认为项目的流式输出机制有三个值得学习的地方。

19.1 用户体验更好

AI 生成代码可能需要几十秒甚至更久。

如果没有流式输出,用户只能看到 loading,无法判断系统是否还在工作。

有了 SSE / Flux 后,用户可以持续看到内容更新或步骤进度,等待体验会好很多。


19.2 后端可以边返回边收集

项目不是简单地把模型输出转发给前端,而是在流处理器中收集完整 AI 回复,并在流结束后保存到对话历史。

这解决了一个很实际的问题:

前端需要实时内容
数据库需要完整记录

SimpleTextStreamHandlerJsonMessageStreamHandler 分别处理简单文本流和复杂 JSON 消息流,使得不同代码生成模式都可以被正确保存。


19.3 简单任务和复杂任务分开处理

HTML / MULTI_FILE 使用简单文本处理器。

VUE_PROJECT 使用 JSON 消息处理器。

这个设计很合理,因为不同任务的流式输出结构不同。如果强行统一处理,代码会变得混乱。


20. 当前实现中需要注意的问题

当前实现也有一些可以进一步优化的地方。

20.1 Controller 包装格式比较简单

目前 AppController 把每个 chunk 包装成:

{
  "d": "内容"
}

这种格式比较轻量,但信息较少。

后续可以考虑增加更明确的事件类型,例如:

{
  "type": "ai_delta",
  "data": "内容",
  "timestamp": 123456789
}

这样前端处理会更清晰。


20.2 AI 内容流和工具调用流可以统一事件协议

当前 Vue 项目通过 JsonMessageStreamHandler 处理工具调用信息,而 HTML / MULTI_FILE 主要是普通文本流。

后续可以设计统一的事件协议,例如:

ai_delta:AI 文本增量
tool_start:工具调用开始
tool_result:工具调用结果
file_saved:文件保存完成
error:错误信息
done:生成完成

这样前端无论面对哪种生成模式,都可以用一套事件处理逻辑。


20.3 工作流模块和正式业务流还可以进一步打通

目前 WorkflowSseController 更像是工作流流式输出的演示接口,而正式 AI 生成接口仍然走 AppController → AppServiceImpl → AiCodeGeneratorFacade 这条链路。

后续可以考虑把工作流进度事件也整合进正式生成接口,让用户不仅看到 AI 生成内容,还能看到:

正在分析需求
正在收集图片
正在增强提示词
正在生成代码
正在检查代码质量
正在构建项目

这样系统会更像一个完整的 AI 应用生成平台。


21. 本期重点理解

这一期最重要的是理解:AI Code Mother 使用流式输出来提升 AI 生成体验。

可以总结为五点:

第一,Flux 是后端内部表示异步数据流的方式。
第二,SSE 是服务端向前端持续推送数据的方式。
第三,AppController 使用 Flux<ServerSentEvent<String>> 返回 AI 生成内容。
第四,StreamHandlerExecutor 根据代码生成类型选择不同流处理器。
第五,工作流模块也提供 Flux 和 SseEmitter 两种进度流式输出方式。

这个设计体现了 AI 应用开发中的一个重要思想:

AI 生成过程不能只关注最终结果,也要关注生成过程的实时反馈。


22. 我的理解

我认为流式输出是这个项目从“普通后端系统”走向“AI 应用系统”的关键设计之一。

普通 CRUD 系统通常只关心:

请求是否成功?
返回结果是什么?

而 AI 代码生成系统还需要关心:

模型现在生成到哪里了?
用户是否能实时看到内容?
中间过程是否能保存?
工具调用是否能展示?
工作流步骤是否能反馈?

所以,SSE / Flux 不只是一个接口写法问题,而是 AI 产品体验的一部分。

这个项目通过 Flux 表示模型输出流,通过 ServerSentEvent 返回前端,通过 StreamHandlerExecutor 保存完整对话,再通过 WorkflowSseController 演示工作流进度输出。整体设计比较清晰,也很适合作为学习 AI 后端流式响应的案例。


23. 本期小结

本期主要分析了 AI Code Mother 的 SSE / Flux 流式输出机制。

在正式 AI 代码生成接口中,AppController 返回 Flux<ServerSentEvent<String>>,并将 Service 层返回的 Flux<String> 包装成 SSE 事件;AppServiceImpl 负责生成代码流,并交给 StreamHandlerExecutor 根据代码类型选择处理器;SimpleTextStreamHandler 处理 HTML 和 MULTI_FILE 的普通文本流;JsonMessageStreamHandler 处理 VUE_PROJECT 的复杂 JSON 消息流和工具调用信息。与此同时,工作流模块还提供了 executeWorkflowWithFluxexecuteWorkflowWithSse 两种进度流式输出方式。

这一期可以用一句话总结:

SSE / Flux 的作用,是让 AI 代码生成过程从“等待最终结果”,变成“实时看到生成内容和执行进度”。

下一期继续分析:

AI Code Mother 项目学习笔记(九):代码解析与文件保存机制

下一期会重点分析项目如何把大模型输出的代码内容解析成 HTML、多文件或 Vue 工程文件,并保存到对应的应用目录中。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐