Spring Boot 中使用 WebClient + Flux + SseEmitter 实现流式响应
目录
- 前言
- 一、什么是流式响应?
- 二、什么是 SSE?
- 三、SseEmitter 是什么?
- 四、为什么会出现 WebClient + Flux + SseEmitter?
- 五、什么是 Spring WebFlux?
- 六、WebFlux 和 Reactor 的关系
- 七、什么是 Flux?
- 八、Flux 的三个核心阶段
- 九、为什么不调用 subscribe(),Flux 可能不会执行?
- 十、Flux 的执行信号
- 十一、subscribe() 是什么?
- 十二、map() 是什么?
- 十三、flatMap() 是什么?
- 十四、map 和 flatMap 的区别
- 十五、doOnNext() 是什么?
- 十六、doOnNext() 和 map() 的区别
- 十七、onErrorResume() 是什么?
- 十八、引入 WebClient
- 十九、创建 WebClient
- 二十、WebClient 普通请求
- 二十一、WebClient 流式请求
- 二十二、完整的 WebClient + Flux + SseEmitter 示例
- 二十三、Controller 使用 SseEmitter
- 二十四、这段代码到底发生了什么?
- 二十五、为什么这里一定要 subscribe()?
- 二十六、但是为什么 WebFlux Controller 又经常不需要手动 subscribe()?
- 二十七、什么时候选择 WebClient + Flux + SseEmitter?
- 二十八、什么时候选择 Spring WebFlux?
- 二十九、两种方案对比
- 三十、为什么 WebFlux 项目中使用 MyBatis/JDBC 需要谨慎?
- 三十一、WebClient + Flux + SseEmitter 的完整架构
- 三十二、Flux 常用操作总结
- 三十三、一个完整的 Flux 示例
- 三十四、总结
前言
在 AI Agent、ChatGPT、智能问答等应用中,一个非常常见的需求是:
后端调用 AI 服务,AI 服务不会一次性返回完整结果,而是像打字一样,一段一段地返回内容;后端再把这些内容实时推送给前端。
例如用户发送:
请介绍一下 Spring Boot
AI 服务可能不是等待全部内容生成完成后一次返回:
Spring Boot 是一个……
它主要用于……
它具有……
而是:
Spring
Spring Boot
Spring Boot 是
Spring Boot 是一个
Spring Boot 是一个用于
……
这种模式就是流式响应(Streaming Response)。
在 Spring Boot 中实现这种场景,有两种比较常见的方案:
方案一:
Spring MVC + WebClient + Flux + SseEmitter
方案二:
Spring WebFlux + WebClient + Flux
本文将从这两个方案的区别开始,逐步介绍:
- 什么是 Spring WebFlux
- 什么是 Flux
- 什么是 Mono
- WebClient 是什么
- 为什么 WebClient 可以返回 Flux
- 什么情况下选择
WebClient + Flux + SseEmitter - 什么情况下选择 Spring WebFlux
- 如何在 Spring Boot 中引入 WebClient
- 如何调用流式 HTTP 接口
- 如何通过 SseEmitter 将数据转发给前端
subscribe()、map()、flatMap()、doOnNext()、onErrorResume()的作用- 为什么不调用
subscribe(),Flux 可能根本不会执行 - 最后通过一个完整的 AI Agent 流式调用案例串起来
一、什么是流式响应?
传统 HTTP 接口一般是:
客户端
↓
发送请求
↓
服务器处理
↓
服务器等待业务执行完成
↓
返回完整结果
↓
请求结束
例如:
@GetMapping("/hello")
public String hello() {
return "Hello World";
}
客户端必须等待服务器执行完成,然后一次性拿到:
Hello World
而流式响应不同。
服务器可以:
客户端
↓
发送请求
↓
服务器开始处理
↓
返回第一部分数据
↓
继续处理
↓
返回第二部分数据
↓
继续处理
↓
返回第三部分数据
↓
……
↓
处理完成
↓
关闭连接
例如 AI:
你好
你好,我
你好,我是
你好,我是一个
你好,我是一个AI
你好,我是一个AI助手
前端可以边接收边显示。
这就是 AI 对话中经常看到的“打字机效果”。
二、什么是 SSE?
SSE 全称是:
Server-Sent Events
即:
服务器发送事件
它是一种基于 HTTP 的服务器向客户端持续推送数据的技术。
通信方向主要是:
Server
↓
↓
↓
Client
也就是说,它特别适合:
- AI 流式回答
- 实时日志
- 任务进度
- 消息通知
- 股票行情
- 实时状态更新
例如服务器可以持续发送:
data: Hello
data: World
data: !
客户端就可以一条一条接收。
三、SseEmitter 是什么?
如果项目使用的是传统 Spring MVC,可以使用:
SseEmitter
它是 Spring MVC 提供的一个 SSE 工具类。
例如:
@GetMapping("/stream")
public SseEmitter stream() {
SseEmitter emitter = new SseEmitter(60000L);
new Thread(() -> {
try {
emitter.send("Hello");
Thread.sleep(1000);
emitter.send("World");
Thread.sleep(1000);
emitter.send("!");
emitter.complete();
} catch (Exception e) {
emitter.completeWithError(e);
}
}).start();
return emitter;
}
这里的核心是:
emitter.send(...)
每调用一次:
emitter.send(...)
就可以向客户端发送一次数据。
最终:
emitter.complete();
表示 SSE 流结束。
四、为什么会出现 WebClient + Flux + SseEmitter?
假设现在有这样的架构:
前端
↓
Spring Boot 服务
↓
AI Agent 服务
AI Agent 本身支持流式返回。
例如:
Spring Boot
↓
POST /agent/run
↓
AI Agent
AI Agent
↓
Flux
↓
数据1
数据2
数据3
数据4
……
我们的 Spring Boot 服务需要做一个“中转”。
也就是前端调用你的 Spring Boot 服务,Spring Boot 服务作为中转调用其他 AI Agent 服务返回:
流式
Spring Boot ───────────→ AI Agent
↑
│
│ 流式
│
↓
前端
于是可以使用:
WebClient
↓
Flux
↓
SseEmitter
具体来说:
AI Agent
↓
WebClient
↓
Flux<AiAgentStreamRsp>
↓
doOnNext()
↓
SseEmitter.send()
↓
前端
这就是本文的核心方案。
五、什么是 Spring WebFlux?
Spring WebFlux 是 Spring 提供的一套响应式 Web 框架。
传统 Spring Boot Web 项目通常使用:
spring-boot-starter-web
它主要基于:
Spring MVC
Servlet
Tomcat
而 WebFlux 使用:
spring-boot-starter-webflux
其核心思想是:
响应式编程 + 非阻塞 I/O
简单理解:
传统 Spring MVC 更接近:
请求
↓
Tomcat线程
↓
执行任务
↓
等待数据库 / HTTP
↓
继续执行
↓
返回结果
而响应式模型更强调:
请求
↓
发起异步操作
↓
暂时不阻塞等待
↓
数据准备完成
↓
继续处理
↓
返回结果
六、WebFlux 和 Reactor 的关系
理解 WebFlux 时经常会遇到:
Mono
Flux
它们不是 Spring MVC 的概念,而是来自:
Project Reactor
WebFlux 基于 Reactor。
可以简单理解为:
Spring WebFlux
↓
Reactor
↓
┌─────────────┐
│ │
Mono Flux
其中:
Mono<T>
表示:
0 或 1 个数据。
例如:
Mono<User>
表示未来可能得到一个 User。
而:
Flux<T>
表示:
0 到 N 个数据。
例如:
Flux<User>
表示未来可能不断得到:
User1
User2
User3
User4
...
七、什么是 Flux?
这是理解本文最重要的概念。
Flux<String>
可以简单理解为:
一个可以不断产生 String 的数据流。
例如:
Flux<String> flux =
Flux.just(
"Hello",
"World",
"!"
);
可以理解为:
Flux
│
├── Hello
├── World
└── !
但需要注意:
Flux 本身并不是数据。
它更像是:
描述“数据将如何产生、如何处理”的一个数据流定义。
八、Flux 的三个核心阶段
理解 Flux 可以把它分成三个阶段:
创建
↓
处理
↓
订阅执行
例如:
Flux<String> flux = Flux.just("A", "B", "C")
.map(String::toLowerCase)
.doOnNext(System.out::println);
此时只是:
创建了一条流水线
并不代表一定已经执行。
只有:
flux.subscribe();
之后,数据才真正开始流动。
九、为什么不调用 subscribe(),Flux 可能不会执行?
这是 Reactor 初学者最容易遇到的问题。
例如:
Flux<String> flux = Flux.just("A", "B", "C")
.map(item -> {
System.out.println("处理:" + item);
return item.toLowerCase();
});
如果代码到这里结束:
没有 subscribe()
那么:
处理:A
处理:B
处理:C
可能根本不会打印。
因为 Flux 是惰性的(Lazy)。
前面的代码主要是在构建:
数据处理流水线
真正执行需要订阅:
flux.subscribe();
执行流程:
Flux.just()
↓
map()
↓
doOnNext()
↓
subscribe()
↓
开始产生数据
↓
A
↓
B
↓
C
↓
complete
所以可以记住一句话:
Flux 是一条待执行的数据流水线,subscribe() 是启动这条流水线的重要入口。
十、Flux 的执行信号
Flux 最核心的信号可以理解为:
onNext
onError
onComplete
例如:
onNext(data1)
onNext(data2)
onNext(data3)
onComplete()
表示:
数据1
数据2
数据3
正常结束
如果发生异常:
onNext(data1)
onNext(data2)
onError(exception)
表示:
数据1
数据2
发生异常
这三个概念非常重要。
十一、subscribe() 是什么?
最简单:
flux.subscribe();
表示:
订阅这个 Flux,让它开始执行。
完整一点:
flux.subscribe(
data -> System.out.println("数据:" + data),
error -> System.out.println("异常:" + error),
() -> System.out.println("完成")
);
三个参数分别对应:
onNext
onError
onComplete
也就是:
data -> ...
处理数据。
error -> ...
处理异常。
() -> ...
处理完成。
十二、map() 是什么?
map() 用来做:
一对一的数据转换。
例如:
Flux<String> flux =
Flux.just("hello", "world")
.map(String::toUpperCase);
数据:
hello
world
经过 map:
HELLO
WORLD
可以理解为:
A
↓ map
B
所以:
map()
特别适合:
对象转换
字段转换
格式转换
数据加工
例如:
Flux<UserDTO> result =
flux.map(user -> {
UserDTO dto = new UserDTO();
dto.setName(user.getName());
return dto;
});
十三、flatMap() 是什么?
flatMap() 和 map() 是 WebFlux / Reactor 中非常重要的一组概念。
假设:
Mono<User> getUser();
现在我们需要:
先查询 User
↓
再根据 User 查询 Order
例如:
Mono<User> userMono = getUser();
然后:
userMono.flatMap(user -> getOrders(user.getId()));
如果:
getOrders()
返回:
Mono<List<Order>>
那么:
flatMap()
特别适合处理:
一个异步操作完成之后,再继续执行另一个异步操作。
十四、map 和 flatMap 的区别
可以简单记忆:
map
一对一转换
flatMap
一个异步结果 → 另一个异步结果
例如:
Mono<User> userMono;
如果只是:
userMono.map(user -> user.getName());
结果:
Mono<String>
这是普通的数据转换。
而:
userMono.flatMap(user -> userService.querySomething(user));
如果 querySomething() 本身返回:
Mono<Result>
那么就应该使用:
flatMap()
否则容易出现嵌套:
Mono<Mono<Result>>
而 flatMap 可以把它“摊平”为:
Mono<Result>
所以:
map
T → R
flatMap
T → Mono<R>
可以作为初学阶段的一个非常重要的记忆方式。
十五、doOnNext() 是什么?
例如:
Flux<String> flux =
Flux.just("A", "B", "C")
.doOnNext(data ->
System.out.println("收到:" + data)
);
每产生一个数据:
A
↓
doOnNext
B
↓
doOnNext
C
↓
doOnNext
因此:
doOnNext()
特别适合做:
- 日志
- 监控
- 调试
- 发送消息
- 统计
- 副作用操作
例如:
.doOnNext(data -> {
log.info("收到AI数据:{}", data);
})
十六、doOnNext() 和 map() 的区别
这是一个非常容易混淆的问题。
map():
.map(data -> {
return newData;
})
它的主要目的:
修改 / 转换数据。
而:
.doOnNext(data -> {
log.info("{}", data);
})
它主要用于:
观察数据、执行副作用。
例如:
Flux<String> flux =
Flux.just("a", "b", "c")
.map(String::toUpperCase)
.doOnNext(data ->
System.out.println("当前数据:" + data)
);
数据经过:
a
↓ map
A
↓ doOnNext
打印 A
b
↓ map
B
↓ doOnNext
打印 B
十七、onErrorResume() 是什么?
它用于:
发生异常后提供一个备用 Flux。
例如:
Flux<String> flux =
Flux.just("A", "B", "C")
.map(data -> {
if ("B".equals(data)) {
throw new RuntimeException("测试异常");
}
return data;
})
.onErrorResume(error -> {
return Flux.just("发生异常,返回备用数据");
});
执行:
A
↓
B
↓
异常
↓
onErrorResume
↓
发生异常,返回备用数据
所以:
onErrorResume()
可以理解为:
出现异常以后,不一定直接让整个流失败,而是切换到另外一个数据流。
例如远程服务调用:
webClient.get()
.retrieve()
.bodyToFlux(String.class)
.onErrorResume(error -> {
log.error("远程服务异常", error);
return Flux.empty();
});
发生异常后:
远程服务异常
↓
onErrorResume
↓
Flux.empty()
↓
流正常结束
十八、引入 WebClient
如果项目已经是 Spring Boot 项目,可以引入:
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-webflux</artifactId>
</dependency>
即使你的项目 Controller 仍然使用:
Spring MVC
也可以通过这个依赖使用:
WebClient
Reactor
Mono
Flux
需要注意:
引入
spring-boot-starter-webflux并不意味着你的项目必须全部改成 WebFlux。
如果项目同时存在:
spring-boot-starter-web
和:
spring-boot-starter-webflux
具体 Web 应用运行模式还涉及 Spring Boot 的自动配置和项目依赖情况,因此生产项目中最好明确项目到底采用 MVC 还是 WebFlux,而不是无意中混用。
十九、创建 WebClient
最简单:
WebClient webClient = WebClient.builder()
.baseUrl("http://localhost:8081")
.build();
也可以注册成 Spring Bean:
@Configuration
public class WebClientConfig {
@Bean
public WebClient webClient() {
return WebClient.builder()
.baseUrl("http://localhost:8081")
.build();
}
}
然后:
@Resource
private WebClient webClient;
二十、WebClient 普通请求
例如:
public Mono<User> getUser(Long userId) {
return webClient.get()
.uri("/user/{id}", userId)
.retrieve()
.bodyToMono(User.class);
}
这里:
bodyToMono(User.class)
表示:
我预计远程接口最终返回一个 User。
所以:
Mono<User>
二十一、WebClient 流式请求
如果远程服务本身是流式接口:
public Flux<AiAgentStreamRsp> agentStream(
LlmAgentReqDTO request) {
return webClient.post()
.uri("/agent/run")
.bodyValue(request)
.retrieve()
.bodyToFlux(AiAgentStreamRsp.class);
}
这里:
bodyToFlux(AiAgentStreamRsp.class)
表示:
远程 HTTP 响应是一个持续产生
AiAgentStreamRsp的数据流。
最终:
Flux<AiAgentStreamRsp>
二十二、完整的 WebClient + Flux + SseEmitter 示例
下面模拟一个 AI Agent。
假设远程 AI 服务:
POST http://localhost:8081/agent/run
支持流式返回:
数据1
数据2
数据3
数据4
我们的 Spring Boot 服务需要:
前端
↓
/agent-stream
↓
Spring Boot
↓
WebClient
↓
AI Agent
↓
Flux
↓
SseEmitter
↓
前端
1. DTO
@Data
public class LlmAgentReqDTO {
private String question;
}
返回对象:
@Data
public class AiAgentStreamRsp {
private String content;
}
2. Service
@Service
public class LlmApiService {
private final WebClient webClient;
public LlmApiService(WebClient webClient) {
this.webClient = webClient;
}
public Flux<AiAgentStreamRsp> agentStream(
LlmAgentReqDTO request) {
return webClient.post()
.uri("/agent/run")
.bodyValue(request)
.retrieve()
.bodyToFlux(AiAgentStreamRsp.class);
}
}
这里最核心的一行:
.bodyToFlux(AiAgentStreamRsp.class)
最终得到:
Flux<AiAgentStreamRsp>
二十三、Controller 使用 SseEmitter
@RestController
@RequestMapping("/api")
public class AgentController {
private final LlmApiService llmApiService;
public AgentController(LlmApiService llmApiService) {
this.llmApiService = llmApiService;
}
@PostMapping("/agent-stream")
public SseEmitter agentStream(
@RequestBody LlmAgentReqDTO request) {
// 60秒超时
SseEmitter emitter = new SseEmitter(60000L);
llmApiService.agentStream(request)
.doOnNext(data -> {
try {
emitter.send(
SseEmitter.event()
.name("message")
.data(data)
);
} catch (IOException e) {
emitter.completeWithError(e);
}
})
.doOnError(emitter::completeWithError)
.doOnComplete(emitter::complete)
.subscribe();
return emitter;
}
}
二十四、这段代码到底发生了什么?
重点分析:
llmApiService.agentStream(request)
返回:
Flux<AiAgentStreamRsp>
然后:
.doOnNext(...)
表示:
Flux 每产生一个
AiAgentStreamRsp,执行一次这里的代码。
例如:
AI Agent
↓
数据1
执行:
doOnNext(data1)
然后:
emitter.send(data1);
接着:
AI Agent
↓
数据2
执行:
doOnNext(data2)
然后:
emitter.send(data2);
所以:
Flux<AiAgentStreamRsp>
↓
data1 → emitter.send(data1)
↓
data2 → emitter.send(data2)
↓
data3 → emitter.send(data3)
↓
data4 → emitter.send(data4)
二十五、为什么这里一定要 subscribe()?
因为:
llmApiService.agentStream(request)
.doOnNext(...)
.doOnError(...)
.doOnComplete(...)
只是构建了一条 Reactor 流水线。
真正启动:
.subscribe();
所以:
.subscribe();
相当于:
开始订阅远程 Agent 的数据流。
没有:
subscribe()
可能出现:
创建 Flux
↓
配置 doOnNext
↓
配置 doOnError
↓
配置 doOnComplete
↓
代码结束
但没有真正订阅。
于是:
AI Agent 请求
可能根本没有真正执行
这就是为什么 Reactor 被称为具有惰性执行特征。
二十六、但是为什么 WebFlux Controller 又经常不需要手动 subscribe()?
这是一个非常重要的问题。
例如 WebFlux:
@GetMapping(
value = "/stream",
produces = MediaType.TEXT_EVENT_STREAM_VALUE
)
public Flux<String> stream() {
return Flux.just(
"A",
"B",
"C"
);
}
这里没有:
.subscribe();
为什么仍然可以执行?
因为:
Spring WebFlux 会替你订阅 Controller 返回的 Publisher。
也就是说:
Controller
↓
return Flux
↓
Spring WebFlux
↓
subscribe()
↓
开始执行
所以:
普通 Reactor 代码
Flux<String> flux = ...;
flux.subscribe();
需要自己订阅。
WebFlux Controller
public Flux<String> stream() {
return flux;
}
通常不要自己:
flux.subscribe();
而是:
return Flux
让 WebFlux 框架负责订阅。
这是一个非常重要的原则:
在 WebFlux Controller 中,通常不要手动 subscribe();把 Flux 返回给框架即可。
二十七、什么时候选择 WebClient + Flux + SseEmitter?
对于已经存在的 Spring MVC 项目,这个方案非常实用。
例如你的项目目前:
Spring Boot
Spring MVC
Tomcat
MyBatis
MyBatis-Plus
Feign
Servlet Filter
Interceptor
项目整体都是传统 MVC 架构。
现在突然增加:
AI Agent 流式问答
这个时候没有必要为了一个接口把整个项目改成 WebFlux。
可以:
原有项目
↓
Spring MVC
↓
SseEmitter
远程 AI
↑
WebClient
↑
Flux
也就是:
Spring MVC
+
WebClient
+
Flux
+
SseEmitter
这种方式改造成本比较低。
特别适合:
- 老项目
- Spring MVC 项目
- 只有少量 SSE 接口
- AI 流式接口
- 需要调用第三方流式 HTTP API
- 不希望修改整个项目架构
二十八、什么时候选择 Spring WebFlux?
如果你的系统本身就是:
高并发 + 大量 I/O + 非阻塞编程
那么可以考虑 WebFlux。
例如:
API Gateway
实时数据系统
高并发流式服务
实时消息服务
大量 HTTP 调用
大量异步 I/O
如果整个项目从设计之初就是响应式架构:
Controller
↓
Service
↓
WebClient
↓
Reactive Database Driver
那么 WebFlux 的优势更加明显。
例如:
@GetMapping(
value = "/agent-stream",
produces = MediaType.TEXT_EVENT_STREAM_VALUE
)
public Flux<AiAgentStreamRsp> agentStream(
@RequestBody LlmAgentReqDTO request) {
return llmApiService.agentStream(request);
}
这里不需要:
SseEmitter
也不需要:
doOnNext(data -> emitter.send(data))
更不需要:
subscribe();
整个链路直接是:
AI Agent
↓
WebClient
↓
Flux
↓
WebFlux Controller
↓
SSE
↓
前端
二十九、两种方案对比
| 对比项 | Spring MVC + WebClient + Flux + SseEmitter | Spring WebFlux |
|---|---|---|
| 项目基础 | Spring MVC | WebFlux |
| SSE | SseEmitter | Flux |
| HTTP Client | WebClient | WebClient |
| 响应式 | 局部使用 | 全链路响应式 |
| subscribe | 通常需要手动桥接 | Controller 通常不需要 |
| 改造成本 | 低 | 较高 |
| 老项目适合度 | 很高 | 一般 |
| 少量流式接口 | 很适合 | 没必要强行改 |
| 全链路异步 | 一般 | 更适合 |
| MyBatis/JDBC | 可以继续使用 | 会削弱非阻塞优势 |
| 学习成本 | 相对较低 | 较高 |
三十、为什么 WebFlux 项目中使用 MyBatis/JDBC 需要谨慎?
WebFlux 最大的特点之一是:
非阻塞
但是传统 JDBC:
线程
↓
执行 SQL
↓
等待数据库
↓
数据库返回
属于阻塞式 I/O。
如果:
WebFlux
↓
Service
↓
MyBatis
↓
JDBC
那么数据库查询依然可能阻塞线程。
因此:
使用 WebFlux 并不意味着整个系统自动变成非阻塞。
如果希望完整发挥响应式架构优势,数据库层通常也需要响应式驱动,例如:
R2DBC
也就是说:
完整响应式链路:
WebFlux
↓
Reactor
↓
WebClient
↓
R2DBC
↓
Database
而:
WebFlux
↓
MyBatis
↓
JDBC
则仍然存在阻塞 I/O。
所以对于一个已经大量使用 MyBatis 的传统 Spring Boot 项目,没有必要仅仅为了一个 AI 流式接口就全面切换 WebFlux。
三十一、WebClient + Flux + SseEmitter 的完整架构
最终可以把本文的方案总结成:
┌─────────────────────┐
│ 前端 Browser │
└──────────┬──────────┘
│
│ SSE
▼
┌─────────────────────┐
│ Spring MVC │
│ │
│ SseEmitter │
└──────────┬──────────┘
│
emitter.send()
│
▼
┌─────────────────────┐
│ Flux │
│ │
│ Flux<AiAgentRsp> │
└──────────┬──────────┘
│
WebClient
│
▼
┌─────────────────────┐
│ AI Agent │
│ │
│ Streaming API │
└─────────────────────┘
这里三个核心组件各自负责不同的事情:
WebClient
负责:
调用远程 HTTP 服务
webClient.post()
Flux
负责:
表示连续产生的数据
Flux<AiAgentStreamRsp>
SseEmitter
负责:
在 Spring MVC 中将数据持续推送给浏览器
emitter.send(...)
所以可以简单记忆:
WebClient = 调别人
Flux = 接收/处理一串数据
SseEmitter = 推给前端
三十二、Flux 常用操作总结
| 方法 | 作用 | 典型用途 |
|---|---|---|
subscribe() | 订阅并启动流 | 手动执行 Flux |
map() | 一对一转换 | 对象转换、字段转换 |
flatMap() | 异步展开/串联 | 一个异步操作接另一个异步操作 |
doOnNext() | 观察每条数据 | 日志、监控、副作用 |
doOnComplete() | 正常结束时执行 | 清理资源、结束 SSE |
doOnError() | 发生异常时执行 | 日志、异常通知 |
onErrorResume() | 异常后切换备用流 | 降级、兜底 |
filter() | 过滤数据 | 条件过滤 |
take() | 只取前 N 个 | 限制数据量 |
timeout() | 超时控制 | 防止请求无限等待 |
三十三、一个完整的 Flux 示例
最后通过一个例子把前面的方法串起来:
Flux.just("1", "2", "3")
.map(Integer::parseInt)
.filter(number -> number > 1)
.doOnNext(number ->
System.out.println("当前数据:" + number)
)
.map(number ->
"数字:" + number
)
.onErrorResume(error -> {
System.out.println("发生异常:" + error.getMessage());
return Flux.just("备用数据");
})
.subscribe(
data -> System.out.println("最终结果:" + data),
error -> System.out.println("最终异常:" + error),
() -> System.out.println("数据流结束")
);
执行过程可以理解为:
1
↓ map
1
↓ filter
过滤掉
2
↓ map
2
↓ filter
2
↓ doOnNext
打印
↓ map
数字:2
↓ subscribe
最终结果
3
↓
3
↓
doOnNext
↓
数字:3
↓
subscribe
最终结果
↓
complete
三十四、总结
如果你的项目本身是传统 Spring MVC 项目,现在需要增加一个 AI 流式接口,那么:
WebClient + Flux + SseEmitter
是一个非常实用的方案。
完整链路:
前端
↓
SSE
↓
SseEmitter
↓
Flux
↓
WebClient
↓
AI Agent
其中:
WebClient
负责调用远程 HTTP 服务。
Flux
负责表示和处理连续的数据流。
SseEmitter
负责在 Spring MVC 中持续向客户端发送 SSE 数据。
而如果整个项目本身就是响应式架构,并且希望从:
Controller
↓
Service
↓
HTTP Client
↓
Database
都采用非阻塞模型,那么可以直接选择:
Spring WebFlux
此时 Controller 可以直接返回:
public Flux<AiAgentStreamRsp> agentStream(...) {
return llmApiService.agentStream(...);
}
而不需要:
SseEmitter
doOnNext()
subscribe()
最后一定要理解 Reactor 的一个核心思想:
Flux 是一条数据流,操作符是在构建这条数据流的处理链,而 subscribe() 才是订阅这条数据流、触发执行的重要入口。
不过在 Spring WebFlux Controller 中,通常不需要自己调用 subscribe(),因为 WebFlux 会负责订阅 Controller 返回的 Publisher。
因此,理解:
Flux
↓
subscribe
↓
onNext
↓
onError / onComplete
是学习 WebClient、WebFlux、SSE 以及 AI 流式接口的基础。
更多推荐

所有评论(0)