Spring Reactive代码异常:反应式流后续处理逻辑未触发排查
反应式流分组处理未执行后续逻辑问题
问题描述
业务场景为将预订列表按年龄分组后逐个处理,因Flux的group和collectMap不符合需求,采用先collectList收集列表再用Stream分组的方式。期望所有处理完成后返回Mono.empty(),但当前代码仅执行到processBookings的打印语句,processBooking及后续Booking::3、Booking::4的打印逻辑完全未触发,且不想通过subscribe手动订阅解决。
原代码
@GetMapping("/test") public Mono<Void> test() { return prepareForBooking(); } @Data public class Booking { String name; String age; } private Mono<Void> prepareForBooking() { System.out.println("Booking::1"); Booking booking = new Booking(); booking.setName("San"); booking.setAge("18"); Booking booking1 = new Booking(); booking1.setName("Man"); booking1.setAge("19"); Booking booking2 = new Booking(); booking2.setName("Dan"); booking2.setAge("18"); Booking booking3 = new Booking(); booking3.setName("Can"); booking3.setAge("17"); List<Booking> bookings = new ArrayList<>(); bookings.add(booking); bookings.add(booking1); bookings.add(booking2); bookings.add(booking3); return Flux.fromIterable(bookings) .collectList() .map( bookingList -> { var bookingsByAge = bookingList.stream() .collect(Collectors.groupingBy(Booking::getAge)); return bookingsByAge.keySet().stream() .map( age -> { if (age.equals("18") || age.equals("17")) { return processBookings(bookingsByAge.get(age)); } else { return Mono.empty(); } }); }) .flatMapMany(Flux::fromStream) .then(); } private Mono<Void> processBookings(List<Booking> bookings) { System.out.println("Booking::2" + bookings); return Flux.fromIterable(bookings).flatMap(booking -> this.processBooking(booking)).collectList().flatMap(bookingNames -> { System.out.println("Booking::4" + bookingNames); return Mono.empty(); }); } private Mono<String> processBooking(Booking booking) { System.out.println("Booking::3" + booking.getName()); return Mono.just(booking.getName()); }
实际输出
Booking::1 Booking::2[ProviderIntegrationBookingController.Booking(name=Can, age=17)] Booking::2[ProviderIntegrationBookingController.Booking(name=San, age=18), ProviderIntegrationBookingController.Booking(name=Dan, age=18)]
预期输出
Booking::1 Booking::2[ProviderIntegrationBookingController.Booking(name=Can, age=17)] Booking::2[ProviderIntegrationBookingController.Booking(name=San, age=18), ProviderIntegrationBookingController.Booking(name=Dan, age=18)] Booking::3 San Booking::3 Dan Booking::3 Can Booking::4 [Can] Booking::4 [San, Dan]
问题根源
核心问题出在prepareForBooking的流处理链中:
map操作返回的是Stream<Mono<Void>>,随后flatMapMany(Flux::fromStream)将其转换为Flux<Mono<Void>>——这里的Flux仅会发射Mono<Void>对象本身,而非订阅并执行这些Mono的内部逻辑。- 最终的
.then()仅等待Flux发射完所有Mono<Void>对象就完成,不会触发这些Mono内部的Flux.fromIterable(...).flatMap(...)执行,导致processBooking和后续打印逻辑完全未运行。
解决方案
修改prepareForBooking的流处理逻辑,确保每个processBookings返回的Mono<Void>被订阅执行:
优化后的prepareForBooking方法
private Mono<Void> prepareForBooking() { System.out.println("Booking::1"); Booking booking = new Booking(); booking.setName("San"); booking.setAge("18"); Booking booking1 = new Booking(); booking1.setName("Man"); booking1.setAge("19"); Booking booking2 = new Booking(); booking2.setName("Dan"); booking2.setAge("18"); Booking booking3 = new Booking(); booking3.setName("Can"); booking3.setAge("17"); List<Booking> bookings = new ArrayList<>(); bookings.add(booking); bookings.add(booking1); bookings.add(booking2); bookings.add(booking3); return Flux.fromIterable(bookings) .collectList() .flatMapMany(bookingList -> { // 按年龄分组 Map<String, List<Booking>> bookingsByAge = bookingList.stream() .collect(Collectors.groupingBy(Booking::getAge)); // 转换为Flux并过滤目标年龄组,同时flatMap订阅每个processBookings返回的Mono return Flux.fromIterable(bookingsByAge.entrySet()) .filter(entry -> entry.getKey().equals("18") || entry.getKey().equals("17")) .flatMap(entry -> processBookings(entry.getValue())); }) .then(); }
可选优化:processBookings方法简化
将flatMap替换为doOnSuccess,逻辑更清晰:
private Mono<Void> processBookings(List<Booking> bookings) { System.out.println("Booking::2" + bookings); return Flux.fromIterable(bookings) .flatMap(this::processBooking) .collectList() .doOnSuccess(bookingNames -> System.out.println("Booking::4 " + bookingNames)) .then(); }
优化后执行结果
会完全符合预期输出,所有分组处理逻辑都会被触发,包括processBooking的打印和最终的结果收集打印。
内容的提问来源于stack exchange,提问作者Sandeep Nair
相关产品推荐
相关产品推荐

