Java 响应式编程 - Spring WebFlux 简介及基本案例

Spring WebFlux 是 Spring 团队为了应对高并发、高吞吐、低延迟的现代云原生/微服务场景,基于 Project Reactor 核心打造的、完全非阻塞(Non-Blocking)、支持背压(Backpressure)的全新一代 Web 框架。


和传统web的对比

生态位上的对应

为了看清它为什么能以极少的资源吞下海量的并发,我们需要对比它与传统 Spring MVC 的底层物理模型:

  • 旧世界(Spring MVC):采用 Thread-per-request(一请求一线程) 模型。Tomcat 默认开 200 个线程。一个 HTTP 请求进来,就死死咬住一个线程不放。如果这个请求要去调一个耗时 3 秒的微服务,这个线程就会原地阻塞(卡死) 3 秒。并发请求一旦超过 200,后续请求只能排队,甚至直接超时崩溃。
  • 新世界(Spring WebFlux):默认基于 Netty。它采用 Event Loop(事件循环) 模型。后台通常只有和 CPU 核心数相同的物理线程(比如 8 核就只有 8 个线程)。请求进来时,线程只负责“登记这个事件(注册回调)”,随后立刻去接待下一个请求。当耗时 3 秒的远程调用结果返回时,Netty 会发出一个通知,Event Loop 线程再回过头来把结果吐给前端。线程永远在高速运转,绝不为任何 I/O 阻塞一分一秒。

按照上图,虽然生态位完全对等,但在底层物理形态上,新世界的演进其实比旧世界更激进、更彻底:

  • 旧世界的 Servlet 容器是有“边界感”的:在传统架构中,Tomcat(容器)和 Spring MVC(框架)是两套独立的源码体系。Tomcat 是一个独立运行的 Java 进程,Spring MVC 是一个打成 WAR 包塞进 Tomcat 里的应用。Tomcat 内部维护着自己的 Http11Processor 线程池,Spring 只是在这个线程池里运行的业务代码。
  • 新世界的 Reactor 和 WebFlux 是“水乳交融”的:在响应式微服务中(基于 Netty),Netty、Reactor、WebFlux 在运行时被高度抽象、并彻底融合成了一条密不可分的响应式管道。Netty 的 EventLoop 线程,在通过 Reactor 转换后,会直接一路无缝穿透到你的 WebFlux Controller,甚至穿透到你的 Context 链路追踪和 R2DBC 数据库层。它们之间不再有“容器线程”和“业务框架”的明显割裂。


组件上的对应

在这两套组件的对齐中,有两处地方发生了颠覆性的底层质变,在开发中需要特别注意:

  • HttpServletRequest –> ServerWebExchange:
    • 在 Spring MVC 中,你如果想在 Controller 之外传递一些请求元数据(比如权限、TraceId),通常会祭出 ThreadLocal(如 RequestContextHolder)。
    • 在 WebFlux 中,不要再碰 ThreadLocal。所有的请求、响应、甚至是 Session 属性以及我们在 Reactor Context,全部被死死封装在 ServerWebExchange 这一个对象里。它就像一个传输舱,顺着响应式管道一路向下游流动,任何位置的操作符都能向它存取数据。
  • RestTemplate –> WebClient:
    • RestTemplate 的底层物理模型是。发起 HTTP 连接的那个线程,必须原地死等对方服务器返回数据。
    • WebClient 底层直接基于 Netty 的无阻塞网络通道。当它发送请求后,当前线程立刻收工去干别的事。当对方服务器响应时,底层网络驱动会触发一个事件,由 EventLoop 线程把结果打包成 Mono 顺着水管冲下来。在构建微服务集群(如调用三方支付、查询用户中心)时,换用 WebClient 能让你的服务器 QPS 发生质的飞跃。

结合我们之前写的 用 Reactor API 写 Server 。实际上,Spring WebFlux 在底层做的事情,本质上就是用注解和 IoC 容器,把我们上面手写的这段 HttpServer.create().route(…) 代码给优雅地包装了起来。 我们的 @RestController 最终被解析成了 routes.get() 或 routes.post()。


主要组件介绍

在前面的对比中,我们从宏观层面对齐了组件的生态位。现在,我们聚焦 Spring WebFlux 框架内部,深度拆解它是如何各司其职、在 Netty 之上搭建起整条非阻塞响应式流水线的。Spring WebFlux 的底层组件设计高度借鉴了 Spring MVC 的经典架构,但它们全部被“响应式化”改造,所有的核心方法都返回 MonoFlux


$$\text{Netty 接收字节} \longrightarrow \text{HttpHandler 抽象封装} \longrightarrow \text{WebFilter 拦截过滤} \longrightarrow \text{DispatcherHandler 核心统筹}$$

$$\downarrow$$

$$\text{HandlerResultHandler 编码写出} \longleftarrow \text{HandlerAdapter 执行参数绑定} \longleftarrow \text{HandlerMapping 匹配路由}$$


  • 核心入口枢纽层(前端控制器)
    • HttpHandler
      • 职责:它是 WebFlux 最底层的通用网关入口,不属于 Spring 的核心业务逻辑,而是用来抹平底层服务器物理差异的抽象层。
      • 原理:无论是 Netty、Tomcat、Jetty 还是 Undertow 接收到请求,都会统一被包装并转换为 HttpHandler 的底层请求和响应,随后无缝递交给上层的 Spring 框架。
    • WebHandler / DispatcherHandler - 核心总暴风眼
      • 职责:相当于 Spring MVC 里的 DispatcherServlet。它是整个框架的前端控制器(Front Controller),也是整个响应式 Web 请求的核心枢纽。
      • 原理:它本身不处理任何具体的业务。它的工作是组合并调度流中的三大战略副官(HandlerMapping、HandlerAdapter、HandlerResultHandler),负责请求的接收、分发、适配和结果写出。
  • 请求路由寻址层(找谁来干活)
    • HandlerMapping (地址簿)
      • 职责:负责寻找处理该请求的负责人(Handler)。
      • 原理:当一个请求 URL(如 /api/chat)进来时,HandlerMapping 会在内存的映射表里检索。如果是注解派,它会找出对应的 @RestController 里的 Method。如果是函数式派,它会定位到对应的 RouterFunction 路由。最终会把请求和拦截器打包,返回一个 Mono 信号,代表未来会吐出一个处理器。
  • 业务执行适配层(真正的打工人)
    • HandlerAdapter (通用转换器)
      • 职责:负责真正去调用和执行上面找到的那个处理器方法(Handler)。
      • 原理:因为 WebFlux 支持注解控制器、函数式 Lambda 甚至原生的底层自定义处理器,它们的入参和执行方式各不相同。HandlerAdapter 就像一个多功能插头适配器,屏蔽了这些差异。它会负责参数绑定(把请求体 JSON 解析为 Java 对象),然后开足马力执行你的业务代码,并最终返回一个统一的 HandlerResult 对象。
  • 结果序列化写出层(出水口管道)
    • HandlerResultHandler (成品包装工)
      • 职责:负责把业务方法执行完返回的结晶(如 Mono<User>、Flux<String>、视图模板),序列化、包装并写入 HTTP 响应体。
      • 核心子类生态: WebFlux 会根据你方法的返回值类型和 produces 媒体类型,自动派发不同的包装工:
        • ResponseBodyResultHandler:处理带有 @ResponseBody 的注解或常规 JSON 返回值。
        • ServerResponseResultHandler:专门处理函数式路由流派返回的 ServerResponse。
        • ViewResolutionResultHandler:如果你用了 Thymeleaf,它负责把数据塞给 HTML 模板,渲染出动态网页。
    • HttpMessageReader / HttpMessageWriter (响应式编解码器)
      • 职责:负责 Java 对象与网络字节流(DataBuffer)之间的非阻塞编解码(Codec)。
      • 原理:在底层,它们利用 Jackson 响应式读写器,当 Netty 抓到一段 TCP 字节块时,HttpMessageReader 顺着 Flux 流边读边解析;当要吐出 ChatGPT 打字机字符时,HttpMessageWriter 负责实时把 String 字符打上 data:\n\n 并无阻塞地写入网卡。
  • WebFilter 与 WebExceptionHandler
    • WebFilter(响应式过滤器):拦截一切的前哨站。通过 exchange.getLog().trace(…) 或者 Context 的注入,它在 DispatcherHandler 运行前和运行后织入切面。它是完全非阻塞的链式调用(filterChain.filter(exchange))。
    • WebExceptionHandler(响应式异常总监):只要整个响应式管道内部(从读取参数、查库、到结果写出)发生任何暴雷抛出异常,都会被它兜底捕获。它会将异常转化为体面的错误 JSON 或 404 页面,确保即使业务崩溃,底层的 Netty EventLoop 也绝不卡死。


独特的发送事件SSE

由于底层天然是持久的、非阻塞的响应式管道,WebFlux 做“数据流式长连接推送(SSE, Server-Sent Events)” 具有降维打击般的优势。典型场景如股票大屏不间断刷新、或者是类似于 ChatGPT 那种打字机流式字符吐出。在传统 Spring MVC 里实现这个极其痛苦且耗费线程,而在 WebFlux 里,只需要 3 行代码:

1
2
3
4
5
6
@GetMapping(value = "/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<String> streamingData() {
// 🌟 每隔 1 秒,自动向前端源源不断地异步推送一条实时动态时钟数据
return Flux.interval(Duration.ofSeconds(1))
.map(tick -> "⏰ 实时数据流播报,当前第 " + tick + " 次滴答");
}


实际该怎么选

很多人一听 WebFlux 吞吐量大,就盲目地把所有老项目往上迁移,结果往往被碰得头破血流。关于选型:

  • 全链路非阻塞才是真理:响应式管道具有“短板效应”。如果你的 Controller 是 WebFlux,Service 是 Reactor,但你底层的数据库用的依然是传统的、阻塞的 JDBC(如 MyBatis + 传统塑料 MySQL 驱动),那么当请求遇到数据库查盘时,线程依然会被卡死。正宗搭配是:WebFlux + WebClient + R2DBC(响应式数据库驱动)/ Reactive Mongo / Redis。
  • 它不一定会让单个请求“变快”,但它能承受“巨量并发”:如果一个请求查库耗时 50ms,在 Spring MVC 里是 50ms,在 WebFlux 里可能还是 50ms,甚至因为线程上下文切换和算力开销还会变成 52ms。但是,当 10 万个并发请求同时涌入时,Spring MVC 早已瘫痪,而 WebFlux 依然能面不改色地用极低的内存把所有请求平稳消化。
  • 选型的经验:
    • 适合 WebFlux 的场景:高并发的分布式微服务网关(如 Spring Cloud Gateway)、长连接聊天服务器(IM)、高频数据同步监控中心、流式 AI 文本生成接口。
    • 适合 Spring MVC 的场景:常规的大型企业级 ERP 系统、传统的 CRUD 管理后台等。


两种开发模式

Spring WebFlux 非常温情地提供了两套完全不同的编程范式,无论你是保守派还是激进派,都能找到舒服的姿势。

注解驱动派

它的写法和传统的 Spring MVC 一模一样,保留了 @RestController、@GetMapping、@PostMapping 等全部全家桶注解。唯一的物理区别是:方法的返回值不再是具体对象,而是包装成了 Mono<T> 或 Flux<T>。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
@RestController
@RequestMapping("/api/users")
public class UserReactiveController {

@Autowired
private UserService userService; // 假设这是一个响应式的 Service

@GetMapping("/{id}")
public Mono<User> getUserById(@PathVariable String id) {
// 🌟 返回一个 Mono 代理,Spring 会在底层异步拧开水龙头放水
return userService.findUserById(id)
.defaultIfEmpty(new User("guest")) // 空值兜底
.contextWrite(ctx -> ctx.put("TRACE_ID", "REQ-" + id)); // 链路追踪注入
}
}


函数式路由派

这是新世界的专属写法语流派。它把“路由隐式路由(由谁来处理这个请求)”和“业务逻辑处理”彻底解耦。通过 RouterFunction 和 HandlerFunction,像拼乐高积木一样声明式地编排你的路由表。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
import static org.springframework.web.reactive.function.server.RequestPredicates.*;
import static org.springframework.web.reactive.function.server.RouterFunctions.*;

@Configuration
public class UserRouter {

@Bean
public RouterFunction<ServerResponse> route(UserHandler handler) {
return route(GET("/f/users/{id}"), handler::getUser)
.andRoute(POST("/f/users"), handler::createUser);
}
}

@Component
class UserHandler {
public Mono<ServerResponse> getUser(ServerRequest request) {
String id = request.pathVariable("id");
Mono<User> userMono = Mono.just(new User(id, "KJ"));
return ServerResponse.ok().body(userMono, User.class);
}
}


实际的选择

实际上绝大多数(约 85% 以上)的团队在实际开发中,依然选择注解式开发。因为:

  • 对大多数从 Spring MVC 转型过来的 Java 程序员来说,@RestController、@PostMapping、@PathVariable 早就刻进了 DNA 里。
  • 现有的很多框架(如 Knife4j/Swagger 接口文档生成、Spring Security 早期配置、各类旧的验证注解 @Valid)对注解模式的支持是最完美的。如果是老项目重构,改改返回值(把 User 变成 Mono<User>)就能直接跑起来。
  • 在一个大型团队中,往往同时存在 Spring MVC 项目和 WebFlux 项目(如网关)。让所有人统一使用一套基于注解的控制层规范,能极大降低人员跨项目的协作沟通成本。

那么函数式路由式的开发模式就没人用了吗?并不是,它在特定的领域也非常受欢迎。它的核心优势在于:

  • 启动性能与内存开销的极致追求:传统的注解流派在系统启动时,需要利用反射疯狂扫描所有的类、解析注解、并建立复杂的 URL 映射树(HandlerMapping)。这会导致启动变慢,且常驻内存偏大。
    而函数式流派是纯 Java 代码编排,不需要反射和注解扫描,配合 GraalVM 打包成 AOT(原生可执行文件/Native Image) 时,启动速度可以从 5 秒直接飙升到 50 毫秒,内存占用缩减 80%。
  • 真正的端到端响应式声明美感:在函数式路由里,你的路由链路逻辑、跨域配置、拦截过滤可以写在同一个文件里,这非常契合函数式编程(FP)思想,在编写轻量级网关或单一职责的微服务时极其丝滑。


WebFlux Hello World

我们现在使用最纯正的 Spring Boot 3.x + Spring WebFlux,从零手写一个最简单的响应式微服务。我们这里开始构建的模式奕然采取大家普遍使用的注解式模式。(最后也会展示函数式路由怎么写)

项目依赖

构建 WebFlux 项目,引入的是 spring-boot-starter-webflux,它底层不包含 Tomcat,而是原生自带了高效的 Netty 服务器。

1
2
3
4
5
6
7
8
9
10
11
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-webflux</artifactId>
</dependency>
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<optional>true</optional>
</dependency>
</dependencies>


启动类

1
2
3
4
5
6
7
8
9
10
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;

@SpringBootApplication
public class WebFluxApplication {
public static void main(String[] args) {
SpringApplication.run(WebFluxApplication.class, args);
System.out.println("KooHub WebFlux 服务已在 Netty 上全速运转...");
}
}


业务类

UserDemoController

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
@RestController
@RequestMapping("/api/users")
public class UserDemoController {

/**
* 🧱 关卡 1:返回单个对象的异步包装 (Mono)
* 访问地址:http://localhost:8080/api/users/1
* 这个接口返回一个包含单数据的 Mono。前端访问时,会瞬间拿到结果。
*/
@GetMapping("/{id}")
public Mono<User> getUserById(@PathVariable String id) {
// 声明式包装:代表“未来某一刻会吐出一个 User 对象”的代理
return Mono.just(new User(id, "KJ"));
}

/**
* 🌊 关卡 2:返回多个对象的异步流 (Flux)
* 访问地址:http://localhost:8080/api/users
* 这个接口返回包含多个数据的 Flux。虽然有多个数据,但对前端来说表现像普通的 JSON 数组。
*/
@GetMapping
public Flux<User> getAllUsers() {
// 依次喷射两个数据,前端表现形式为标准的 JSON 数组 [{}, {}]
return Flux.just(
new User("101", "张三"),
new User("102", "李四")
);
}

/**
* ⏳ 关卡 3:ChatGPT 同款流式文本推流 (Server-Sent Events)
* 访问地址:http://localhost:8080/api/users/stream
* 注意:必须指定 produces 为 TEXT_EVENT_STREAM_VALUE
*
* 这是 WebFlux 真正炫技的舞台。我们利用 Flux.interval 和响应式非阻塞特性,每隔 1 秒向前端吐一条数据,
* 完美模拟 ChatGPT 的打字机输出或者实时大屏。在传统 Web 里,一个线程会被这几秒钟死死卡住;而在 WebFlux
* 里,执行这个接口期间,Netty 线程依然可以去接待另外上万个用户请求。
*/
@GetMapping(value = "/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<String> streamData() {
return Flux.interval(Duration.ofSeconds(1)) // 每隔1秒发出一个递增数字 (0, 1, 2...)
.map(seq -> "[Owlias 实时同步] 当前正在推送第 " + seq + " 秒的系统监控日志...") // 加工转换
.take(5); // 好就收:推完 5 条后,流自动发送 onComplete 信号安全闭环
}
}

User

1
2
3
4
5
6
7
@Data
@AllArgsConstructor
@NoArgsConstructor
public class User {
private String id;
private String username;
}


改造成函数式路由模式

为了让你对这两种皮囊都有绝对的掌控力,我们这就把刚刚写的最简项目里的 “单体返回” 和 “批量返回” 两个接口,用函数式路由重写一遍,让你一眼看穿它们的映射关系。


编写业务处理器

它相当于把原本 Controller 里的方法体抽离出来,入参统一变成 ServerRequest,返回值统一变成 Mono。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
package com.koohub.webflux.handler;

import com.zdemo.scloud.user.User;
import org.springframework.http.MediaType;
import org.springframework.stereotype.Component;
import org.springframework.web.reactive.function.server.ServerRequest;
import org.springframework.web.reactive.function.server.ServerResponse;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;

@Component
public class UserHandler {

// 对应原来的 GET /api/users/{id}
public Mono<ServerResponse> getUserById(ServerRequest request) {
String id = request.pathVariable("id");
Mono<User> userMono = Mono.just(new User(id, "KJ"));

// 用 ServerResponse 的流式 API 将 Mono 响应体打包吐出
return ServerResponse.ok()
.contentType(MediaType.APPLICATION_JSON)
.body(userMono, User.class);
}

// 对应原来的 GET /api/users
public Mono<ServerResponse> getAllUsers(ServerRequest request) {
Flux<User> userFlux = Flux.just(
new User("101", "koohub-functional-1"),
new User("102", "owlias-functional-2")
);
return ServerResponse.ok()
.contentType(MediaType.APPLICATION_JSON)
.body(userFlux, User.class);
}
}


编写乐高路由表

在这里,用完全声明式的代码把请求的 URL、Method 与上面的 Handler 方法死死绑定在一起。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
package com.koohub.webflux.router;

import com.zdemo.scloud.handler.UserHandler;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.http.MediaType;
import org.springframework.web.reactive.function.server.RouterFunction;
import org.springframework.web.reactive.function.server.ServerResponse;
import static org.springframework.web.reactive.function.server.RequestPredicates.GET;
import static org.springframework.web.reactive.function.server.RequestPredicates.accept;
import static org.springframework.web.reactive.function.server.RouterFunctions.route;

@Configuration
public class UserRouter {

@Bean
public RouterFunction<ServerResponse> userRoutes(UserHandler handler) {
// 像拼积木一样,把请求谓词和处理器函数编排成一张清晰的路由映射表
return route(GET("/f/users/{id}").and(accept(MediaType.APPLICATION_JSON)), handler::getUserById)
.andRoute(GET("/f/users").and(accept(MediaType.APPLICATION_JSON)), handler::getAllUsers);
}
}


WebFlux 模拟 ChatGPT 交互

前面我们使用原生 Reactor Netty 写了一个模拟 ChatGPT 交互的 小案例,相信你已经被原生 API 折磨得够呛。现在我们使用 WebFlux 实现这个简单的案例。

后端代码

我们不需要手动去管理 HttpServerResponse,更不需要手动去写 data: \n\n 这样的协议拼装。WebFlux 内部的 HttpMessageWriter 只要看到你声明的 produces = MediaType.TEXT_EVENT_STREAM_VALUE,就会自动把 Flux 里的每一个元素包装成标准的 SSE 事件流格式输送给前端。

ChatStreamController

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
import org.springframework.http.MediaType;
import org.springframework.web.bind.annotation.*;
import reactor.core.publisher.Flux;
import java.time.Duration;

@RestController
@CrossOrigin(origins = "*") // 允许跨域,方便前后端无缝联调
@RequestMapping("/api")
public class ChatStreamController {

/**
* ChatGPT 同款流式对话接口
* 直接返回 Flux<String>,WebFlux 自动将其转化为高效的异步 SSE 流式长连接
*/
@PostMapping(value = "/chat", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<String> chatWithAi(@RequestBody String userPrompt) {
System.out.println("🤖 [WebFlux 收到提问]: " + userPrompt);

// 1 模拟大模型生成的硬核回答文本
String aiAnswer = "你好!我是基于 Spring WebFlux 架构构建的现代 AI 助手。"
+ "在 WebFlux 的世界里,传统的组件如 DispatcherServlet 进化为了全新的异步响应式 DispatcherHandler。"
+ "当请求涌入时,Netty 的 Event Loop 线程只负责注册事件,随后立刻去接待下一个请求。"
+ "现在你看到的每一个字,都是顺着响应式管道无阻塞、带背压地实时喷射出来的。这就是高并发的魅力!";

// 2 将一整段长文本切碎成单个字符数组
String[] chars = aiAnswer.split("");

// 3 利用高频时钟将字符像机关枪一样源源不断地吐给前端
return Flux.interval(Duration.ofMillis(40)) // 每 40ms 发出一个时钟滴答
.take(chars.length) // 字数喷完,流自动发送 onComplete 信号安全闭环
.map(index -> chars[index.intValue()]); // 极致简洁:直接返回字符即可,WebFlux 自动打包成标准 SSE 协议格式!
}
}


前端代码

在 Spring Boot 项目中,静态文件有其正宗的安家之处。 将下方的 index.html 直接放进项目的 src/main/resources/public/ 文件夹下。WebFlux 底层会自动完成无阻塞的零拷贝映射,完全不需要你手写任何一行读取代码!

1
2
3
4
5
6
7
8
9
10
src/main/resources
├── public
│ └── index.html 👈 [主页/HTML阵地] 干净、纯粹,直接访问 http://localhost:8080/
└── static
├── css
│ └── main.css 👈 [样式阵地]
├── js
│ └── chat.js 👈 [脚本阵地]
└── images
└── logo.png 👈 [多媒体阵地]

index.html

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
<!DOCTYPE html>
<html lang="zh-CN">
<head>
<meta charset="UTF-8">
<title>KooHub AI - WebFlux 舱</title>
<style>
body { font-family: 'Segoe UI', Tahoma, Geneva, Verdana, sans-serif; background: #eef2f3; padding: 40px; }
.chat-container { max-width: 750px; margin: 0 auto; background: white; padding: 30px; border-radius: 16px; box-shadow: 0 10px 30px rgba(0,0,0,0.08); }
h2 { color: #2c3e50; border-bottom: 2px solid #f1f2f6; padding-bottom: 15px; margin-top: 0; }
.input-box { display: flex; margin-bottom: 20px; }
input { flex: 1; padding: 14px; border: 1px solid #dcdde1; border-radius: 8px; font-size: 16px; outline: none; transition: border 0.3s; }
input:focus { border-color: #3498db; }
button { background: #3498db; color: white; border: none; padding: 0 28px; margin-left: 12px; border-radius: 8px; cursor: pointer; font-size: 16px; font-weight: bold; transition: background 0.3s; }
button:hover { background: #2980b9; }
.output-zone { background: #1e272e; color: #f5f6fa; padding: 25px; border-radius: 10px; min-height: 180px; font-size: 16px; line-height: 1.8; white-space: pre-wrap; word-break: break-all; font-family: 'Courier New', Courier, monospace; }
.cursor { display: inline-block; width: 3px; height: 18px; background: #3498db; animation: blink 1s infinite; vertical-align: middle; margin-left: 3px; }
@keyframes blink { 0%, 100% { opacity: 0; } 50% { opacity: 1; } }
</style>
</head>
<body>

<div class="chat-container">
<h2>🤖 KooHub AI - Spring WebFlux 生产级推流舱</h2>
<div class="input-box">
<input type="text" id="promptInput" value="请剖析 Spring WebFlux 的非阻塞核心组件原理" placeholder="请输入你的硬核问题...">
<button onclick="sendWebFluxStream()">发射问题</button>
</div>
<div class="output-zone" id="outputZone">等待提问中...</div>
</div>

<script>
async function sendWebFluxStream() {
const inputEl = document.getElementById('promptInput');
const outputEl = document.getElementById('outputZone');
const prompt = inputEl.value;

if (!prompt) return;

// 闪烁光标初始化
outputEl.innerHTML = '<span class="cursor"></span>';

try {
// 1 请求 Spring WebFlux 后端接口
const response = await fetch('/api/chat', {
method: 'POST',
body: prompt,
headers: { 'Content-Type': 'text/plain;charset=UTF-8' }
});

const reader = response.body.getReader();
const decoder = new TextDecoder('utf-8');
outputEl.innerHTML = "";

// 2 抠干非阻塞管道字节流
while (true) {
const { value, done } = await reader.read();
if (done) break;

const chunk = decoder.decode(value, { stream: true });

// 3 解析由 WebFlux 自动组装的标准 SSE 报文
const lines = chunk.split('\n');
for (const line of lines) {
if (line.startsWith('data:')) {
// 提取真正的字符内容
const actualChar = line.replace('data:', '');
outputEl.innerHTML += actualChar;
}
}
}
// 完美收尾打勾
outputEl.innerHTML += ' <b style="color:#2ecc71;">✓</b>';

} catch (error) {
outputEl.innerHTML = "发生致命错误: " + error.message;
}
}
</script>
</body>
</html>

如果你是从老旧的 SpringMVC (WAR 包部署) 转过来的,可能会问:“以前那种不能被外部直接访问、安全性极高的 WEB-INF/jsp/ 或者 WEB-INF/templates/ 目录呢?” 其实,在新世界的 WebFlux 架构中:

  • 物理上抹除了 WAR 包和 Servlet 容器:WebFlux 默认作为独立 JAR 包运行在 Netty 上,Netty 根本没有传统 Tomcat 那套特殊的 WEB-INF 安全沙箱隔离机制。
  • 全面倒向 Thymeleaf / Freemarker 等模板引擎:如果你不想让前端直接拉取裸 HTML 静态文件,而是希望后端动态渲染,通常会引入 spring-boot-starter-thymeleaf。此时,所有的动态网页模板会妥善安置在 /src/main/resources/templates/ 目录下。这个目录天然受到保护,前端直接敲 URL 是绝对访问不到的,必须通过 Controller 路由转发才能渲染输出。


WebFilter 实现全链路 TraceId

在传统的 Spring MVC 生态中,我们做全链路日志追踪(TraceId)通常会求助于组件 HandlerInterceptor 配合老朋友 Mapped Diagnostic Context(MDC,底层基于 ThreadLocal)。

但在 Spring WebFlux 的世界里,绝对不要在多线程交替执行的异步管道中直接使用传统的 MDC/ThreadLocal!因为一个响应式流在执行过程中,可能会高频发生线程切换(比如从 Netty 线程切到 boundedElastic 线程),这会导致传统的 ThreadLocal 发生严重的 TraceId 丢失、或者错乱污染的灾难级事故。

Project Reactor 官方提供的救世主方案是 Context API。它是一个绑定在响应式流下游的只读元数据上下文,会随着每一次 map/flatMap 异步长征一路向下游传送。下面我们直接手写一套基于 WebFilter 与 Reactor Context 的非阻塞全链路 TraceId 追踪系统。

定义 TraceContext

为了实现代码的极简和优雅,我们需要在 org.slf4j.MDC 无法发挥作用的地方,封装一个专为 Reactor 服务的上下文钥匙。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
import org.slf4j.MDC;
import reactor.core.publisher.Signal;
import java.util.UUID;
import java.util.function.Consumer;

public class TraceContext {

// 链路追踪的日志 Key
public static final String TRACE_ID_KEY = "traceId";

/**
* 让日志门面(SLF4J/MDC)与 Reactor 管道中的 Context 完成 “量子纠缠”。
* 利用 Reactor 的 doOnEach 操作符,当流中吐出信号(数据/异常/完成)时,动态向当前线程写入/清除 MDC。
*/
public static <T> Consumer<Signal<T>> logWithTrace(Consumer<T> logAction) {
return signal -> {
if (signal.isOnNext()) { // 仅对正常吐出数据的信号进行拦截
// 1 从当前的 Reactor 信号上下文里抠出 TraceId
String traceId = signal.getContextView().getOrDefault(TRACE_ID_KEY, "UNKNOWN");
try {
// 2 极其短暂地绑定到当前线程的 MDC
MDC.put(TRACE_ID_KEY, traceId);
// 3 触发真正的业务日志打印逻辑
logAction.accept(signal.get());
} finally {
// 4 穿过火线后立刻粉碎它,绝对不污染下一个和该线程擦肩而过的陌生请求
MDC.remove(TRACE_ID_KEY);
}
}
};
}

/**
* 极简生成器
*/
public static String generateId() {
return UUID.randomUUID().toString().replace("-", "").substring(0, 16);
}
}


编写 TraceFilter

接下来,我们在寻址执行(DispatcherHandler)之前,架设一个全局的 WebFilter。 它负责把从前端请求头里抓来的 TraceId(或者由系统自动生成的最新 TraceId)注入到整个 WebFlux 响应式流水线的最底部。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
import com.zdemo.scloud.model.TraceContext;
import org.springframework.core.annotation.Order;
import org.springframework.stereotype.Component;
import org.springframework.web.server.ServerWebExchange;
import org.springframework.web.server.WebFilter;
import org.springframework.web.server.WebFilterChain;
import reactor.core.publisher.Mono;

@Component
@Order(-1) // 🌟 级别拉满!确保全网最高优先级,在跨域和路由触发前就把 TraceId 打进去
public class TraceFilter implements WebFilter {

@Override
public Mono<Void> filter(ServerWebExchange exchange, WebFilterChain chain) {
// 1 提取:先看看前端请求头里有没有上游微服务发过来的 TraceId
String traceId = exchange.getRequest().getHeaders().getFirst("X-Trace-Id");
if (traceId == null || traceId.isEmpty()) {
traceId = TraceContext.generateId(); // 如果没有(比如浏览器直接访问),立刻生成一个
}

// 2 吐回:顺手把这个 TraceId 塞进 HTTP 响应头里,方便前端查盘报错调试
exchange.getResponse().getHeaders().add("X-Trace-Id", traceId);

// 3 核心杀招:向下游推进。利用 contextWrite 将 TraceId 打入整个响应式流的生命周期中
final String finalTraceId = traceId;
return chain.filter(exchange)
.contextWrite(context -> context.put(TraceContext.TRACE_ID_KEY, finalTraceId));
}
}


配置日志打印模板

创建或修改你的 logback-spring.xml。为了让日志里能自动带上我们的 TraceId,必须要在日志的 Layout 模板中,加上配置项 [%X{traceId}]

1
2
3
4
5
6
7
8
9
10
11
12
<?xml version="1.0" encoding="UTF-8"?>
<configuration>
<appender name="CONSOLE" class="ch.qos.logback.core.ConsoleAppender">
<encoder>
<pattern>%d{yyyy-MM-dd HH:mm:ss.SSS} [%thread] %highlight(%-5level) [%X{traceId}] %cyan(%logger{15}) - %msg%n</pattern>
</encoder>
</appender>

<root level="INFO">
<appender-ref ref="CONSOLE" />
</root>
</configuration>


业务控制层日志的打印

现在,我们把这套全链路追踪应用到之前写的那个 ChatGPT 控制层。 传统的 log.info(…) 无法读取到响应式 Context,我们需要换用我们刚刚编写的 TraceContext.logWithTrace 黄金管道过滤器。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
@Slf4j
@RestController
@CrossOrigin(origins = "*") // 允许跨域,方便前后端无缝联调
@RequestMapping("/api")
public class ChatStreamController {

/**
* ChatGPT 同款流式对话接口
* 直接返回 Flux<String>,WebFlux 自动将其转化为高效的异步 SSE 流式长连接
*/
@PostMapping(value = "/chat", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<String> chatWithAi(@RequestBody String userPrompt) {
// 传统的日志打印在 WebFlux 中无法感知 Context,输出的 traceId 为空
log.info("⚠️ [传统日志拦截] 收到前端盲发请求,内容:{}", userPrompt);

String aiAnswer = "你好!我是基于 Spring WebFlux 架构构建的现代 AI 助手。"
+ "在 WebFlux 的世界里,传统的组件如 DispatcherServlet 进化为了全新的异步响应式 DispatcherHandler。"
+ "当请求涌入时,Netty 的 Event Loop 线程只负责注册事件,随后立刻去接待下一个请求。"
+ "现在你看到的每一个字,都是顺着响应式管道无阻塞、带背压地实时喷射出来的。这就是高并发的魅力!";
String[] chars = aiAnswer.split("");
return Flux.interval(Duration.ofMillis(40))
.take(chars.length)
.map(index -> chars[index.intValue()])
// 🌟 利用 doOnEach 配合 logWithTrace,让每一行吐出的字符都在日志里拥有完美的身份编码!
.doOnEach(TraceContext.logWithTrace(word ->
log.info("🎯 [打字机动态吐字发射]: {}", word)
));
}
}

点击运行服务,打开浏览器点击“发射问题”,接着观察你的控制台。会观察到控制台打印:

1
2
3
4
5
6
7
2026-06-19 14:10:12.587 [reactor-http-nio-2] INFO  [] c.z.s.c.ChatStreamController - ⚠️ [传统日志拦截] 收到前端盲发请求,内容:请剖析 Spring WebFlux 的非阻塞核心组件原理
2026-06-19 14:10:12.647 [parallel-1] INFO [933445c78f904885] c.z.s.c.ChatStreamController - 🎯 [打字机动态吐字发射]: 你
2026-06-19 14:10:12.687 [parallel-1] INFO [933445c78f904885] c.z.s.c.ChatStreamController - 🎯 [打字机动态吐字发射]: 好
2026-06-19 14:10:12.727 [parallel-1] INFO [933445c78f904885] c.z.s.c.ChatStreamController - 🎯 [打字机动态吐字发射]: !
2026-06-19 14:10:12.766 [parallel-1] INFO [933445c78f904885] c.z.s.c.ChatStreamController - 🎯 [打字机动态吐字发射]: 我
2026-06-19 14:10:12.807 [parallel-1] INFO [933445c78f904885] c.z.s.c.ChatStreamController - 🎯 [打字机动态吐字发射]: 是
...


整合 R2DBC

在传统架构中,无论你把 Web 层优化得多么惊天地泣鬼神,只要底层数据库驱动用的是旧时代的 JDBC(如标准的 MyBatis 或 Hibernate),那么每次查库都会死死卡住一个 Netty 的物理线程。

R2DBC 的核心革命:它用响应式的协议重写了数据库通信协议(比如 MySQL 的文字/二进制流协议)。当 Java 向 MySQL 发出一条复杂的 SQL 之后,当前线程立刻释放,等到 MySQL 把数据算好、通过 TCP 发回通知时,R2DBC 驱动再用响应式流(Flux/Mono)把数据顺着管道一路喷射给最前端的浏览器。真正实现 “全链路零阻塞”。

下面我们基于 Spring Boot 3.x + Spring Data R2DBC + MySQL(响应式版),演示包含单表 CRUD 与多表关联的生产最佳实践。


引入依赖

1
2
3
4
5
6
7
8
9
10
11
12
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-r2dbc</artifactId>
<version>3.3.0</version>
</dependency>
<dependency>
<groupId>io.asyncer</groupId>
<artifactId>r2dbc-mysql</artifactId>
<version>1.1.3</version>
</dependency>
</dependencies>


配置文件

在 application.yml 中,数据库连接协议全面升级为 r2dbc 开头:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
server:
port: 8080

spring:
application:
name: zdemo-webflux
r2dbc:
url: r2dbc:mysql://192.168.1.251:3306/zdemo?useSSL=false&characterEncoding=UTF-8
username: root
password: 123456
pool:
enabled: true # 默认为 true
initial-size: 10 # 生产建议:初始化 10 个物理长连接
max-size: 30 # 生产建议:响应式下 30 个连接就足够榨干几千 QPS
max-idle-time: 30m # 生产建议:空闲连接半小时自动释放,防止卡死 MySQL 进程
max-life-time: 1h # 生产建议:一个连接最长存活 1 小时,到期自动重连,免疫网络抖动

响应式连接池是弹性复用的,一定要加 r2dbc-pool 相关的配置,这样单台机器用极小的物理开销就能扛住惊人的并发量。在传统基于 Tomcat 的多线程 Servlet 生态(Spring MVC + Druid/HikariCP)里,连接池的 max-size 动辄配到 200 甚至 500。 你可能疑惑为什么在新世界的 WebFlux + R2DBC 里,只建议你配到 30 呢?原因是:

  • 传统 JDBC(HikariCP):它的模式是 “一个连接奴役一个线程”。如果并发来了 100 个请求,在查库期间,这 100 个线程必须死死咬住 100 个物理数据库连接在原地坐牢。连接数不够,线程就得排队等死。
  • 响应式 R2DBC(r2dbc-pool):它的模式是 “长连接管道化多路复用”。当 1000 个请求涌入时,Netty 通过非阻塞的方式把 1000 条 SQL 发给 R2DBC 连接池。由于 R2DBC 内部是通过事件驱动流式返回结果的,一条物理连接可以在 1 毫秒内帮 A 请求发完 SQL,转头零切换开销立刻去帮 B 请求发 SQL,然后等 MySQL 吐出数据时再异步分发给 A 和 B。

在生产测试中,R2DBC 的 20 个连接所能释放出来的吞吐量,足以打爆传统 JDBC 200 个连接的并发极限。把连接数配得太大,反而会加重 MySQL 服务端本身的内存与线程上下文切换负担。


仓储层最佳实践

库表和实体

首先是库表的准备:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
-- 1. 创建用户表 (t_user)
CREATE TABLE IF NOT EXISTS `t_user` (
`id` BIGINT NOT NULL AUTO_INCREMENT COMMENT '主键ID',
`username` VARCHAR(50) NOT NULL COMMENT '用户名',
`email` VARCHAR(100) NOT NULL COMMENT '电子邮箱',
PRIMARY KEY (`id`),
UNIQUE KEY `uk_username` (`username`) COMMENT '用户名唯一索引'
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_general_ci COMMENT='用户表';

INSERT INTO `t_user` (`id`, `username`, `email`) VALUES
(1, 'koohub_admin', 'admin@koohub.com'),
(2, 'zhangqingli', 'zhang@owlias.com'),
(3, 'reactor_geek', 'geek@flux.io');

-- 2. 创建文章表 (t_article)
CREATE TABLE IF NOT EXISTS `t_article` (
`id` BIGINT NOT NULL AUTO_INCREMENT COMMENT '主键ID',
`user_id` BIGINT NOT NULL COMMENT '作者用户ID',
`title` VARCHAR(150) NOT NULL COMMENT '文章标题',
`content` TEXT NOT NULL COMMENT '文章正文',
PRIMARY KEY (`id`),
KEY `idx_user_id` (`user_id`) COMMENT '作者ID普通索引'
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_general_ci COMMENT='文章表';

INSERT INTO `t_article` (`user_id`, `title`, `content`) VALUES
(2, 'Owlias v1.3 架构升级演进史', '在 Owlias v1.3 中,我们全面重构了前端交互,引入了主题动态切换、实时流式搜索以及极为硬核的关键词高亮过滤。今天我们来拆解它的底层比例调优...'),
(2, 'Hexo 与 Markdown 动静分离的完美工程实践', '如何让你的静态博客拥有动态大厂的工程美感?答案是利用原生的 HTML 标签直接内嵌 Tailwind 流水线组件,适配手机端与电脑端的弹性切换...'),
(1, 'Spring Boot 3.x 响应式 R2DBC 性能压测报告', '在 30 个物理连接的限制下,基于 R2DBC 的响应式架构吞吐量直接打爆了传统 JDBC 200 个连接的并发极限,上下文切换开销降低了 70%...'),
(3, 'Project Reactor 核心背压机制(Backpressure)源码剖析', '当下游的 request(n) 信号逆流而上时,底层的物理驱动是如何做到精准控流的?本文将带你一行行看懂 BaseSubscriber 的核心设计...'),
(2, '从 Hume 到 Kant:因果律在响应式编程中的哲学启示', 'David Hume 认为因果律只是习惯性的联想,而 Kant 在《纯粹理性批判》中将其升格为先验范畴。这正如同响应式中的声明式管道:图纸早已存在,唯有订阅那一刻,现实才真正发生...'),
(1, '手把手教你用 WebFilter 降维打击全链路日志 TraceId', 'ThreadLocal 在多线程交替的 EventLoop 中已经寿终正寝。唯有将元数据注入 Reactor Context,配合信号过滤器的动态擦除,才是唯一的生产解...'),
(3, 'WebFlux 的前哨站:HttpHandler 与 DispatcherHandler 执行内幕', '当 Netty 抓到 TCP 字节块时,HttpHandler 会瞬间抹平容器差异,将其升格转换为标准的 ServerWebExchange 丢给总调度暴风眼...'),
(2, 'MySQL 零拷贝技术在高性能文件推流中的应用', '利用 routes.file() 一行代码直接调用操作系统的 NIO 异步通道,不经过 JVM 内存中转,直接把 index.html 喷射向远端网卡...'),
(1, '微服务解耦:Service 层内存动态组装 vs 底层大 SQL 选型指南', '在分库分表的未来演进中,方案 A 的 flatMap 内存缝合虽然多了几行代码,却能带给你坚不可摧的架构解耦能力。方案 B 则在单体合并时拥有极致的性能...'),
(3, 'Reactive Stream TCK 兼容性测试完全通关指南', '如何证明你写的响应式驱动是真正的非阻塞?你必须通过符合 Reactive Streams 规范的全部边界条件测试,包括空流、异常流以及疯狂压测下的背压稳健性...');

其次是 Entity 和 DTO:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
import org.springframework.data.annotation.Id;
import org.springframework.data.relational.core.mapping.Table;

@Table("t_user")
@Data
@AllArgsConstructor
@NoArgsConstructor
public class User {
@Id
private Long id;
private String username;
private String email;
}

@Table("t_article")
@Data
public class Article {
@Id
private Long id;
private Long userId;
private String title;
private String content;
}

@Data
public class ArticleDetailsDTO {
private Long articleId;
private String title;
private String content;
private Long userId;
private String authorName; // 来自关联的 t_user 表
}


仓储层 Repository

UserRepository

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
import com.zdemo.scloud.model.User;
import org.springframework.data.domain.Pageable;
import org.springframework.data.r2dbc.repository.Query;
import org.springframework.data.repository.reactive.ReactiveCrudRepository;
import org.springframework.stereotype.Repository;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;

@Repository
public interface UserRepository extends ReactiveCrudRepository<User, Long> {

/**
* spring data 衍生查询:
* 自动解析为 SELECT * FROM t_user WHERE username = ?
*/
Mono<User> findByUsername(String username);

/**
* 分页方案 A:通过方法名自动衍生分页查询
* @param pageable 包含 pageNumber(页码) 和 pageSize(每页条数) 以及 Sort(排序) 的对象
* @return 返回当前页的响应式数据流
*/
Flux<User> findAllBy(Pageable pageable);

/**
* 分页方案 B:手写纯正原生 SQL 配合 Pageable
*/
@Query("SELECT * FROM t_user WHERE email LIKE '%.com'")
Flux<User> findAllBy2(Pageable pageable);

/**
* 游标过滤分页
*/
@Query("SELECT * FROM t_user WHERE id > :lastId ORDER BY id ASC LIMIT :pageSize")
Flux<User> streamNextPage(Long lastId, int pageSize);
}

ArticleRepository

1
2
3
public interface ArticleRepository extends ReactiveCrudRepository<Article, Long> {
// 基础的单表增删改查由 ReactiveCrudRepository 自动提供
}


业务层 service

在传统 MVC 里,执行完 Service 意味着数据库的操作已经尘埃落定;但在 WebFlux 生产标准下,执行完 Service 只代表“业务流水线的图纸已经设计完毕”。真正的查库、逻辑触发、网络传输,直到 Controller 把这个 Mono/Flux 抛给底层的 Netty 容器并引发订阅时,才会真正非阻塞地爆发执行。

关于响应式开发中,我们越来越倾向于废弃传统的 Page 对象。因为既然全链路都打通了非阻塞的 Flux 流,我们更倾向于直接利用流的 take() 操作符或者前端的背压(request(n) 背压信号)来实现分页。现代的大数据大屏、博客的文章列表、或者 ChatGPT 的对话历史,最标准的做法是“无限下滚加载(Infinite Scroll)。如果一定要分页,实际上是分页退化成了纯粹为了配合前端 UI 视觉展现而做的“流量闸门控制”。

UserService 和 UserServiceImpl

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
public interface UserService {
Mono<User> createUser(User user);
Mono<User> register(User user); // 包含重复校验的“增”
Mono<User> findById(Long id); // “查单条”
Flux<User> findAllStream(); // “查流式列表”
Mono<Page<User>> getUserPage(Pageable pageable); // 分页查询
Mono<User> update(Long id, User user); // “改”
Mono<Void> delete(Long id); // “删”
}

@Service
public class UserServiceImpl implements UserService {

@Resource
private UserRepository userRepository;

@Override
@Transactional
public Mono<User> createUser(User user) {
return userRepository.save(user);// 传 ID 则更新,不传则新增
}

@Override
@Transactional // 响应式环境下依然使用声明式事务
public Mono<User> register(User user) {
return userRepository.findByUsername(user.getUsername())
// 如果找到了,利用 Mono.error 抛出响应式异常,中断管道并回滚
.flatMap(existing -> Mono.<User>error(new IllegalArgumentException("用户名已存在!")))
// 如果没找到,switchIfEmpty 才会触发真正的安全落库
.switchIfEmpty(Mono.defer(() -> userRepository.save(user)));
}

@Override
public Mono<User> findById(Long id) {
return userRepository.findById(id);
}

@Override
public Flux<User> findAllStream() {
return userRepository.findAll();
}

/**
* 非阻塞分页查询:在新世界,有了 Flux 的流式传输与弹性背压(Backpressure)防护,
* 内存已经金刚不坏。分页退化成了纯粹为了配合前端 UI 视觉展现而做的 “流量闸门控制”。
*/
public Mono<Page<User>> getUserPage(Pageable pageable) {
// 1 异步编排任务 A:去查当前页的数据列表,收拢为 List
Mono<List<User>> contentMono = userRepository.findAllBy(pageable) // 需要在 Repository 定义该方法返回 Flux<User>
.collectList(); // 将 Flux 流在内存中聚拢为单个 Mono<List>

// 2 异步编排任务 B:去查总记录数
Mono<Long> countMono = userRepository.count();

// 3 核心合流:利用 Mono.zip 类似多线程并发,让 MySQL 同时去算 COUNT 和 LIMIT
// 两个都算完后,非阻塞地组装成标准的 PageImpl 对象返回
return Mono.zip(contentMono, countMono)
.map(tuple -> new PageImpl<>(tuple.getT1(), pageable, tuple.getT2()));
}

@Override
@Transactional
public Mono<User> update(Long id, User user) {
return userRepository.findById(id)
.flatMap(existingUser -> {
existingUser.setUsername(user.getUsername());
existingUser.setEmail(user.getEmail());
return userRepository.save(existingUser);
});
}

@Override
@Transactional
public Mono<Void> delete(Long id) {
return userRepository.deleteById(id);
}
}

ArticleService 和 ArticleServiceImpl

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
public interface ArticleService {
Mono<Article> saveArticle(Article article);
// 多表联合方案 A:代码层编排组合
Mono<ArticleDetailsDTO> getArticleDetailsByMemory(Long articleId);
// 多表联合方案 B:DatabaseClient 跑原生大 SQL 映射
Flux<ArticleDetailsDTO> getArticleDetailsStreamBySql();
}

@Service
public class ArticleServiceImpl implements ArticleService {

@Resource
private ArticleRepository articleRepository;

@Resource
private UserRepository userRepository;

@Resource
private DatabaseClient databaseClient; // 注入用于执行大 SQL 的客户端

@Override
public Mono<Article> saveArticle(Article article) {
return articleRepository.save(article);
}

/**
* 🌟 方案 A 实现:纯代码层响应式操作符动态缝合(无 SQL 耦合,微服务标准)
*/
@Override
public Mono<ArticleDetailsDTO> getArticleDetailsByMemory(Long articleId) {
return articleRepository.findById(articleId)
// 嵌套调用 UserRepository 查出对应的作者
.flatMap(article -> userRepository.findById(article.getUserId())
.map(user -> {
// 将两层单表的数据拼装进 DTO 吐给上层
ArticleDetailsDTO dto = new ArticleDetailsDTO();
dto.setArticleId(article.getId());
dto.setTitle(article.getTitle());
dto.setContent(article.getContent());
dto.setUserId(user.getId());
dto.setAuthorName(user.getUsername());
return dto;
})
);
}

/**
* 🌟 方案 B 实现:DatabaseClient 跑多表 LEFT JOIN 原生大 SQL(极致性能)
*/
@Override
public Flux<ArticleDetailsDTO> getArticleDetailsStreamBySql() {
String complexSql = "SELECT a.id AS aid, a.title, a.content, u.id AS uid, u.username " +
"FROM t_article a " +
"LEFT JOIN t_user u ON a.user_id = u.id";

return databaseClient.sql(complexSql)
.map((row, rowMetadata) -> {
// 非阻塞映射回调,手动抽干 ResultRow 组装 DTO
ArticleDetailsDTO dto = new ArticleDetailsDTO();
dto.setArticleId(row.get("aid", Long.class));
dto.setTitle(row.get("title", String.class));
dto.setContent(row.get("content", String.class));
dto.setUserId(row.get("uid", Long.class));
dto.setAuthorName(row.get("username", String.class));
return dto;
})
.all(); // 喷射出去变成 Flux<ArticleDetailsDTO>
}
}


控制层

UserReactiveController

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
@RestController
@RequestMapping("/api/user")
public class UserReactiveController {

@Resource
private UserService userService;

// 🟢 1.1 增
@PostMapping
public Mono<User> createUser(@RequestBody User user) {
return userService.createUser(user);
}

// 🟢 1.2 增
@PostMapping("/register")
@ResponseStatus(HttpStatus.CREATED)
public Mono<User> registerUser(@RequestBody User user) {
return userService.register(user);
}

// 🔵 2. 查(单条与流式列表)
@GetMapping("/{id}")
public Mono<ResponseEntity<User>> getUser(@PathVariable Long id) {
return userService.findById(id)
.map(ResponseEntity::ok)
.defaultIfEmpty(ResponseEntity.notFound().build());
}

@GetMapping(produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<User> streamAllUsers() {
// 降维打击:MySQL 里的数据会顺着网络驱动,以流的形式持续、非阻塞地滑向前端
return userService.findAllStream();
}

// 🟡 3. 改
@PutMapping("/{id}")
public Mono<User> updateUser(@PathVariable Long id, @RequestBody User user) {
return userService.update(id, user);
}

// 🔴 4. 删
@DeleteMapping("/{id}")
public Mono<Void> deleteUser(@PathVariable Long id) {
return userService.delete(id);
}
}

ArticleReactiveController

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
@RestController
@RequestMapping("/api/article")
public class ArticleReactiveController {

@Resource
private ArticleService articleService;

@PostMapping
public Mono<Article> createArticle(@RequestBody Article article) {
return articleService.saveArticle(article);
}

// 测试多表联合方案 A
@GetMapping("/memory/{id}")
public Mono<ArticleDetailsDTO> getDetailsByMemory(@PathVariable Long id) {
return articleService.getArticleDetailsByMemory(id);
}

// 测试多表联合方案 B:以标准打字机 SSE 形式实时将大联表数据喷向前端
@GetMapping(value = "/stream-sql", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<ArticleDetailsDTO> getDetailsStreamBySql() {
return articleService.getArticleDetailsStreamBySql();
}
}


加入 AOP 逻辑

响应式的 Service 还需要之前切面吗?答案是极度需要!但实现逻辑发生了翻天覆地的变化。如果你直接用传统的 Spring AOP(如 @Around 环绕通知)去拦截一个响应式 Service 方法,并在切面里写 System.currentTimeMillis() 来统计耗时,你会得到一个极其荒谬的结论:每个复杂的查库方法耗时都是 0 毫秒。

为什么传统切面在响应式中会“失效”呢?这是因为响应式 Service 方法在被调用时,它根本没有真正去查数据库,它只是秒级组装了一条 Mono 管道并把它返回给了 Controller。真正的数据库 I/O 动作,要等到最后前端浏览器发起订阅(Subscribe)时才会真正触发。 传统的 AOP 只能拦截到“管道组装完毕”的那一瞬间,根本拦截不到“数据在管道里流动”的真实业务阶段。

要在 WebFlux 中做切面(如业务日志、权限检查、耗时统计),我们必须顺应响应式的游戏规则,利用管道的操作符(如 doOnSuccess、defer)来织入切面逻辑。


引入依赖

你只需要在依赖管理中加入 spring-boot-starter-aop,它会自动为你拉取 aspectjweaver 等全套核心切面组件。

1
2
3
4
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-aop</artifactId>
</dependency>


切面示例

如果你希望保留原汁原味的 @Aspect 注解,你必须在切面里对返回的 Mono/Flux 进行二次代理和编排:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
import lombok.extern.slf4j.Slf4j;
import org.aspectj.lang.ProceedingJoinPoint;
import org.aspectj.lang.annotation.Around;
import org.aspectj.lang.annotation.Aspect;
import org.springframework.stereotype.Component;
import reactor.core.publisher.Mono;

/**
* 响应式编程另一种AOP的优雅实现:利用响应式自带的 transform 操作符
*
* 由于 WebFlux 的响应式流天生自带极强的扩展性,架构师更倾向于抛弃外部的 AOP 框架,
* 直接在 Service 的管道尾端使用通用的响应式组件来完成切面功能。
*/
@Slf4j
@Aspect
@Component
public class ReactiveLogAspect {

@Around("execution(* com.zdemo.scloud.service..*.*(..))")
public Object profile(ProceedingJoinPoint point) throws Throwable {
long startTime = System.currentTimeMillis();

// 1 执行原方法,拿到的其实是一个 Mono 管道
Object result = point.proceed();

if (result instanceof Mono) {
// 2 核心杀招:不能直接算时间,必须挂载到 Mono 的生命周期钩子(doOnSuccess/doOnError)里
return ((Mono<?>) result)
.doOnSuccess(data -> {
long endTime = System.currentTimeMillis();
log.info("[AOP 耗时统计] Service 执行成功,真实非阻塞耗时: {} ms", endTime - startTime);
})
.doOnError(err -> {
long endTime = System.currentTimeMillis();
log.error("[AOP 异常捕获] Service 发生暴雷,耗时: {} ms", endTime - startTime);
});
}
return result;
}
}


更优雅的实现

很多走纯粹响应式路线的架构师,在 WebFlux 项目中会故意不引入 spring-boot-starter-aop 这个外部依赖。因为引入它意味着引入了传统的 AspectJ 动态代理字节码机制。如果你不想引入这个胖依赖,Reactor 本身就内置了一个能够平替 AOP 的高阶操作符——.transform()。你可以写一个纯粹的函数式拦截器,在 Service 的流中直接调用:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
import reactor.core.publisher.Mono;
import java.util.function.Function;

public class ReactiveTimer {
/**
* 原生切面代理:无需 ProceedingJoinPoint 依赖
*/
public static <T> Function<Mono<T>, Mono<T>> logTime(String methodName) {
return idMono -> Mono.defer(() -> {
long startTime = System.currentTimeMillis();
return idMono
.doOnSuccess(data -> {
long endTime = System.currentTimeMillis();
System.out.println("✨ [原生流控耗时] " + methodName + " 耗时: " + (endTime - startTime) + "ms");
});
});
}
}

在 Service 层使用时,连注解都不用写,直接挂在管道末端,用起来极其丝滑:

1
2
3
4
5
public Mono<User> getUserById(Long id) {
return userRepository.findById(id)
// 利用 transform 动态织入无依赖的响应式切面
.transform(ReactiveTimer.logTime("getUserById"));
}


蜂鸟空间项目

关于更多的响应式编程项目,可以参考最近写的一个小项目 《蜂鸟空间》。主要的业务涉及博客系统的认证授权、发布,以及实时评论等功能。主要的技术栈:

  • 核心框架: Spring Boot 3.x / WebFlux (响应式编程)
  • 实时通信: Netty + WebSocket
  • 数据存储: MySQL + Redis (缓存与消息总线)
  • 身份认证: JWT (JSON Web Token)
  • 通信协议: Redis Pub/Sub (集群消息总线)