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

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函数逻辑与拆分步骤完全一致:

  1. 从Mono<CarrierDto>中提取List<Daily>
  2. 将列表转换为包含所有元素的Flux<Daily>
  3. 映射每个Daily对象为Price对象

由于Reactor的冷流特性,每次订阅都会重新执行流逻辑,因此无论使用transform整合调用还是拆分步骤调用,都会输出全部4条Price数据。

内容的提问来源于stack exchange,提问作者D-Dᴙum

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 20:19:49