Java 响应式编程:Reactor 框架深度解析
核心概念
响应式编程是一种编程范式,关注数据的异步流和变化传播。在 Java 中,Reactor 框架提供了强大的响应式编程支持,基于 Reactive Streams 规范实现。
Reactor 核心组件
- Mono:表示 0 或 1 个元素的异步序列
- Flux:表示 0 到 N 个元素的异步序列
- Scheduler:调度器,控制任务执行线程
- Operator:操作符,用于转换和处理数据流
创建流
// 创建 Flux
Flux<String> flux = Flux.just("Hello", "World", "Reactor");
// 创建 Mono
Mono<String> mono = Mono.just("Hello");
// 从集合创建
List<String> list = Arrays.asList("A", "B", "C");
Flux<String> fromList = Flux.fromIterable(list);
// 从数组创建
String[] array = {"X", "Y", "Z"};
Flux<String> fromArray = Flux.fromArray(array);
// 创建范围
Flux<Integer> range = Flux.range(1, 10);
// 创建空流
Flux<String> empty = Flux.empty();
Mono<String> monoEmpty = Mono.empty();
// 创建错误流
Flux<String> error = Flux.error(new RuntimeException("Something went wrong"));
订阅流
// 订阅并消费数据
Flux.just("A", "B", "C")
.subscribe(
item -> System.out.println("Received: " + item),
error -> System.err.println("Error: " + error.getMessage()),
() -> System.out.println("Completed")
);
// 简化订阅
Flux.just(1, 2, 3)
.subscribe(System.out::println);
// 使用 Disposable 控制订阅
Disposable disposable = Flux.interval(Duration.ofSeconds(1))
.subscribe(tick -> System.out.println("Tick: " + tick));
// 5 秒后取消订阅
Thread.sleep(5000);
disposable.dispose();
操作符
// 转换操作符
Flux.just("hello", "world")
.map(String::toUpperCase)
.subscribe(System.out::println); // HELLO, WORLD
// 过滤操作符
Flux.range(1, 10)
.filter(n -> n % 2 == 0)
.subscribe(System.out::println); // 2, 4, 6, 8, 10
// 映射操作符
Flux.just("a", "b", "c")
.flatMap(s -> Flux.just(s.toUpperCase(), s.toLowerCase()))
.subscribe(System.out::println); // A, a, B, b, C, c
// 组合操作符
Flux<String> flux1 = Flux.just("A", "B");
Flux<String> flux2 = Flux.just("X", "Y");
Flux.concat(flux1, flux2)
.subscribe(System.out::println); // A, B, X, Y
Flux.zip(flux1, flux2, (a, b) -> a + "-" + b)
.subscribe(System.out::println); // A-X, B-Y
// 聚合操作符
Flux.range(1, 10)
.reduce(0, Integer::sum)
.subscribe(sum -> System.out.println("Sum: " + sum)); // Sum: 55
Flux.range(1, 10)
.collectList()
.subscribe(list -> System.out.println("List: " + list)); // List: [1, 2, ..., 10]
错误处理
// 错误处理操作符
Flux.error(new RuntimeException("Original error"))
.onErrorReturn("Fallback value")
.subscribe(System.out::println); // Fallback value
Flux.error(new RuntimeException("Original error"))
.onErrorResume(e -> Flux.just("Recovered from: " + e.getMessage()))
.subscribe(System.out::println); // Recovered from: Original error
Flux.just(1, 2, 0, 3)
.map(n -> 10 / n)
.onErrorContinue((e, value) -> System.out.println("Skipping " + value))
.subscribe(System.out::println); // 10, 5, 3
Flux.just(1, 2, 3)
.doOnError(e -> System.err.println("Error: " + e.getMessage()))
.subscribe();
// 使用 retry
Flux.error(new RuntimeException("Retry error"))
.retry(3)
.subscribe(
System.out::println,
e -> System.err.println("Failed after retries: " + e.getMessage())
);
调度器
// 使用调度器
Flux.range(1, 10)
.subscribeOn(Schedulers.parallel()) // 在并行线程池执行
.subscribe(System.out::println);
Flux.range(1, 10)
.subscribeOn(Schedulers.boundedElastic()) // 在弹性线程池执行
.observeOn(Schedulers.single()) // 在单线程池观察结果
.subscribe(System.out::println);
// 创建自定义调度器
Scheduler customScheduler = Schedulers.newParallel("custom", 4);
Flux.range(1, 10)
.subscribeOn(customScheduler)
.subscribe(System.out::println);
// 使用虚拟线程(Java 21+)
Flux.range(1, 10)
.subscribeOn(Schedulers.fromExecutor(Executors.newVirtualThreadPerTaskExecutor()))
.subscribe(System.out::println);
背压处理
// 背压策略
Flux.range(1, 1000)
.onBackpressureBuffer(100) // 缓冲最多 100 个元素
.subscribe();
Flux.range(1, 1000)
.onBackpressureDrop(dropped -> System.out.println("Dropped: " + dropped))
.subscribe();
Flux.range(1, 1000)
.onBackpressureLatest() // 只保留最新元素
.subscribe();
// 使用 request 控制流速
Flux.range(1, 10)
.subscribe(new BaseSubscriber<Integer>() {
@Override
protected void hookOnSubscribe(Subscription subscription) {
request(2); // 初始请求 2 个元素
}
@Override
protected void hookOnNext(Integer value) {
System.out.println("Received: " + value);
request(1); // 处理完一个再请求一个
}
});
实际应用示例
// 响应式 REST 客户端
@Service
public class UserService {
private final WebClient webClient;
public UserService(WebClient.Builder webClientBuilder) {
this.webClient = webClientBuilder.baseUrl("https://api.example.com").build();
}
public Mono<User> getUserById(Long id) {
return webClient.get()
.uri("/users/{id}", id)
.retrieve()
.onStatus(HttpStatus::isError, response ->
response.bodyToMono(String.class)
.flatMap(body -> Mono.error(new UserNotFoundException(body)))
)
.bodyToMono(User.class);
}
public Flux<User> getUsersByDepartment(String department) {
return webClient.get()
.uri("/users?department={department}", department)
.retrieve()
.bodyToFlux(User.class);
}
public Mono<User> createUser(UserCreateRequest request) {
return webClient.post()
.uri("/users")
.bodyValue(request)
.retrieve()
.onStatus(HttpStatus::isError, response ->
response.bodyToMono(String.class)
.flatMap(body -> Mono.error(new UserCreationException(body)))
)
.bodyToMono(User.class);
}
}
// 响应式数据访问
@Repository
public class ReactiveUserRepository {
private final MongoClient mongoClient;
public ReactiveUserRepository(MongoClient mongoClient) {
this.mongoClient = mongoClient;
}
public Mono<User> findById(String id) {
return getCollection().findById(id);
}
public Flux<User> findByDepartment(String department) {
return getCollection().find(query(where("department").is(department)));
}
public Mono<User> save(User user) {
return getCollection().save(user);
}
public Mono<Void> deleteById(String id) {
return getCollection().deleteById(id);
}
private MongoCollection<User> getCollection() {
return mongoClient.getDatabase("example").getCollection("users", User.class);
}
}
// 组合多个异步操作
@Service
public class OrderService {
private final UserService userService;
private final ProductService productService;
private final OrderRepository orderRepository;
public OrderService(UserService userService, ProductService productService, OrderRepository orderRepository) {
this.userService = userService;
this.productService = productService;
this.orderRepository = orderRepository;
}
public Mono<Order> createOrder(Long userId, Long productId, int quantity) {
return Mono.zip(
userService.getUserById(userId),
productService.getProductById(productId)
)
.flatMap(tuple -> {
User user = tuple.getT1();
Product product = tuple.getT2();
if (product.getStock() < quantity) {
return Mono.error(new InsufficientStockException());
}
Order order = new Order();
order.setUserId(userId);
order.setProductId(productId);
order.setQuantity(quantity);
order.setTotalPrice(product.getPrice() * quantity);
order.setStatus("PENDING");
return orderRepository.save(order);
})
.doOnSuccess(order -> {
// 更新库存
productService.updateStock(productId, -quantity).subscribe();
})
.doOnError(e -> {
// 记录错误日志
System.err.println("Failed to create order: " + e.getMessage());
});
}
}
测试响应式代码
// 测试 Mono
@Test
void monoTest() {
Mono<String> mono = Mono.just("Hello");
StepVerifier.create(mono)
.expectNext("Hello")
.expectComplete()
.verify();
}
// 测试 Flux
@Test
void fluxTest() {
Flux<Integer> flux = Flux.range(1, 5);
StepVerifier.create(flux)
.expectNext(1, 2, 3, 4, 5)
.expectComplete()
.verify();
}
// 测试错误处理
@Test
void errorTest() {
Flux<String> errorFlux = Flux.error(new RuntimeException("Test error"));
StepVerifier.create(errorFlux)
.expectError(RuntimeException.class)
.verify();
}
// 测试转换操作
@Test
void transformTest() {
Flux<String> flux = Flux.just("a", "b", "c")
.map(String::toUpperCase);
StepVerifier.create(flux)
.expectNext("A", "B", "C")
.expectComplete()
.verify();
}
最佳实践
- 避免阻塞:不要在响应式链中调用阻塞方法
- 正确处理背压:根据消费者能力控制流速
- 使用合适的调度器:根据任务类型选择调度器
- 组合操作符:使用 compose 和 transform 组合操作符
- 错误处理:在适当的位置处理错误
- 资源管理:使用 using 管理资源生命周期
- 避免嵌套订阅:使用 flatMap 代替嵌套 subscribe
- 测试验证:使用 StepVerifier 测试响应式代码
实际应用场景
- 高并发 API:处理大量并发请求
- 实时数据流:如 WebSocket 通信
- 数据处理管道:批量数据处理
- 异步组合:组合多个异步操作
总结
Reactor 框架为 Java 提供了强大的响应式编程能力。通过 Mono 和 Flux 两个核心类型,可以优雅地处理异步数据流。合理使用操作符和调度器,可以构建高性能、高并发的应用程序。
别叫我大神,叫我 Alex 就好。这其实可以更优雅一点,响应式编程让异步操作变得更加简洁和高效。
转载自 CSDN-专业IT技术社区
原文链接:https://blog.csdn.net/alex_goden/article/details/160913902



