Reactor Transform函数执行结果不符合预期问题排查
Reactor Transform函数整合执行与拆分执行结果不一致问题分析
问题描述
创建的Reactor Transform函数在整合执行时仅返回一条Price数据,但拆分步骤单独执行时能返回全部4条数据,两种执行方式输出结果不一致。
相关代码定义
CarrierDto与Price实体
public record CarrierDto(MetaData metaData, List<Daily> dailySeries) { public record MetaData(String name, String information) {}; public record Daily(String date, Integer price) {}; }
@Table("price") public record Price ( @Column("p_date") String date, @Column("p_price") Integer price) { }
存在问题的Transform函数
static Function<Mono<CarrierDto>, Flux<Price>> transform = carrierDtoMono - carrierDtoMono.map(CarrierDto::dailySeries) .flatMapMany(Flux::fromIterable) .map(daily -> new Price(daily.date(), daily.price()));
测试代码
import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import java.util.Arrays; import java.util.List; import java.util.function.Function; public class TestMain { static Function<Mono<CarrierDto>, Flux<Price>> transform = carrierDtoMono - carrierDtoMono.map(CarrierDto::dailySeries) .flatMapMany(Flux::fromIterable) .map(daily -> new Price(daily.date(), daily.price())); public static void main(String[] args) { final CarrierDto dto = new CarrierDto(new CarrierDto.MetaData("Test", "Test information"), Arrays.asList(new CarrierDto.Daily("2023-01-01", 100), new CarrierDto.Daily("2023-01-02", 101), new CarrierDto.Daily("2023-01-03", 102), new CarrierDto.Daily("2023-01-04", 103))); System.out.println("Using Transform function."); Mono.just(dto) .transform(transform) .subscribe(price -> System.out.println(String.format("Date: %s Price: %d", price.date(), price.price()))); System.out.println("Without transform function."); final Mono<List<CarrierDto.Daily>> monoDailySeries = Mono.just(dto).map(CarrierDto::dailySeries); final Flux<CarrierDto.Daily> dailyFlux = monoDailySeries.flatMapMany(Flux::fromIterable); final Flux<Price> priceFlux = dailyFlux.map(daily -> new Price(daily.date(), daily.price())); priceFlux.subscribe(price -> System.out.println(String.format("Date: %s Price: %d", price.date(), price.price()))); } }
问题原因
核心问题是Transform函数的Lambda表达式语法错误:将Lambda的箭头运算符->误写为减号-。
这段错误的代码会被编译器错误解析,尝试对carrierDtoMono参数与后续的Flux流进行无效的减法运算,最终导致流逻辑异常,仅输出第一条数据。
修复方案
将Lambda表达式的减号-修正为箭头->,确保函数逻辑符合预期:
static Function<Mono<CarrierDto>, Flux<Price>> transform = carrierDtoMono -> carrierDtoMono.map(CarrierDto::dailySeries) .flatMapMany(Flux::fromIterable) .map(daily -> new Price(daily.date(), daily.price()));
验证说明
修正后的transform函数逻辑与拆分步骤完全一致:
- 从
Mono<CarrierDto>中提取List<Daily> - 将列表转换为包含所有元素的
Flux<Daily> - 映射每个
Daily对象为Price对象
由于Reactor的冷流特性,每次订阅都会重新执行流逻辑,因此无论使用transform整合调用还是拆分步骤调用,都会输出全部4条Price数据。
内容的提问来源于stack exchange,提问作者D-Dᴙum
相关产品推荐
相关产品推荐

