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

Reactor响应式代码改造中excelParser.readFile重复调用问题求助

问题:Reactor响应式代码中excelParser.readFile被调用两次

我是Reactor新手,正在将一段同步代码改造为响应式流实现,功能可正常运行,但发现excelParser.readFile被调用了两次,不清楚其中原因,特向大家求助。

原同步代码

public String magnetCreation(String name, MultipartFile file) {

    // XLS file save
    var xls = fileStorage.save(file);
    // Excel Parsing
    var data = excelParser.readFile(xls.toAbsolutePath().toString());

    // Delete excel file
    fileStorage.deleteFile(xls);

    // Create one CSV file for each BH curve
    List<MagnetBHCurveURL> urls = new ArrayList<>();

    final MagnetDetails details = new MagnetDetails(name, false, List.of());

    data.bhCurves().forEach(curve -> {

        var csvFile = generateCSV(curve);

        // Now the file name is generated on the fileserver
        var savedName = fileServer.uploadFile(csvFile, uploadPath).block();
        urls.add(new MagnetBHCurveURL(savedName, curve.temp));
    });

    log.info("URL csv upload: {}", urls);

    // thing creation
    var thingJsonObject = magnetUtility.createThing(details, data.summary(), urls);
    var id = twinUtility.create(thingJsonObject);
    // Return thingId
    return id;
}

改造后的响应式代码

public Mono<String> magnetCreation(String name, Mono<FilePart> file) {
var xlsElaborated = fileStorage.save(file).map(path -> excelParser.readFile(path.toAbsolutePath().toString()));

var urls = xlsElaborated.flatMapMany(data -> Flux.fromIterable(data.bhCurves()))
        .flatMap(curve -> {
            var csv = generateCSV(curve);
            return fileServer.uploadFile(csv, uploadPath)
                    .map(url -> new MagnetBHCurveURL(url, curve.temp));
        }).collect(Collectors.toList());

return Mono.zip(xlsElaborated, urls).map(x -> {
    var mData = x.getT1();
    var urlsLists = x.getT2();
    final MagnetDetails details = new MagnetDetails(name, false, List.of());
    return magnetUtility.createThing(details, mData.summary(), urlsLists);
}).flatMap(thingJsonObject -> twinUtility.createAsync(thingJsonObject));
}

原因分析

Reactor中的Mono/Flux是惰性求值的,只有当存在订阅者时才会执行流水线内的操作。你的代码里xlsElaborated被触发了两次订阅:

  • 第一次是构建urls时,xlsElaborated.flatMapMany(...)会订阅上游的xlsElaborated,导致fileStorage.save(file).map(...)整条流水线执行,包括调用excelParser.readFile;
  • 第二次是Mono.zip(xlsElaborated, urls)时,xlsElaborated再次被订阅,流水线重复执行,所以excelParser.readFile被调用了第二次。

解决办法

方案1:用cache()缓存执行结果

给xlsElaborated添加cache()操作符,让多次订阅共享同一次执行的结果:

var xlsElaborated = fileStorage.save(file)
    .map(path -> excelParser.readFile(path.toAbsolutePath().toString()))
    .cache(); // 缓存结果,后续订阅复用第一次的执行输出

方案2:调整代码结构,避免重复订阅

将依赖xlsElaborated的逻辑放到同一个流水线中,通过嵌套flatMap避免多次订阅:

public Mono<String> magnetCreation(String name, Mono<FilePart> file) {
    return fileStorage.save(file)
        .map(path -> excelParser.readFile(path.toAbsolutePath().toString()))
        .flatMap(data -> {
            // 在同一个flatMap内处理urls生成和后续逻辑
            var urls = Flux.fromIterable(data.bhCurves())
                .flatMap(curve -> {
                    var csv = generateCSV(curve);
                    return fileServer.uploadFile(csv, uploadPath)
                        .map(url -> new MagnetBHCurveURL(url, curve.temp));
                })
                .collect(Collectors.toList());
            
            return urls.map(urlsLists -> {
                final MagnetDetails details = new MagnetDetails(name, false, List.of());
                return magnetUtility.createThing(details, data.summary(), urlsLists);
            });
        })
        .flatMap(thingJsonObject -> twinUtility.createAsync(thingJsonObject));
}

补充:恢复文件删除逻辑

原同步代码中有删除临时Excel文件的逻辑,你改造后的代码漏掉了,建议添加进去确保资源释放。可以用finallyDo(无论成功失败都会执行):

return fileStorage.save(file)
    .finallyDo(path -> fileStorage.deleteFile(path)) // 确保文件被删除
    .map(path -> excelParser.readFile(path.toAbsolutePath().toString()))
    // 后续逻辑...

如果fileStorage.deleteFile是同步方法,直接传入即可;如果是响应式方法,可以用finallyDo(Mono::fromRunnable)适配。

内容的提问来源于stack exchange,提问作者Mex

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 19:35:22