Spring AI 2.0:Flux
方法作用Spring AI 常见用途map一对一转换内容脱敏、包装 SSE、格式转换filter过滤元素忽略空片段或无效事件doOnNext观察元素日志、统计输出片段doOnError观察异常记录模型或网络异常doFinally观察最终结束统计耗时、识别完成/失败/取消onErrorResume切换备用流返回友好错误、服务降级timeout无信号超时防止模型长期无响应retryWhen按策略重试对连接阶段瞬时错误有限重试take截取元素调试、预览、主动提前结束concatWith追加流添加结束事件或尾部提示collectList收集为列表完整输出后批量处理reduce聚合为一个值拼接完整回答后存库scan持续输出累计值维护当前完整回答快照flatMap并发异步转换并发处理多个独立请求concatMap顺序异步转换按顺序处理问题或工具结果bufferTimeout按数量或时间分批合并片段、降低处理频率把Mono类比为“未来可能得到一个结果”把Flux类比为“未来可能陆续得到多个结果”。类型元素数量常见用途MonoT0 或 1 个查询单条数据、一次性 AI 回答、保存结果FluxT0 到 N 个AI 流式回答、消息流、列表查询、SSE 推送just将少量元素封装成流FluxStringfluxFlux.just(hello,world,flux);fromFlux.from(...)接收 Reactive Streams 的Publisher或者 Flowable 对象。PublisherStringpublisherobtainPublisher();FluxStringfluxFlux.from(publisher);FlowableApplicationResultnewApplication().streamCall(param);map同步的一对一转换FluxStringupperCaseFlux.just(java,spring).map(String::toUpperCase);Spring AI 中可以用map处理每个流式片段FluxStringcontentchatClient.prompt().user(message).stream().content().map(chunk-chunk.replace(敏感词,***));map中不能返回null。如果某个元素不需要保留应该使用filter或handle。Flux.deferFlux.defer(() - publisher)把 Publisher 的创建推迟到订阅时刻每次新订阅都会重新执行一遍生成逻辑。用于延迟 IO、延迟副作用、保证每次订阅拿到全新数据流。FluxStringfluxFlux.defer(()-{// 只有 subscribe() 的时候才会走到这里执行 getData()returnFlux.just(getData());});// 此时还没调用 getData()// flux.subscribe(); // 订阅触发才执行 lambda、调用getData典型使用场景延迟执行耗时 / IO 逻辑把数据库查询、接口调用放到 defer 的 lambda 内部有订阅才触发 IO避免提前浪费资源。每次订阅都重新生成新的数据流封装有副作用的代码比如获取连接、打开文件放到 defer 内部订阅才打开取消订阅就释放避免提前执行副作用。doOnXxx事件doFirst /doOnError/doOnComplete /doOnCancel 这一组属于回调钩子side‑effect 副作用方法不修改数据流元素只做事件监听。doOnXxx 只是监听不会捕获异常异常依旧向下传播。doFirst订阅发生之前执行流还没有开始发射数据doOnComplete流正常走完全部元素发射完毕正常完成doOnError流发生异常抛出错误doOnCancel订阅被手动取消主动终止流没有走完也没有报错。场景1正常走完无异常不取消FluxIntegerflux1Flux.just(1,2,3).doFirst(()-System.out.println(✅ doFirst准备订阅还没发数据)).doOnComplete(()-System.out.println(✅ doOnComplete流正常结束全部数据发射完成)).doOnError(e-System.out.println(❌ doOnError发生异常e.getMessage())).doOnCancel(()-System.out.println(⚠️ doOnCancel流被手动取消));flux1.subscribe(System.out::println);✅ doFirst准备订阅还没发数据123✅ doOnComplete流正常结束全部数据发射完成场景2中间抛出异常FluxIntegerflux2Flux.just(1,2,3).map(i-{if(i2){thrownewRuntimeException(模拟业务报错);}returni;}).doFirst(()-System.out.println(✅ doFirst准备订阅)).doOnComplete(()-System.out.println(✅ doOnComplete正常完成报错不会进这里)).doOnError(e-System.out.println(❌ doOnError捕获异常e.getMessage())).doOnCancel(()-System.out.println(⚠️ doOnCancel流被手动取消报错不会进这里));flux2.subscribe(System.out::println,error-System.out.println(【subscribe收到异常】error.getMessage()));✅ doFirst准备订阅1❌ doOnError捕获异常模拟业务报错 【subscribe收到异常】模拟业务报错场景3手动cancel取消流没有报错、没有走完completeFluxIntegerflux3Flux.just(1,2,3,4,5).doFirst(()-System.out.println(✅ doFirst准备订阅)).doOnComplete(()-System.out.println(✅ doOnComplete取消不会进)).doOnError(e-System.out.println(❌ doOnError取消不会进)).doOnCancel(()-System.out.println(⚠️ doOnCancel流被手动取消));// subscribe 返回 Disposable可以手动取消订阅vardisposableflux3.subscribe(System.out::println);// 模拟业务收到部分数据后手动取消流disposable.dispose();✅ doFirst准备订阅1⚠️ doOnCancel流被手动取消场景4区分 doFirst 位置doFirst在链不同位置的效果// doFirst是越靠下游越先执行从订阅点向上执行// 链式语法调用是有顺序的顺序不同效果不同。Flux.just(10,20).doFirst(()-System.out.println(A doFirst)).map(i-i*2).doFirst(()-System.out.println(B doFirst)).subscribe(System.out::println);BdoFirstAdoFirst2040Reactor Flux doFirst / doOnComplete / doOnError / doOnCancel方法触发时机说明doFirst(Runnable)订阅subscribe发生时数据流还未发出任何元素可以做日志、初始化注意位置越靠近下游越优先执行不生产数据仅副作用doOnComplete(Runnable)流正常完整结束所有元素发射完毕没有异常、没有取消异常、cancel时不会执行适合正常结束日志、统计doOnError(ConsumerThrowable)流发生异常向上传播错误信号只在异常场景触发不会捕获异常异常继续往下传递打印异常日志doOnCancel(Runnable)订阅被手动dispose()取消流既没有正常complete也没有报错SSE聊天场景用户点停止输出就会触发 doOnCancel用来做资源清理信号互斥规则doOnComplete和doOnError互斥正常完成就不会进error抛异常就不会进complete。doOnCancel和 complete / error 互斥手动取消流既不会complete也不会error。doFirst只要发生订阅就执行无论后续是complete/error/cancel。实际业务场景SSE流式聊天doFirst记录开始流式日志doOnComplete大模型正常输出完毕记录成功结束doOnError大模型调用报错打印异常doOnCancel用户点击【停止输出】按钮dispose取消Flux触发doOnCancel做资源清理。业务场景类比对应 Spring AI SSE用户请求进来开始订阅流 → doFirstAI 完整输出全部 token服务端正常结束 → doOnComplete调用大模型接口抛异常 → doOnError用户前端点【停止回答】AbortController后端 Flux 被 cancel → doOnCancel。doOnCancel 非常适合做 SSE 聊天的资源释放这个钩子只有主动取消才进报错、正常完成不会进入。takeWhiletakeWhile(Predicate)满足条件就继续接收元素一旦条件不满足直接终止流发送onComplete信号。只判断每一个下发出来的元素条件返回 true 就下发一旦返回 false当前这个不满足的元素直接丢弃流直接结束触发 doOnComplete不会触发 doOnError、不会触发 doOnCancel// 当4到来44false条件不成立流终止4、5、6 全部不再下发。// 导致终止的那一条元素不会向下游传递Flux.just(1,2,3,4,5,6).takeWhile(num-num4).doFirst(()-System.out.println(doFirst)).doOnComplete(()-System.out.println(doOnComplete 流结束)).doOnCancel(()-System.out.println(doOnCancel)).subscribe(System.out::println);doFirst123doOnComplete 流结束在 Spring AI 流式聊天场景可以检测输出内容包含某个敏感词takeWhile直接终止大模型输出。// 一旦chunk包含敏感词直接终止流// ⚠️注意此时是正常 complete走doOnComplete不是 cancel不走 doOnCancel 钩子。flux.takeWhile(chunk-!chunk.contains(敏感词))concatWithconcatWith(Publisher? extends T other)先执行当前流当前流正常 onComplete 完成之后再去订阅第二个流把第二个流的数据接续发射出来。串行执行不是并行必须等第一个流全部结束才跑第二个。第一个流异常第二个不会执行takeWhile 提前 completeconcatWith 会执行后面流// 发射 1,2,3 → flux1 正常 complete// 才去订阅 flux2发射 10,20,30FluxIntegerflux1Flux.just(1,2,3).doOnComplete(()-System.out.println(【flux1完成】));FluxIntegerflux2Flux.just(10,20,30).doFirst(()-System.out.println(开始订阅flux2));flux1.concatWith(flux2).subscribe(System.out::println);123【flux1完成】开始订阅flux2102030在 Spring AI 流式聊天中的业务举例场景先输出一段前置提示文本再输出大模型流式回答FluxStringprefixFlux.just(【知识库检索完成开始回答】\n);FluxStringaiStreamchatClient.prompt().user(question).stream().content();//先输出prefix等prefix完成之后再输出AI流式tokenreturnprefix.concatWith(aiStream);doOnXxx Spring AI 示例业务场景SSE 流式聊天接口记录各个生命周期日志doOnCancel对应前端点击停止输出AbortController取消请求。关键点chatClient.stream().content() 返回 FluxdoOnCancel前端断开 / 点停止按钮后端 Flux 会触发 cancel适合清理资源、计数doOnError大模型调用异常、限流、鉴权失败触发doOnCompleteAI 完整把回答输出完毕正常结束doFirst订阅发生开始推送 token 之前执行Slf4jRestControllerRequestMapping(/rag)publicclassRagStreamController{privatefinalChatClientchatClient;publicRagStreamController(ChatClientchatClient){this.chatClientchatClient;}DatapublicstaticclassChatQueryDTO{privateStringsid;privateStringquestion;}/** * SSE流式接口produces text/event‑stream */PostMapping(value/streamChat,producesMediaType.TEXT_EVENT_STREAM_VALUE)publicFluxStringstreamChat(RequestBodyChatQueryDTOdto){Stringsiddto.getSid();Stringquestiondto.getQuestion();FluxStringcontentFluxchatClient.prompt().user(question).advisors(a-a.param(ChatMemory.CONVERSATION_ID,sid)).stream().content();// 挂上Reactor生命周期钩子returncontentFlux.doFirst(()-{// 订阅发生即将开始返回tokenlog.info([doFirst] 会话sid{},开始流式问答用户问题:{},sid,question);}).doOnComplete(()-{// ✅ AI完整输出完毕正常结束log.info([doOnComplete] 会话sid{},流式回答全部输出完成,sid);}).doOnError(throwable-{// ❌ 发生异常大模型报错、限流、网络异常log.error([doOnError] 会话sid{},流式问答异常error{},sid,throwable.getMessage(),throwable);}).doOnCancel(()-{// ⚠️ 重点前端主动断开连接 / 用户点击【停止输出】触发// 既没有正常complete也没有异常属于人为取消流log.warn([doOnCancel] 会话sid{},流式回答被用户主动取消,sid);});}}各个钩子触发时机结合前端microsoft/fetch‑event‑sourcedoOnComplete、doOnError、doOnCancel 三者互斥只会进入其中一个doFirst只要订阅就会执行不管后面结局是 complete /error/canceldoFirst前端发送请求后端开始订阅Flux还没有返回任何 token。用途打印请求日志、埋点计数。doOnCompleteAI 把全部 token 全部推送给前端流正常结束。用途统计成功会话、记录完成时间。doOnError大模型 API 调用失败、密钥错误、限流超时代码内部抛出异常。⚠️注意报错不会执行 doOnComplete /doOnCancel。doOnCancel【SSE 最关键】Flux 收到取消信号进入doOnCancel。用途释放临时资源、记录用户中途终止问答埋点。两种场景会触发用户点击前端停止按钮调用 abortController.abort()用户直接关闭浏览器标签页网络连接断开// microsoft/fetch-event-sourceletabortController:AbortController|nullnull;asyncfunctionchat(){abortControllernewAbortController();awaitfetchEventSource(/rag/streamChat,{method:POST,signal:abortController.signal,// 这个signal abort会传递到后端Flux触发doOnCancel// ...省略headers body})}// 用户点击停止按钮functionstopAnswer(){if(abortController){abortController.abort();// ← 后端进入 doOnCancel}}注意doOnCancel 不会捕获异常只是监听信号如果后端直接返回 Flux.error()只会进doOnError不会进doOnCancelSpring SSE 场景只有客户端主动断开才会触发 doOnCancel服务端主动结束流触发doOnComplete。

相关新闻

OpenRouter Ori DeepSeek Harness:本地部署大模型实战指南

OpenRouter Ori DeepSeek Harness:本地部署大模型实战指南

在探索大模型应用落地的过程中,开发者们常常面临一个核心矛盾:如何在享受云端强大模型能力的同时,又能保证数据安全、降低延迟,并实现成本可控?OpenRouter 最新推出的 Ori DeepSeek Harness 正是为解决这一痛点而生。…

2026/8/18 20:50:26 阅读更多 →
AI搜索时代SEO实战:从GEO地理意图到AEO答案引擎的优化策略

AI搜索时代SEO实战:从GEO地理意图到AEO答案引擎的优化策略

在实际的搜索引擎优化工作中,我们正面临一个根本性的转变。传统的SEO策略,如关键词堆砌、外链建设,其效果正被以Google SGE为代表的AI原生搜索体验所稀释。当用户的问题可以直接在搜索结果页通过AI生成摘要得到解答时,仅仅“出现在…

2026/8/18 20:50:26 阅读更多 →
保时捷电动Macan技术解析:PPE平台、800V高压与性能重塑

保时捷电动Macan技术解析:PPE平台、800V高压与性能重塑

1. 从“燃油图腾”到“电动先锋”:Macan转型的行业背景与深层逻辑 保时捷要推出纯电动版Macan,这个消息在汽车圈里炸开了锅。很多人第一反应可能是:“保时捷也扛不住了?” 或者 “连Macan都要电动化,燃油性能车是不是没…

2026/8/18 20:50:26 阅读更多 →

最新新闻

基于Qwen与RAG+LoRA的法律大模型实战:罪名识别与刑期预测

基于Qwen与RAG+LoRA的法律大模型实战:罪名识别与刑期预测

这次我们来看一个面向法律领域的实战项目:基于通义千问(Qwen)大模型,结合RAG(检索增强生成)和LoRA(低秩适应)微调技术,构建一个能够进行罪名识别、刑期预测和司法解释生成…

2026/8/18 22:02:18 阅读更多 →
OpenClaw端口通信失效分析与优化实践

OpenClaw端口通信失效分析与优化实践

1. OpenClaw端口通信失效问题全景分析 OpenClaw作为新兴的开源通信中间件,其端口通信失效问题在实际部署中频繁出现。根据我们团队在三个大型企业项目中的实测数据,约67%的OpenClaw部署初期都会遭遇不同形式的端口通信问题。这些故障往往表现为以下几种典…

2026/8/18 22:02:18 阅读更多 →
大模型上下文长度技术解析:从原理到实践,2026年发展趋势与部署指南

大模型上下文长度技术解析:从原理到实践,2026年发展趋势与部署指南

这次我们来看一个关于大模型上下文长度的技术话题。如果你正在选型大模型、开发AI应用,或者被“上下文不足”的问题困扰,这篇文章可以直接收藏。上下文长度(Context Window)直接决定了大模型能“记住”多长的对话、处理多长的文档…

2026/8/18 22:02:18 阅读更多 →
基于Qwen与RAG+LoRA构建法律大模型:从知识库到微调部署实战

基于Qwen与RAG+LoRA构建法律大模型:从知识库到微调部署实战

最近在参与一个法律科技相关的项目,核心需求是让大模型能够理解复杂的案情描述,并给出专业的法律分析,比如判断可能涉及的罪名、预测可能的刑期范围,甚至生成相关的司法解释摘要。这听起来像是法律专家的活儿,但团队希…

2026/8/18 22:01:17 阅读更多 →
网页视频保存不了?VideoDownloadHelper 免费开源下载插件快速上手完整指南

网页视频保存不了?VideoDownloadHelper 免费开源下载插件快速上手完整指南

网页视频保存不了?VideoDownloadHelper 免费开源下载插件快速上手完整指南 【免费下载链接】VideoDownloadHelper Chrome Extension to Help Download Video for Some Video Sites. 项目地址: https://gitcode.com/gh_mirrors/vi/VideoDownloadHelper 深夜想…

2026/8/18 22:01:17 阅读更多 →
基于RAG与LoRA微调通义千问,构建专业法律AI助手实战

基于RAG与LoRA微调通义千问,构建专业法律AI助手实战

如果你正在开发一个法律咨询或案件分析系统,是否遇到过这样的困境:通用大模型对法律条文的理解似是而非,回答“可能构成XX罪”但不敢给出明确判断;或者,面对复杂的案情描述,模型无法准确关联到具体的司法解…

2026/8/18 22:01:17 阅读更多 →

日新闻

告别逐帧截图:用 extract-video-ppt 快速提取视频中的 PPT 并一键导出 PDF

告别逐帧截图:用 extract-video-ppt 快速提取视频中的 PPT 并一键导出 PDF

告别逐帧截图:用 extract-video-ppt 快速提取视频中的 PPT 并一键导出 PDF 【免费下载链接】extract-video-ppt extract the ppt in the video 项目地址: https://gitcode.com/gh_mirrors/ex/extract-video-ppt 如果你还停留在"看网课 不停暂停 截图 …

2026/8/18 0:00:57 阅读更多 →
思源宋体TTF一站式上手:7个字重免费商用,从下载到上线的完整走查

思源宋体TTF一站式上手:7个字重免费商用,从下载到上线的完整走查

思源宋体TTF一站式上手:7个字重免费商用,从下载到上线的完整走查 【免费下载链接】source-han-serif-ttf Source Han Serif TTF 项目地址: https://gitcode.com/gh_mirrors/so/source-han-serif-ttf 你是不是也经历过这种时刻:设计稿里…

2026/8/18 0:00:58 阅读更多 →
华硕笔记本控制权回收指南:GHelper 如何用一个 10MB 文件替代 Armoury Crate

华硕笔记本控制权回收指南:GHelper 如何用一个 10MB 文件替代 Armoury Crate

华硕笔记本控制权回收指南:GHelper 如何用一个 10MB 文件替代 Armoury Crate 【免费下载链接】g-helper Lightweight Armoury Crate alternative for Asus laptops with nearly the same functionality. Works with ROG Zephyrus, Flow, TUF, Strix, Scar, ProArt, …

2026/8/18 0:00:59 阅读更多 →

周新闻

基于阿里云与通义千问(Qwen)构建AI应用:从模型调用到生产部署的完整实践指南

基于阿里云与通义千问(Qwen)构建AI应用:从模型调用到生产部署的完整实践指南

如果你是一名开发者,最近可能已经感受到了AI大模型正在从“玩具”变成“生产力工具”的强烈信号。从代码补全到智能Agent,从本地部署到云端API,我们正处在一个技术栈快速重构的节点。然而,面对层出不穷的模型、框架和工具&#xf…

2026/8/18 9:15:35 阅读更多 →
工业通信系统底层逻辑:04 反射——高频能量撞墙之后会发生什么?

工业通信系统底层逻辑:04 反射——高频能量撞墙之后会发生什么?

第四篇:反射——高频能量撞墙之后会发生什么? —— 你以为信号已经过去了,其实它正在回来打你 老Q的现场笔记 第五季,我们正式进入工业神经系统层。这里不再是单个设备的战斗,而是整个工厂“经脉”层面的秩序之战。从这一篇开始,你将第一次看清:看似简单的信号传播,背…

2026/8/18 9:06:28 阅读更多 →
【文章复现】非线性值迭代自适应动态规划(ADP):离散时间非线性系统的策略迭代自适应动态规划算法研究附Matlab代码

【文章复现】非线性值迭代自适应动态规划(ADP):离散时间非线性系统的策略迭代自适应动态规划算法研究附Matlab代码

✅作者简介:热爱科研的Matlab仿真开发者,擅长毕业设计辅导、数学建模、数据处理、建模仿真、程序设计、完整代码获取、论文复现及科研仿真。🍎 往期回顾关注个人主页:Matlab科研工作室👇 关注我领取海量matlab电子书和…

2026/8/18 9:04:56 阅读更多 →

月新闻

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南 【免费下载链接】BaiduNetdiskPlugin-macOS For macOS.百度网盘 破解SVIP、下载速度限制~ 项目地址: https://gitcode.com/gh_mirrors/ba/BaiduNetdiskPlugin-macOS 还在为百度网盘macOS版的龟速下…

2026/8/17 18:54:37 阅读更多 →
终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换

终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换

终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换 【免费下载链接】ncmdump 项目地址: https://gitcode.com/gh_mirrors/ncmd/ncmdump 还在为网易云音乐下载的NCM格式文件无法在其他播放器播放而烦恼吗?ncmdump解密工具帮你轻松解决这个困…

2026/8/17 18:55:16 阅读更多 →
HarmonyOS 应用开发《掌上英语》第81篇: 智能体卡片:为英语学习 App 打造桌面级学习助手

HarmonyOS 应用开发《掌上英语》第81篇: 智能体卡片:为英语学习 App 打造桌面级学习助手

AgentCard 智能体卡片:为英语学习 App 打造桌面级学习助手适用平台:HarmonyOS 7.0 (API 26 Beta)一、引言 HarmonyOS 7.0(API 26 Beta)新增了 AgentCard 智能体卡片能力,这是继 HMAF(鸿蒙智能体框架&#x…

2026/8/17 18:55:55 阅读更多 →