1. 什么是响应式编程?
1.1 核心定义
响应式编程(Reactive Programming)是一种基于异步数据流、数据变化自动推送、自动响应的声明式编程范式。
1.2 通俗举例
Excel 公式:
=A1\+B1,修改 A1/B1 的值,结果单元格自动更新。核心逻辑:数据变动 → 自动推送 → 自动更新。
1.3 命令式 VS 响应式
| 编程范式 | 核心特点 | 执行逻辑 |
|---|---|---|
| 命令式 | 主动拉取、阻塞执行 | a = b \+ c,b 变更后 a 不会自动更新 |
| 响应式 | 被动推送、非阻塞 | a 永久依赖 b+c,b 变更自动触发 a 更新 |
1.\4 和AKKA的比较
- 做Spring 生态的异步 Web / 微服务、数据管道 → 选 Reactor。
- 做高可用、分布式、强容错的复杂系统 → 选 Akka。
2. 响应式编程四大核心要素
-
数据流(Stream):一切皆流,网络请求、消息、定时任务、用户操作都属于数据流。
-
发布订阅模型:生产者发布数据,消费者监听并处理数据。
-
异步非阻塞:不占用线程阻塞等待,适配高并发、高吞吐场景。
-
声明式链式编程:仅定义处理规则,无需手动管控执行流程。
3. Java 主流响应式框架
3.1 RxJava
-
第三方开源响应式工具库
-
核心对象:
Observable(无限数据流) -
常用场景:Android 开发、中间件、通用异步工具
3.2 Project Reactor(重点)
-
Spring 官方自研响应式库,Spring WebFlux 底层依赖
-
核心两大对象:
-
Mono:0~1 个元素,用于单个接口返回、单个对象查询
-
Flux:0~N 个元素,用于流式列表、持续消息推送
-
3.3 Mono 是什么?(超通俗比喻)
现实场景:
你去餐厅吃饭
- 你点单(调用方法)
- 服务员给你一张小票
- 你不用站在厨房等
- 厨房做好了,叫号取餐
对应代码:
- 点单 = 调用 dao.findById(id)
- 小票 = Mono
- 不用等 = 非阻塞、线程释放
- 做好通知 = 回调、IO 完成
结论:
Mono 不是数据,不是线程,不是执行。 Mono = 未来数据的凭证。
3.4 Mono 的核心作用(3 条就够)
1. 告诉 Reactor:这是异步 IO 操作
只要你看到返回 Mono
Reactor 就知道:
- 后面要等数据
- 现在可以释放线程
- 等好了再通知我
这就是 非阻塞的关键。
2. 封装“成功、失败、完成”三种信号
Mono 内部只会发送 3 种信号:
- onNext(数据) → 成功
- onError(异常) → 失败
- onComplete() → 结束 Reactor 只认这 3 种信号。
3. 让你可以链式编排逻辑(声明式)
mono
.doOnNext(日志)
.map(转换)
.switchIfEmpty(兜底)
.timeout(超时)
4. 总结
调用方法返回 Mono 时,代码根本没执行!
Mono<Order> mono = orderDao.findById("123");
// 执行到这里,根本没查库!!!
4. 实战示例一:RxJava3 基础案例
4.1 Maven 依赖
<dependency>
<groupId>io.reactivex.rxjava3</groupId>
<artifactId>rxjava</artifactId>
<version>3.1.8</version>
</dependency>
4.2 可运行代码
import io.reactivex.rxjava3.core.Observable;
import io.reactivex.rxjava3.schedulers.Schedulers;
public class RxJavaDemo {
public static void main(String[] args) throws InterruptedException {
// 1. 创建数据流(发布者)
Observable<String> dataStream = Observable.create(emitter -> {
for (int i = 1; i <= 5; i++) {
Thread.sleep(500);
emitter.onNext("数据->" + i);
}
emitter.onComplete();
});
// 2. 流式处理 + 订阅消费
dataStream
.subscribeOn(Schedulers.newThread()) // 异步线程执行
.filter(data -> !data.equals("数据->3")) // 过滤指定数据
.map(data -> data + "【处理完成】") // 数据转换
.subscribe(
data -> System.out.println("接收:" + data), // 正常消费
err -> System.err.println("异常:" + err), // 异常捕获
() -> System.out.println("数据流结束") // 流完成回调
);
Thread.sleep(3000);
}
}
4.3 运行结果
接收:数据->1【处理完成】
接收:数据->2【处理完成】
接收:数据->4【处理完成】
接收:数据->5【处理完成】
数据流结束
5. 实战示例二:Project Reactor 案例
5.1 Maven 依赖
<dependency>
<groupId>io.projectreactor</groupId>
<artifactId>reactor-core</artifactId>
<version>3.6.5</version>
</dependency>
5.2 可运行代码
import reactor.core.publisher.Flux;
import reactor.core.scheduler.Schedulers;
public class ReactorDemo {
public static void main(String[] args) throws InterruptedException {
// 1. 创建流式发布者
Flux<Integer> dataStream = Flux.create(emitter -> {
for (int i = 1; i <= 5; i++) {
Thread.sleep(500);
System.out.println("生产数据:" + i);
emitter.next(i);
}
emitter.complete();
});
// 2. 链式处理 + 异步订阅
dataStream
.subscribeOn(Schedulers.boundedElastic()) // 异步线程池
.filter(num -> num != 3) // 过滤数据
.map(num -> num * 10) // 数据转换
.doOnNext(d -> System.out.println("中间处理:" + d))
.subscribe(
data -> System.out.println("最终消费:" + data),
err -> System.err.println("异常:" + err),
() -> System.out.println("数据流结束")
);
Thread.sleep(3500);
System.out.println("主线程执行完毕");
}
}
5.3 运行结果
生产数据:1
中间处理:10
最终消费:10
生产数据:2
中间处理:20
最终消费:20
生产数据:3
生产数据:4
中间处理:40
最终消费:40
生产数据:5
中间处理:50
最终消费:50
数据流结束
主线程执行完毕
5.4 map和flatMap的差异
map函数实现普通的同步计算。
flatMap会异步订阅内部其他publisher并返回相关结果进行计算。
6. 实战示例三:高并发 API 接口(网关 / 商品 / 订单)
什么时候用?
- 接口 QPS 高(1000+)
- 调用第三方服务、数据库、Redis
- 不想服务器被大量线程卡死 传统 servlet 问题(同步阻塞) 1000 个请求 → 开 1000 个线程 线程等待 IO 时,占着坑不干活 → 服务器卡死 Reactor 优势(非阻塞) 1000 个请求 → 2~4 个核心线程 等待 IO 时线程去处理别的请求 → 吞吐量提升 3~10 倍 controller层代码:
@RestController
public class OrderController {
@GetMapping("/order/{id}")
public Mono<Order> getOrder(@PathVariable String id) {
// Mono:0/1 个结果 → 非阻塞查询 DB/Redis
return orderService.findById(id)
.doOnNext(order -> log.info("查询订单:{}", id));
}
@GetMapping("/orders")
public Flux<Order> listOrders() {
// Flux:多个结果 → 流式返回
return orderService.findAll();
}
}
service层代码:
import org.springframework.stereotype.Service;
import reactor.core.publisher.Mono;
@Service
public class OrderServiceImpl implements OrderService {
// 注入响应式 DAO(ReactiveMongo/ReactiveRedis/R2DBC)
private final OrderReactiveDao orderDao;
// 构造器注入(Spring推荐)
public OrderServiceImpl(OrderReactiveDao orderDao) {
this.orderDao = orderDao;
}
/**
* 响应式根据ID查询订单
* 重点:返回 Mono<Order>,非阻塞!
*/
@Override
public Mono<Order> findById(String orderId) {
// 非阻塞查询数据库/Redis
return orderDao.findById(orderId)
// 数据到达后自动执行(不阻塞线程)
.doOnNext(order -> System.out.println("查询到订单:" + orderId))
// 空值处理
.switchIfEmpty(Mono.error(new RuntimeException("订单不存在:" + orderId)));
}
}
DAO层代码:
import org.springframework.data.r2dbc.repository.Query;
import org.springframework.data.r2dbc.repository.R2dbcRepository;
import org.springframework.stereotype.Repository;
import reactor.core.publisher.Mono;
@Repository
public interface OrderDao extends R2dbcRepository<Order, String> {
/**
* 根据订单ID查询订单(响应式)
*/
@Query("SELECT * FROM t_order WHERE order_id = :orderId")
Mono<Order> findByOrderId(String orderId);
}
7. 实战示例四:批量调用 N 个第三方接口(并发聚合)
业务需求 查询用户详情,需要同时调用:
- 用户服务
- 订单服务
- 优惠券服务
- 积分服务 传统写法 串行调用:总耗时 = A + B + C + D Reactor 写法 并行非阻塞调用:总耗时 = 耗时最长的那个
public Mono<UserInfo> getUserInfo(String userId) {
// 同时调用 4 个服务,非阻塞并行执行
Mono<User> userMono = userService.getUser(userId);
Mono<List<Order>> orderMono = orderService.getOrders(userId);
Mono<Coupon> couponMono = couponService.getCoupon(userId);
Mono<Points> pointsMono = pointsService.getPoints(userId);
// 并行聚合结果
return Mono.zip(userMono, orderMono, couponMono, pointsMono)
.map(tuple -> {
// 组装返回
return new UserInfo(tuple.getT1(), tuple.getT2(), tuple.getT3(), tuple.getT4());
});
}
8. 实战示例四:实时数据流推送(消息 / 监控 / 日志 / IoT)
业务需求
- 实时日志推送
- 设备上报数据
- 聊天消息、通知流
- 大屏实时数据 Reactor 天然支持 持续数据流、自动推送、背压保护、不 OOM
@GetMapping(value = "/monitor/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<ServerData> monitorData() {
// 持续推送实时数据(SSE 前端实时接收)
return Flux.interval(Duration.ofSeconds(1))
.map(tick -> {
ServerData data = new ServerData();
data.setCpu(monitorService.getCpuUsage());
data.setMem(monitorService.getMemUsage());
return data;
});
}
9. 核心总结
9.1 为什么叫响应式?
-
数据主动推送,无需循环轮询拉取数据
-
数据变更自动触发下游业务逻辑
-
开发者仅定义规则,无需手动管控执行流程
9.2 响应式编程优势
-
非阻塞异步:线程资源利用率高,吞吐量优异
-
链式编码简洁:规避回调地狱,代码可读性高
-
内置背压机制:防止生产过快导致内存溢出OOM
-
统一异常处理:无需多处编写 try-catch
9.3 生产适用场景
-
高并发网关、微服务(Spring Cloud + WebFlux)
-
实时消息、IM 聊天、消息推送服务
-
大数据日志、IoT 设备数据上报
-
前端异步事件处理(RxJS)
9.4 和servlet的比较
- Servlet 模型:一个请求霸占一个线程,从生到死,等 IO 时啥也不干,纯浪费。
请求 1 → 线程 1(独占)→ 查询 DB(等待,线程空转)→ 结果返回 → 线程 1 释放
请求 2 → 线程 2(独占)→ 查询 DB(等待,线程空转)→ 结果返回 → 线程 2 释放
...
请求 200 → 线程 200(独占)
请求 201 → 排队等待
- Reactor 模型:少量线程无限复用,遇到等待就放手,IO 回来再继续,线程永不空闲。
请求 1 → 线程 1 处理 → 遇到 DB 查询 → 线程 1 释放(去处理请求 3)
请求 2 → 线程 2 处理 → 遇到 Redis 查询 → 线程 2 释放(去处理请求 4)
请求 3 → 线程 1 处理 ...
DB 返回结果 → 空闲线程(可能是任意一个)继续执行请求 1
10. 终极通俗总结
命令式编程:我主动去拿数据。
响应式编程:数据好了主动推给我。
响应式 = 异步数据流 + 自动推送 + 链式处理 + 非阻塞高并发
参考文章
1、《reactor官方文档》 2、《reactor帮助文档》