AI Code Mother 项目学习笔记(八):SSE / Flux 流式输出机制
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 负责处理 HTML 和 MULTI_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_RESPONSE 和 TOOL_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_start、step_completed、workflow_completed 或 workflow_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 回复,并在流结束后保存到对话历史。
这解决了一个很实际的问题:
前端需要实时内容
数据库需要完整记录
SimpleTextStreamHandler 和 JsonMessageStreamHandler 分别处理简单文本流和复杂 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 消息流和工具调用信息。与此同时,工作流模块还提供了 executeWorkflowWithFlux 和 executeWorkflowWithSse 两种进度流式输出方式。
这一期可以用一句话总结:
SSE / Flux 的作用,是让 AI 代码生成过程从“等待最终结果”,变成“实时看到生成内容和执行进度”。
下一期继续分析:
AI Code Mother 项目学习笔记(九):代码解析与文件保存机制
下一期会重点分析项目如何把大模型输出的代码内容解析成 HTML、多文件或 Vue 工程文件,并保存到对应的应用目录中。
更多推荐




所有评论(0)