You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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的流处理链中:

  1. map操作返回的是Stream<Mono<Void>>,随后flatMapMany(Flux::fromStream)将其转换为Flux<Mono<Void>>——这里的Flux仅会发射Mono<Void>对象本身,而非订阅并执行这些Mono的内部逻辑。
  2. 最终的.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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.24 01:39:50