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
相关产品推荐
相关产品推荐

