ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

后端+智谱ai 结合Rxjava,SSE技术向前端推送AI结果

后端+智谱ai 结合Rxjava,SSE技术向前端推送AI结果 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.返回 FlowableModelData对象/** * 通用流式请求 * * param messages * param stream * param temperature * return */ public FlowableModelData doStreamRequest(ListChatMessage 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 FlowableModelData doStreamRequest(String systemMessage, String userMessage, Float temperature) { // 创建聊天完成请求 ListChatMessage 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是自定义的系统信息以上都是用于传递给AiGetMapping(/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流式返回 FlowableModelData 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 - { ListCharacter 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.RxjavaSSE技术向前端推送AI结果核心部分// 建立SSE连接对象, 0 表示无超时时间 SseEmitter sseEmitter new SseEmitter(0L); // AI 生成 SSE流式返回 FlowableModelData 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 - { ListCharacter 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 生成一点后端就推一点前端可以实时看到题目逐道出现的效果类似打字机。技术 1SseEmitter—— Spring 提供的 SSE长连接工具负责保持和前端的连接持续推送数据技术 2FlowableRxJava—— 负责处理 AI 返回的流式数据做转换、拆分、拼接核心业务把 AI 流式吐出的零散文本拼成一道道完整的题目 JSON再逐道推给前端1.新建SseEmitter相当于和前端建立一条「一直开着的传输通道」SseEmitter sseEmitter new SseEmitter(0L);2.从aiManager.doStreamRequest开始进入 RxJava 流式处理得到FlowableModelData对象FlowableModelData 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 - { ListCharacter 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流出现异常时执行这里只打了日志doOnCompleteAI 所有内容生成完、流正常结束时关闭 SSE 连接subscribe()真正启动这条流。RxJava 是「懒执行」的不写这一行前面所有代码都只是定义不会真正运行5.总结Rxjava和SSE是怎么结合在一起使用的RxJava 负责后端内部的流式数据加工SSE 负责把加工好的数据推给前端两者靠sseEmitter.send()这一行代码衔接一个产、一个发天然都是「流式」思想配合起来非常顺滑。RxJava后端内部的异步数据流处理工具管「数据怎么拆、怎么拼、怎么过滤、怎么切线程」SSE后端 ↔ 前端的通信推送技术管「怎么把数据持续推给浏览器」一、先明确各自的分工1. RxJavaFlowable内部数据加工厂它的工作范围完全在后端服务内部从 AI SDK 接收到原始的流式响应一小段一小段的文字碎片用map提取内容、filter过滤空值、flatMap拆成字符用括号计数器拼出完整的单道题目全程在 IO 线程运行不阻塞主线程它只关心「怎么把碎数据加工成可用的业务数据」不关心数据最终是存数据库、返回接口还是推给前端。2. SSESseEmitter对外传输通道它的工作是和前端浏览器打交道建立一条不关闭的 HTTP 长连接后端随时可以调用send()往通道里塞数据前端实时收到调用complete()主动关闭连接它只关心「怎么把数据推给前端」不关心数据是 AI 生成的、数据库查的还是手动拼的。二、核心结合点就这一行代码两者唯一的交集就是doOnNext里的这一行sseEmitter.send(JSONUtil.toJsonStr(stringBuilder.toString()));这就是「加工厂」和「运输通道」的对接窗口RxJava 每加工好一道完整的题目就调用一次send()SSE 接到数据立刻通过长连接推给前端推完继续等 RxJava 加工下一道俗类比 RxJava 包子铺后厨负责揉面、包馅、蒸包子蒸好一个放一个到出餐口 SSE 外卖传送滑道出餐口放一个滑道就送一个到顾客手里 结合点 出餐口的那个窗口后厨放进去滑道传出去通三、完整数据流走一遍全程对应我们把你代码的完整流程按「谁在干活」标出来一眼就能看清前端发起请求 → SSE 建立连接后端创建SseEmitter和前端建立长连接这一步是 SSE 的活。调用 AI 接口 → 拿到 RxJava 流aiManager.doStreamRequest()返回FlowableModelDataAI 生成的文字碎片会源源不断地流进这条 RxJava 管道。链式加工数据 → 全是 RxJava 的活observeOn 切线程 → map 提取内容 → filter 过滤空值 → flatMap 拆字符 → doOnNext 拼题目这一长串全部是 RxJava 在内部处理数据和 SSE 没有任何关系。4.关键衔接拼完一道推一道每当括号计数器归零、一道题拼接完成就执行sseEmitter.send(...)把 RxJava 加工好的成品塞进 SSE 通道推给前端。5.结束 / 异常同步关闭连接RxJava 流正常结束 →doOnComplete里调用sseEmitter.complete()关闭 SSE 连接RxJava 流出错 → 调用sseEmitter.completeWithError()通知前端异常结束部分代码理解为什么要做这一步操作切线程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 状态或者给用户提示「开始生成题目」。
返回列表