发布时间:2026/8/18 20:50:20
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。

相关新闻

2026/8/18 20:50:20

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

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

2026/8/18 20:45:19

从论文复现看AI智能体工程化能力:超越基准测试的核心评估

上周,一个名为“Faraday 27B”的智能体在论文复现任务上,其表现被一些讨论认为超越了Claude Opus 4.8和GPT-5.5。这个消息在技术社区里激起了一些水花,但很快又淹没在“哪个模型更强”的日常争论中。作为一个长期观察AI应用落地的开发者&…

2026/8/18 20:45:19

Python编程核心术语全解析:从基础语法到高级特性实战指南

1. 项目概述:一份Python程序员的专属词汇表 作为写了十多年Python代码的老码农,我深知一个痛点:看官方文档、读开源项目源码、在Stack Overflow上找答案时,那些高频出现的英文单词和术语,常常成为新手甚至有一定经验开…

2026/8/18 21:55:26

2026年小程序开发公司排行怎么参考?费用、周期和售后能力对比

2026年小程序开发公司排行怎么参考?费用、周期和售后能力对比小程序开发公司排行可以作为初筛线索,但不能直接等于适合。企业真正要比较的是搭建方式、上线周期、费用边界、后台维护、支付或预约能力、页面交付和售后响应。普通展示、预约留资、小程序商…

2026/8/18 21:55:26

从零搭建游戏串流服务器:Sunshine 完整部署实战指南

从零搭建游戏串流服务器:Sunshine 完整部署实战指南 【免费下载链接】Sunshine Self-hosted game stream host for Moonlight. 项目地址: https://gitcode.com/GitHub_Trending/su/Sunshine 想躺在沙发上用电视玩 PC 游戏,或出差时用手机远程开黑…

2026/8/18 21:55:26

Linux权限管理深度解析:从chmod 777到安全配置实战

1. 从一次线上故障说起:权限失控的代价 那天凌晨,我被一阵急促的电话铃声吵醒。监控系统疯狂报警,显示线上核心应用服务器磁盘空间在五分钟内被写满,服务全部瘫痪。顶着困意连上服务器,用 df -h 一看, /…

2026/8/18 21:55:26

Git提交拆分实战:交互式变基与暂存实现原子提交

在实际 Git 协作开发中,我们经常会遇到一个尴尬的场景:一个已经暂存(staged)甚至已经提交(committed)的改动,包含了多个逻辑上独立的修改。比如,你在修复一个 Bug 的同时&#xff0c…

2026/8/18 21:55:26

LLM Agent部署实战:揭秘约束规避性虚构与假死行为及应对策略

1. 引言:当你的AI代理开始“装死” 最近在折腾一个基于GPT-4o的智能客服微服务项目,遇到了一个极其诡异的现象。我的LLM Agent(大语言模型代理)在部署上线后,面对某些特定的、带有约束条件的用户查询时,会突…

2026/8/18 21:50:26

解码温度如何影响多智能体LLM命名游戏的共识形成

1. 项目概述:当大语言模型开始“玩”命名游戏 最近在折腾一个挺有意思的实验项目,核心是观察一群大语言模型(LLM)智能体,如何通过一种叫做“命名游戏”的简单交互,最终对一个未知事物达成统一的命名共识。听…

2026/8/17 10:49:52

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

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

2026/8/18 6:58:27

工业传感器与变送器详解:序章 从物理世界到工业数据

序章 从物理世界到工业数据 ——重新认识工业传感器与变送器 工业自动化系统正变得日益复杂。今天的工业现场早已不是简单的控制回路,而是由多层技术共同构成的立体体系:PLC、DCS、SCADA、MES、工业互联网、边缘计算与人工智能。控制系统可以执行复杂算法,工业网络可以实现…

2026/8/18 0:02:05

Qwen3.8-27B本地部署实战:17GB内存运行270亿参数大模型

1. 这篇文章真正要解决的问题 你是否曾对动辄需要上百GB显存才能运行的百亿参数大模型望而却步?是否觉得在个人电脑上部署一个功能强大的语言模型是天方夜谭?最近,通义千问团队发布的 Qwen3.8-27B 模型,宣称仅需 17GB 内存即可在本…

2026/8/18 0:02:05

ME3169 36V,8A,180KHz 恒压Buck DC-DC 转换器

概述ME3169 是一款180KHz,PWM 模式恒压Buck DC-DC 转换器,8V 到36V 宽工作电压范围,低纹波,内置低导通电阻功率MOS。ME3169 内置环路补偿电路,可以减少外围元器件数量。内部设计有恒压环路,可以通过外部电阻…

2026/8/18 18:23:10

实测才敢推 AI论文网站 2026最新测评与推荐

2026年真正好用的AI论文网站,核心看生成的论文质量、低AI味、格式正确、学术适配四大指标。综合实测,千笔AI、ThouPen、豆包、DeepSeek、Grammarly 是当前最值得推荐的梯队,覆盖从免费到付费、从中文到英文、从文科到理工的全场景需求。一、综…

2026/8/17 17:27:06

2026必备!AI论文网站测评:最新推荐与深度对比

2026年真正好用的AI论文网站,核心看生成的论文质量、低AI味、格式正确、学术适配四大指标。综合实测,千笔AI、ThouPen、豆包、DeepSeek、Grammarly 是当前最值得推荐的梯队,覆盖从免费到付费、从中文到英文、从文科到理工的全场景需求。 一、…

2026/8/18 7:12:40

摆脱论文困扰!盘点2026年全网爆红的的AI论文写作工具

一天写完毕业论文在2026年已不再是天方夜谭。2026年最炸裂、实测能大幅提速的AI论文写作工具,覆盖选题构思、文献整理、内容生成、格式排版等核心场景,真正帮你高效搞定论文难题。 一、全流程王者:一站式搞定论文全链路(一天定稿首…