前言

在 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 + SseEmitterSpring WebFlux
项目基础Spring MVCWebFlux
SSESseEmitterFlux
HTTP ClientWebClientWebClient
响应式局部使用全链路响应式
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 流式接口的基础。

Logo

AtomGit AI 社区提供模型库、数据集、Agent、Token等资源

更多推荐