1.创建智谱Ai客户端
创建的客户端交给Spring Bean容器管理
package com.ldy.yudada.config; import ai.z.openapi.ZhipuAiClient; import lombok.Data; import org.springframework.boot.context.properties.ConfigurationProperties; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @Configuration @ConfigurationProperties(prefix = "ai") @Data public class AiConfig { /** * apiKey,需要从平台获取 */ private String apiKey; // 创建Ai客户端: // 从环境变量读取 API Key @Bean public ZhipuAiClient getClient() { return ZhipuAiClient.builder().ofZHIPU().apiKey(apiKey).build(); } }2.定义Ai的使用
这部分展示的是打开了流式对话的写法
1.定义请求(整合消息)
2.获取响应
3.返回 Flowable<ModelData>对象
/** * 通用流式请求 * * @param messages * @param stream * @param temperature * @return */ public Flowable<ModelData> doStreamRequest(List<ChatMessage> messages, Float temperature) { // 创建聊天完成请求 ChatCompletionCreateParams request = ChatCompletionCreateParams.builder() .model("glm-5.2") .stream(Boolean.TRUE) .temperature(temperature) .messages(messages) .build(); // 发送请求 ChatCompletionResponse response = client.chat().createChatCompletion(request); // 获取回复 if (response.isSuccess() && response.getFlowable() != null) { return response.getFlowable(); } else { throw new BusinessException(ErrorCode.SYSTEM_ERROR, response.getMsg()); } } /** * 通用流式请求(简化消息传递) * * @param systemMessage * @param userMessage * @param temperature * @return */ public Flowable<ModelData> doStreamRequest(String systemMessage, String userMessage, Float temperature) { // 创建聊天完成请求 List<ChatMessage> chatMessageList = new ArrayList<>(); ChatMessage systemChatMessage = new ChatMessage(ChatMessageRole.SYSTEM.value(), systemMessage); chatMessageList.add(systemChatMessage); ChatMessage userChatMessage = new ChatMessage(ChatMessageRole.USER.value(), userMessage); chatMessageList.add(userChatMessage); return doStreamRequest(chatMessageList, temperature); } }3.编写接口
其中getGenerateQuestionUserMessage()是自己定义的拼接用户信息的方法
GENERATE_QUESTION_SYSTEM_MESSAGE是自定义的系统信息
以上都是用于传递给Ai
@GetMapping("/ai_generate/sse") public SseEmitter aiGenerateQuestionSSE(AiGenerateQuestionRequest aiGenerateQuestionRequest) { ThrowUtils.throwIf(aiGenerateQuestionRequest == null, ErrorCode.PARAMS_ERROR); // 获取参数 Long appId = aiGenerateQuestionRequest.getAppId(); int questionNumber = aiGenerateQuestionRequest.getQuestionNumber(); int optionNumber = aiGenerateQuestionRequest.getOptionNumber(); // 获取应用信息 App app = appService.getById(appId); ThrowUtils.throwIf(app == null, ErrorCode.NOT_FOUND_ERROR); // 封装 Prompt(生成题目的用户消息的Prompt) String userMessage = getGenerateQuestionUserMessage(app, questionNumber, optionNumber); // 建立SSE连接对象, 0 表示无超时时间 SseEmitter sseEmitter = new SseEmitter(0L); // AI 生成 , SSE流式返回 Flowable<ModelData> modelDataFlowable = aiManager.doStreamRequest(GENERATE_QUESTION_SYSTEM_MESSAGE, userMessage, null); // 左括号计数器,当回归为0时,相当于左括号等于右括号,可以截取 AtomicInteger counter = new AtomicInteger(0); // 拼接完整题目 StringBuilder stringBuilder = new StringBuilder(); modelDataFlowable .observeOn(Schedulers.io()) .map(modelData -> { if (modelData.getChoices() == null || modelData.getChoices().isEmpty()) { return ""; } String content = modelData.getChoices().get(0).getDelta().getContent(); return content != null ? content : ""; }) .map(message -> message.replaceAll("\\s", "")) .filter(StrUtil::isNotBlank) .flatMap(message -> { List<Character> characterList = new ArrayList<>(); for (char c : message.toCharArray()) { characterList.add(c); } return Flowable.fromIterable(characterList); }) .doOnNext(c -> { // 如果是“{”,则计数器加一 ,反之减一 if (c == '{') { counter.addAndGet(1); } if (counter.get() > 0) { stringBuilder.append(c); } if (c == '}') { counter.addAndGet(-1); if (counter.get() == 0) { // 可以拼接题目,并且通过SSE返回给前端 sseEmitter.send(JSONUtil.toJsonStr(stringBuilder.toString())); // 重置,准备拼接下一道题 stringBuilder.setLength(0); } } }) .doOnError((e) -> log.error("sse error" + e)) .doOnComplete(sseEmitter::complete) .subscribe(); return sseEmitter; }4.Rxjava,SSE技术向前端推送AI结果核心部分
// 建立SSE连接对象, 0 表示无超时时间 SseEmitter sseEmitter = new SseEmitter(0L); // AI 生成 , SSE流式返回 Flowable<ModelData> modelDataFlowable = aiManager.doStreamRequest(GENERATE_QUESTION_SYSTEM_MESSAGE, userMessage, null); // 左括号计数器,当回归为0时,相当于左括号等于右括号,可以截取 AtomicInteger counter = new AtomicInteger(0); // 拼接完整题目 StringBuilder stringBuilder = new StringBuilder(); modelDataFlowable .observeOn(Schedulers.io()) .map(modelData -> { if (modelData.getChoices() == null || modelData.getChoices().isEmpty()) { return ""; } String content = modelData.getChoices().get(0).getDelta().getContent(); return content != null ? content : ""; }) .map(message -> message.replaceAll("\\s", "")) .filter(StrUtil::isNotBlank) .flatMap(message -> { List<Character> characterList = new ArrayList<>(); for (char c : message.toCharArray()) { characterList.add(c); } return Flowable.fromIterable(characterList); }) .doOnNext(c -> { // 如果是“{”,则计数器加一 ,反之减一 if (c == '{') { counter.addAndGet(1); } if (counter.get() > 0) { stringBuilder.append(c); } if (c == '}') { counter.addAndGet(-1); if (counter.get() == 0) { // 可以拼接题目,并且通过SSE返回给前端 sseEmitter.send(JSONUtil.toJsonStr(stringBuilder.toString())); // 重置,准备拼接下一道题 stringBuilder.setLength(0); } } }) .doOnError((e) -> log.error("sse error" + e)) .doOnComplete(sseEmitter::complete) .subscribe(); return sseEmitter; }一、这段代码是干嘛的
这是一个AI 生成题目的 SSE 流式推送接口。 前端点击「AI 生成题目」后,后端不会等所有题目全部生成完再一次性返回,而是AI 生成一点,后端就推一点,前端可以实时看到题目逐道出现的效果(类似打字机)。
技术 1:
SseEmitter—— Spring 提供的 SSE长连接工具,负责保持和前端的连接,持续推送数据技术 2:
Flowable(RxJava)—— 负责处理 AI 返回的流式数据,做转换、拆分、拼接核心业务:把 AI 流式吐出的零散文本,拼成一道道完整的题目 JSON,再逐道推给前端
1.新建SseEmitter,相当于和前端建立一条「一直开着的传输通道」
SseEmitter sseEmitter = new SseEmitter(0L);2.从aiManager.doStreamRequest开始,进入 RxJava 流式处理(得到Flowable<ModelData>对象)
Flowable<ModelData> modelDataFlowable = aiManager.doStreamRequest(GENERATE_QUESTION_SYSTEM_MESSAGE, userMessage, null);3.流式处理数据部分
1.切换到 RxJava 内置的 IO 线程池
.observeOn(Schedulers.io())这行代码的作用是:把这行之后所有的流处理逻辑(map、flatMap、doOnNext 等),全部切换到 RxJava 内置的 IO 线程池里执行,避免阻塞原请求线程。
一、拆成两部分理解
1.observeOn:切换下游的执行线程
observeOn是 RxJava 的线程切换操作符,核心规则:
只影响它「之后」的所有下游操作,不影响上游
写在哪里,就从哪里开始切线程;后面再写一次
observeOn可以再次切换
举个直观例子:
上游数据源 .observeOn(线程A) // 从这里开始,后面的逻辑全跑在线程A .map(...) .filter(...) .observeOn(线程B) // 从这里开始,后面的逻辑又切到线程B .doOnNext(...) .subscribe();对应代码:这行写在所有 map、flatMap、doOnNext 之前,所以后面所有数据加工、字符拼接、SSE 推送的逻辑,全部都会跑在 IO 线程里。
2.Schedulers.io():IO 专用调度器
Schedulers是 RxJava 提供的线程池工具,内置了几种常用的调度器
2.获取数据
.map(modelData -> { if (modelData.getChoices() == null || modelData.getChoices().isEmpty()) { return ""; } String content = modelData.getChoices().get(0).getDelta().getContent(); return content != null ? content : "";作用:从 AI 返回的复杂对象ModelData里,提取出真正有用的「生成文本」。
- AI 原始返回结构很深:
modelData → choices[0] → delta → content才是真正的文字 - 做了判空防御,空的就返回空字符串,避免返回 null 触发 RxJava 空指针异常
3.清洗数据
// ③ 去掉所有空白字符 .map(message -> message.replaceAll("\\s", "")) // ④ 过滤掉空字符串 .filter(StrUtil::isNotBlank)作用:清洗数据。
- AI 生成的 JSON 会有换行、空格、缩进,全部去掉,只保留纯字符,方便后面按括号切割
- 空的片段直接过滤掉,不往下游传,减少无效处理
3.将字符串转拆分为字符再流式输出
// ⑤ 把字符串拆成单个字符,逐个发射 .flatMap(message -> { List<Character> characterList = new ArrayList<>(); for (char c : message.toCharArray()) { characterList.add(c); } return Flowable.fromIterable(characterList); })- 比如收到一段文本
{"title":"xxx",会拆成{"title... 一个字符一个字符地流下去 - 为什么要拆这么细?因为后面要逐字符统计括号,精准判断一道题的 JSON 什么时候闭合
flatMap可以把「一个数据」展开成「多个数据」继续流
- 拆分:把一整段字符串,拆成一个个独立的字符,放进列表里
- 输入
"ab{c}"→ 拆成['a','b','{','c','}']
- 输入
- 包装:把字符列表包装成一个
Flowable字符流,交给flatMapFlowable.fromIterable(列表):把列表里的元素,逐个发射出去
flatMap拿到这个字符流之后,会把它「拍扁」接入主管道。最终效果就是:上游下来一段字符串,下游变成一个一个字符依次流过。
flatMap 的核心职责只有一个:扁平化
你必须在函数里返回一个新的流(Publisher/Flowable),这是硬性要求。
4.代码的核心逻辑 —— 括号计数法
// ⑥ 核心:逐字符处理,拼完整题目,SSE 推送 .doOnNext(c -> { // 遇到左括号,计数器+1 if (c == '{') { counter.addAndGet(1); } // 计数器>0,说明在题目JSON内部,把字符拼进去 if (counter.get() > 0) { stringBuilder.append(c); } // 遇到右括号,计数器-1 if (c == '}') { counter.addAndGet(-1); // 计数器回到0,说明一对{}完全闭合 = 一道题拼完了 if (counter.get() == 0) { // 通过SSE把这道完整的题目推给前端 sseEmitter.send(JSONUtil.toJsonStr(stringBuilder.toString())); // 清空缓冲区,准备拼下一道题 stringBuilder.setLength(0); } } })这是整段代码的核心逻辑 —— 括号计数法。 AI 生成的题目是一个 JSON 数组,格式大概是[{...}, {...}, {...}]。 我们的目标是:每拼完一个完整的{...}(一道题),就立刻推给前端,不用等全部生成完。
逻辑通俗讲:
- 遇到
{,计数 +1(进入一层 JSON 对象) - 只要计数 > 0,就把字符往缓冲区里拼
- 遇到
},计数 -1 - 当计数回到 0,说明从第一个
{到这个}刚好闭合,一道题拼完整了 - 调用
sseEmitter.send()把这道题推给前端,清空缓冲区继续拼下一道
5.订阅
// ⑦ 出错时打日志 .doOnError((e) -> log.error("sse error" + e)) // ⑧ AI生成完毕,关闭SSE连接 .doOnComplete(sseEmitter::complete) // ⑨ 订阅流,真正开始执行 .subscribe();doOnError:流出现异常时执行,这里只打了日志doOnComplete:AI 所有内容生成完、流正常结束时,关闭 SSE 连接subscribe():真正启动这条流。RxJava 是「懒执行」的,不写这一行,前面所有代码都只是定义,不会真正运行
5.总结,Rxjava和SSE是怎么结合在一起使用的
RxJava 负责后端内部的流式数据加工,SSE 负责把加工好的数据推给前端;两者靠sseEmitter.send()这一行代码衔接,一个产、一个发,天然都是「流式」思想,配合起来非常顺滑。
- RxJava:后端内部的异步数据流处理工具,管「数据怎么拆、怎么拼、怎么过滤、怎么切线程」
- SSE:后端 ↔ 前端的通信推送技术,管「怎么把数据持续推给浏览器」
一、先明确各自的分工
1. RxJava(Flowable):内部数据加工厂
它的工作范围完全在后端服务内部:
- 从 AI SDK 接收到原始的流式响应(一小段一小段的文字碎片)
- 用
map提取内容、filter过滤空值、flatMap拆成字符 - 用括号计数器拼出完整的单道题目
- 全程在 IO 线程运行,不阻塞主线程
它只关心「怎么把碎数据加工成可用的业务数据」,不关心数据最终是存数据库、返回接口还是推给前端。
2. SSE(SseEmitter):对外传输通道
它的工作是和前端浏览器打交道:
- 建立一条不关闭的 HTTP 长连接
- 后端随时可以调用
send()往通道里塞数据,前端实时收到 - 调用
complete()主动关闭连接
它只关心「怎么把数据推给前端」,不关心数据是 AI 生成的、数据库查的还是手动拼的。
二、核心结合点:就这一行代码
两者唯一的交集,就是doOnNext里的这一行:
sseEmitter.send(JSONUtil.toJsonStr(stringBuilder.toString()));这就是「加工厂」和「运输通道」的对接窗口:
- RxJava 每加工好一道完整的题目,就调用一次
send() - SSE 接到数据,立刻通过长连接推给前端
- 推完继续等 RxJava 加工下一道
俗类比: RxJava = 包子铺后厨,负责揉面、包馅、蒸包子,蒸好一个放一个到出餐口 SSE = 外卖传送滑道,出餐口放一个,滑道就送一个到顾客手里 结合点 = 出餐口的那个窗口:后厨放进去,滑道传出去通
三、完整数据流走一遍,全程对应
我们把你代码的完整流程,按「谁在干活」标出来,一眼就能看清:
前端发起请求 → SSE 建立连接后端创建
SseEmitter,和前端建立长连接,这一步是 SSE 的活。调用 AI 接口 → 拿到 RxJava 流
aiManager.doStreamRequest()返回Flowable<ModelData>,AI 生成的文字碎片会源源不断地流进这条 RxJava 管道。链式加工数据 → 全是 RxJava 的活
observeOn 切线程 → map 提取内容 → filter 过滤空值 → flatMap 拆字符 → doOnNext 拼题目这一长串全部是 RxJava 在内部处理数据,和 SSE 没有任何关系。
4.关键衔接:拼完一道,推一道每当括号计数器归零、一道题拼接完成,就执行:
sseEmitter.send(...)把 RxJava 加工好的成品,塞进 SSE 通道,推给前端。
5.结束 / 异常,同步关闭连接
- RxJava 流正常结束 →
doOnComplete里调用sseEmitter.complete(),关闭 SSE 连接 - RxJava 流出错 → 调用
sseEmitter.completeWithError(),通知前端异常结束
- RxJava 流正常结束 →
部分代码理解
为什么要做这一步操作(切线程):
modelDataFlowable .observeOn(Schedulers.io())首先这行代码会把它之后所有的流处理逻辑,全部切换到 RxJava 内置的 IO 线程池里运行,避免阻塞原请求线程,也让后续的异步推送逻辑在可控的后台线程执行。
后面所有的 map、filter、flatMap、doOnNext、doOnError、doOnComplete 以及 subscribe 的回调,全部都会跑在 IO 线程里。
怎么避免阻塞的呢?
一、先搞懂 Tomcat 的请求处理模型
Tomcat 处理 HTTP 请求,用的是「线程池 + 一请求一线程」的模型,这是理解所有问题的基础。
核心规则
Tomcat 内部维护了一个工作线程池(默认最大 200 个线程),这是服务处理请求的全部 “人手”。
每进来一个 HTTP 请求,Tomcat 就从池子里分配一条空闲线程,专门用来执行你的 Controller 方法。
只有当 Controller 方法执行完(return 了),这条线程才会被归还回线程池,才能继续处理下一个新请求。
线程池总数是有限的,全部占满后,新请求就只能排队等待;排队也满了就会直接拒绝,表现为网站卡死、请求超时。
阻塞的本质:方法没执行完,线程就一直被占用,不能干别的。
代码是按顺序从上往下执行的,只要你的 Controller 方法还没 return,分配给它的 Tomcat 线程就必须一直等着,不能释放。
如果方法里有耗时操作(比如等待 AI 生成结果、同步调用网络接口、Thread.sleep),线程就会卡在这一步,原地等待任务完成
三、回到你的 SSE 场景:不切线程会发生什么
你的接口是 SSE 流式推送,AI 生成题目可能需要 5~30 秒。如果不用 RxJava 切线程,把所有逻辑都写在 Controller 主线程里同步执行,就会出现:
请求进来,Tomcat 分配一条工作线程执行
aiGenerateQuestionSSE方法调用 AI 接口,开始流式生成
线程原地等待:等 AI 返回一段、拼一道题、推一道,循环往复
直到所有题目全部生成完、SSE 连接关闭,方法才 return
这 5~30 秒里,这条 Tomcat 线程被全程占死,啥也干不了
一个用户占一条线程,100 个同时在线就占 100 条,200 个用户就把线程池占满。后面再来新用户,连页面都打不开。
四、为什么observeOn就能解决阻塞
核心逻辑:把耗时的流式处理逻辑,从 Tomcat 请求线程,转移到 RxJava 的后台 IO 线程去跑,让请求线程可以立刻 return。
对应你的代码执行顺序:
Tomcat 线程执行 Controller 方法:
校验参数、查应用、组装 Prompt
创建
SseEmitter拿到 AI 的 Flowable 流
写好链式处理逻辑,执行
.subscribe()启动流
重点:
subscribe()是异步非阻塞的它只是告诉流 “可以开始了”,然后立刻返回,不会等流全部执行完。 加上.observeOn(Schedulers.io())之后,后续所有 map、拼接、推送逻辑,全部会被放到 RxJava 的 IO 线程池里执行,和当前 Tomcat 线程没关系了。Controller 方法立刻
return sseEmitter方法执行结束,Tomcat 线程立刻归还回线程池,可以马上去处理下一个新请求。后续 AI 生成、拼接题目、SSE 推送,全程在后台 IO 线程里默默运行,不占用任何 Tomcat 工作线程。
6.前端接收SSE推送内容
这段代码是前端 SSE 客户端的原生实现,作用是和后端 AI 生成题目的接口建立长连接,实时接收后端流式推送的题目数据。全程用浏览器自带的EventSourceAPI,不需要安装任何第三方依赖。
const eventSource = new EventSource( "http://localhost:8101/api/question/ai_generate/sse" + `?appId=${props.appId}&optionNumber=${form.optionNumber}&questionNumber=${form.questionNumber}` ); // 接收信息 eventSource.onmessage = function (event) { console.log(event.data); }; // 报错或者连接关闭时出发 eventSource.onerror = function (event) { if (event.eventPhase == EventSource.CLOSED) { console.log("连接关闭"); eventSource.close(); } }; // 连接打开时触发 eventSource.onopen = function (event) { console.log("建立连接"); };const eventSource = new EventSource( "http://localhost:8101/api/question/ai_generate/sse" + `?appId=${props.appId}&optionNumber=${form.optionNumber}&questionNumber=${form.questionNumber}` );1.建立长连接:
- 执行
new EventSource(地址)时,浏览器会自动向这个地址发起一个GET 请求,请求头自动带上Accept: text/event-stream,告诉后端「我要接收事件流」。 - 连接成功后会一直保持不关闭,后端可以随时往这条通道里推数据,前端被动接收。
- 后面拼接的是查询参数:把应用 ID、生成题目数、每题选项数传给后端,AI 会按这个要求生成题目。
- 注意:这里写了完整的后端地址(带
localhost:8101),和前端页面端口不一致,属于跨域请求,后端必须配置 CORS 放行,否则浏览器会直接拦截连接。
2.接收后端推送数据
eventSource.onmessage = function (event) { console.log(event.data); };这是最核心的回调,对应你后端的sseEmitter.send():
- 后端每推送一段数据,这个函数就会自动执行一次。
event.data就是后端发过来的具体内容,也就是你后端拼好的单道题目 JSON 字符串。- 后端每拼完一道完整题目就推一次,所以这里会一道接一道地陆续收到数据。
- 目前代码只做了控制台打印,实际业务里需要把
event.data解析成 JSON 对象,追加到页面的题目列表里,实现「边生成边显示」的效果。
3. 错误与连接关闭处理
eventSource.onerror = function (event) { if (event.eventPhase == EventSource.CLOSED) { console.log("连接关闭"); eventSource.close(); } };- 只要连接出现异常、网络中断、后端主动关闭,都会触发这个回调。
EventSource.CLOSED是浏览器内置常量(值为 2),表示连接已经处于关闭状态。- 逻辑:检测到连接关闭了,就手动调用
close()彻底终止连接。
补充一个原生特性:
原生
EventSource默认自带自动重连。如果是网络波动导致的意外断开,浏览器会隔几秒自动尝试重新连接,不需要你手写重连逻辑;只有后端主动正常关闭、或者手动调用close(),才会彻底停止。
4. 连接建立成功回调
eventSource.onopen = function (event) { console.log("建立连接"); };- 当后端返回了正确的响应头(
Content-Type: text/event-stream)、长连接正式建立成功时触发。 - 一般用来关闭加载按钮的 loading 状态,或者给用户提示「开始生成题目」。