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

Spring Boot升级到3后Flux.concat失效,文件写入失败求助

问题分析与解决方案

核心原因

Spring Boot 3(对应Spring WebFlux 6.x)使用的Reactor版本(3.5+)对异步操作的完成信号传递逻辑进行了调整。原代码中Flux.concat串行执行两个文件写入操作后,then触发后续文件读取时,可能因写入操作的完成信号未正确等待磁盘写入完成就发出,导致读取到空文件或不完整文件,进而引发XML解析错误。而直接调用subscribe()会强制等待整个异步写入流程完成,因此能正常写入。

解决方案

1. 改用Mono.when确保写入完成

将串行执行的Flux.concat改为并行等待所有写入操作完成的Mono.when,确保两个文件都完全写入磁盘后再执行后续读取操作:

return Mono.when(
        SmsTemplate.transferTo(SmsTemplateFile),
        lSMSVariables.transferTo(SMSVariablesFile)
)
.then(importService.previewSmsTemplateImport(clientId, spaceKey,
        new FileSystemResource(SmsTemplateFile),
        new FileSystemResource(SMSVariablesFile)))
.map(mapper::map)
.doFinally($ -> {
    FileUtils.deleteQuietly(SmsTemplateFile.toFile());
    FileUtils.deleteQuietly(SMSVariablesFile.toFile());
});

Mono.when会并行执行两个写入操作,仅当两者都发出完成信号后才进入后续的then逻辑,从根源上避免了提前读取文件的问题。

2. 确保transferTo方法正确等待写入完成

如果自定义的transferTo方法未正确处理写入完成的信号,需调整其实现,确保返回的Mono<Void>仅在文件写入、缓冲区刷新且资源释放后才完成:

public Mono<Void> transferTo(Path target) {
    return Mono.using(
            // 打开输出流
            () -> Files.newOutputStream(target, StandardOpenOption.CREATE, StandardOpenOption.TRUNCATE_EXISTING),
            // 写入数据并强制刷新缓冲区
            outputStream -> DataBufferUtils.write(dataBufferPublisher, outputStream)
                    .then(Mono.fromRunnable(() -> {
                        try {
                            outputStream.flush();
                        } catch (IOException e) {
                            throw new UncheckedIOException(e);
                        }
                    })),
            // 关闭流释放资源
            outputStream -> {
                try {
                    outputStream.close();
                } catch (IOException e) {
                    // 忽略或记录日志
                }
            }
    );
}

该实现通过Mono.using管理流的生命周期,显式刷新缓冲区并关闭流,确保数据完全写入磁盘后才发出完成信号。

3. 同步读取方法的适配(可选)

如果previewSmsTemplateImport是同步阻塞方法,需用Mono.fromCallable包裹,避免阻塞WebFlux的事件循环:

.then(Mono.fromCallable(() -> 
        importService.previewSmsTemplateImport(clientId, spaceKey,
                new FileSystemResource(SmsTemplateFile),
                new FileSystemResource(SMSVariablesFile))
))

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 12:13:12