程序员鸭梨头像
关注

Java 响应式编程:Reactor 框架深度解析

Java 响应式编程:Reactor 框架深度解析

核心概念

响应式编程是一种编程范式,关注数据的异步流和变化传播。在 Java 中,Reactor 框架提供了强大的响应式编程支持,基于 Reactive Streams 规范实现。

Reactor 核心组件

  1. Mono:表示 0 或 1 个元素的异步序列
  2. Flux:表示 0 到 N 个元素的异步序列
  3. Scheduler:调度器,控制任务执行线程
  4. 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();
}

最佳实践

  1. 避免阻塞:不要在响应式链中调用阻塞方法
  2. 正确处理背压:根据消费者能力控制流速
  3. 使用合适的调度器:根据任务类型选择调度器
  4. 组合操作符:使用 compose 和 transform 组合操作符
  5. 错误处理:在适当的位置处理错误
  6. 资源管理:使用 using 管理资源生命周期
  7. 避免嵌套订阅:使用 flatMap 代替嵌套 subscribe
  8. 测试验证:使用 StepVerifier 测试响应式代码

实际应用场景

  • 高并发 API:处理大量并发请求
  • 实时数据流:如 WebSocket 通信
  • 数据处理管道:批量数据处理
  • 异步组合:组合多个异步操作

总结

Reactor 框架为 Java 提供了强大的响应式编程能力。通过 Mono 和 Flux 两个核心类型,可以优雅地处理异步数据流。合理使用操作符和调度器,可以构建高性能、高并发的应用程序。

别叫我大神,叫我 Alex 就好。这其实可以更优雅一点,响应式编程让异步操作变得更加简洁和高效。

转载自 CSDN-专业IT技术社区

原文链接:https://blog.csdn.net/alex_goden/article/details/160913902

文章来源转载

评论

赞0

评论列表

微信小程序
QQ小程序

关于作者

点赞数:0
关注数:0
粉丝:0
文章:0
关注标签:0
加入于:--