Flux中常用的方法
在Spring AI里,Flux 是处理AI流式响应的核心。像 doOnNext、window 等方法并非Spring AI独有,它们都来自 Project Reactor——这是Spring生态中实现响应式编程(Reactive Programming)的基石。
理解这些操作符,是掌握Spring AI流式编程的关键。下面我把常用的操作符分分类,建立一个清晰的使用框架。
🧭 核心方法速查表
我把Flux中常用的方法整理成了下面这个表格,方便你快速查阅和对比:
1.副作用与监听
| 方法 | 作用 | Spring AI 常见用途 |
|---|---|---|
| doOnNext | 当数据流发出每个元素时,触发一个回调动作 | 记录日志、统计输出片段数量、将每个流式片段推送给客户端 |
| doOnComplete | 当数据流正常完成(无错误)时,触发一个回调动作 | 通知客户端流式响应已结束,或执行收尾清理工作 |
| doOnError | 当数据流发生错误时,触发一个回调动作 | 记录错误日志,或向客户端返回一个友好的错误提示 |
| doFinally | 无论流是正常完成、发生错误还是被取消,都会在最后执行一次回调 | 统计一次完整调用的总耗时,或执行无论如何都需要做的清理工作 |
| doOnCancel | 当数据流被手动取消时,触发一个回调动作 | 记录取消事件,或执行取消相关的资源清理 |
| doFirst | 在订阅发生之前执行一个回调动作 | 在流开始发射数据前执行一些初始化操作 |
| doOnSubscribe | 当有订阅者订阅时触发 | 可以在这里记录订阅信息 |
2.转换与过滤
| 方法 | 作用 | Spring AI 常见用途 |
|---|---|---|
| map | 对数据流中的每一个元素进行同步的“一对一”转换 | 对AI生成的每个文本片段进行脱敏、格式转换或内容包装 |
| filter | 根据条件过滤掉不符合要求的元素 | 忽略AI返回的空片段或无效事件 |
| flatMap | 将每个元素转换为一个Publisher(发布者),然后将这些新的流“压平”合并 | 并发处理多个独立的请求 |
| concatMap | 和flatMap类似,但会保持顺序,一个接一个地处理 | 按顺序处理问题或工具调用结果 |
3.错误处理
| 方法 | 作用 | Spring AI 常见用途 |
|---|---|---|
| onErrorResume | 当发生错误时,提供一个备用的数据流来“顶替 | 当AI服务出错时,返回一个包含友好错误信息的备用流,实现服务降级 |
| retryWhen | 当发生错误时,根据指定的策略(如重试次数、间隔)进行重试 | 对网络连接等瞬时错误进行有限次数的重试 |
| timeout | 为数据流设置一个超时时间,超时则触发错误 | 防止AI模型长时间无响应 |
4.流控制与截取
| 方法 | 作用 | Spring AI 常见用途 |
|---|---|---|
| take | 从数据流中截取前N个元素 | 调试、预览,或在满足条件时主动提前结束流 |
| skip | 跳过前N个元素 | - |
| limitRate | 控制背压(Backpressure)的速率 | - |
| sample | 在指定时间间隔内,采样最新发出的元素 | - |
5.分组与聚合
| 方法 | 作用 | Spring AI 常见用途 |
|---|---|---|
| window | 将数据流按数量、时间等条件,拆分成多个子数据流(Flux) | 将连续的流式片段按句子或段落分组 |
| buffer | 和window类似,但拆分后是将元素收集到一个集合(如List) 中 | 将片段合并成一块再处理,降低处理频率 |
| collectList | 收集所有元素到一个List中,在流完成时发出 | 等待完整回答后,再对整个结果进行批量处理 |
| reduce / scan | 对流中的元素进行累积计算 | 拼接出完整的AI回答后存入数据库 |
6.组合与合并
| 方法 | 作用 | Spring AI 常见用途 |
|---|---|---|
| concatWith | 将另一个数据流拼接在当前流之后 | 在AI回答结束后,追加一段结束语或提示信息 |
| mergeWith | 将多个数据流合并,元素会交叉发出 | - |
| zipWith | 将两个流的元素按“一对一”的方式配对组合 | - |
7.条件与分支
| 方法 | 作用 | Spring AI 常见用途 |
|---|---|---|
| zipWith | 将两个流的元素按“一对一”的方式配对组合 | - |
| defaultIfEmpty | 如果流是空的,则发出一个默认值 | - |
| switchIfEmpty | 如果流是空的,则切换到另一个备用的流 | - |
🧩 重点方法详解
window vs buffer:分组处理的两种思路
当需要将连续的流式数据分割成块处理时,window 和 buffer 非常有用。它们的主要区别在于输出的类型:
- window:将原始 Flux<T> 拆分成多个新的 Flux<T> 子流。每个子流都像一个独立的“小水管”,可以单独进行订阅和操作。用途:适合需要对每个分组进行复杂异步处理的场景。例如,将AI生成的文本按句子分组后,对每个句子进行独立的语法检查或翻译。
- buffer:将原始 Flux<T> 拆分成多个集合(如 List<T>) 。用途:适合需要批量处理数据的场景。例如,将AI生成的多个文本片段缓存到一个列表中,当累积到一定数量或超时时,再一次性写入数据库或发送给客户端。
doOnNext、doOnComplete、doOnError:无处不在的“观察者”
这组以 doOn 开头的方法,是响应式编程中的副作用(Side-effect) 操作符。它们不会修改数据本身,而是让你能在数据流的特定事件发生时,插入一些额外的操作。这就像是在数据流的管道上安装了几个“观察窗”。
- doOnNext:每当有一个新元素流过时,“观察窗”就会被触发。
- doOnComplete:当整个数据流正常结束时(好比水流完了),“观察窗”会被触发。
- doOnError:当数据流发生错误时(好比管道破裂了),“观察窗”会被触发。
注意:doOnXxx 方法只负责观察和触发动作,不会捕获或处理异常。如果 doOnError 的回调中抛出了新的异常,这个新异常会继续向下游传播。真正的错误恢复应使用 onErrorResume 等方法。
💡 在 Spring AI 中的典型应用
理解这些方法后,看看它们在Spring AI中是如何协同工作的。一个典型的流式AI调用流程如下:
- 发起流式请求:通过 ChatClient 的 stream() 方法发起请求,获得一个 Flux<ChatResponse> 或 Flux<String> 的数据流。
- 处理流式数据:使用 doOnNext 处理每个到达的文本片段。你可以在这里将数据实时推送给前端(如通过SSE),或进行日志记录。使用 map 对每个片段进行转换,例如脱敏或格式调整。使用 filter 过滤掉无效或不需要的片段。使用 window 或 buffer 对片段进行分组,以实现更复杂的批处理逻辑。
- 处理结束与错误:使用 doOnComplete 在流正常结束时,通知客户端或执行收尾工作。使用 doOnError 记录错误日志。使用 onErrorResume 在发生错误时,返回一个包含友好提示的备用流,实现优雅降级。
- 启动执行:最后,必须通过 subscribe() 方法订阅这个数据流,整个流程才会真正开始执行。