如何将CSV反序列化为Either<实体类, Exception>类型的Flux流?
改造方案:返回Either<实体类, Exception>的流式处理实现
完全可以改造为返回Flux<Either<T, Exception>>的形式,同时保持大文件友好的流式处理特性,不需要改用不适合大文件的parse()方法。核心思路是在流式解析每条CSV记录时捕获异常,并将正常记录和异常分别包装到Either的左右分支中,同时确保资源安全和内存高效。
步骤说明
- 自定义/引入Either类型:如果不想引入第三方库(如Vavr),可以自己实现一个极简的Either密封接口;若使用Vavr则直接用
io.vavr.control.Either。 - 关闭CsvToBean自动抛异常:通过
withThrowExceptions(false)让解析器将异常暂存,而非直接中断流。 - 流式捕获并包装结果:遍历解析流时,检查每条记录对应的异常,分别包装为Either的Right(正常实体)或Left(异常)。
- 安全管理资源:用
Flux.using确保输入流和Reader在处理完成后自动关闭,避免资源泄漏。
代码实现
1. 自定义极简Either接口(可选)
如果不想依赖第三方库,可先定义一个简单的Either类型:
public sealed interface Either<L, R> { record Left<L, R>(L value) implements Either<L, R> {} record Right<L, R>(R value) implements Either<L, R> {} }
2. 改造后的deserialize方法
import org.springframework.web.multipart.MultipartFile; import reactor.core.publisher.Flux; import com.opencsv.bean.CsvToBean; import com.opencsv.bean.CsvToBeanBuilder; import com.opencsv.bean.HeaderColumnNameMappingStrategy; import com.opencsv.bean.MappingStrategy; import java.io.BufferedReader; import java.io.InputStreamReader; import java.io.IOException; public <T> Flux<Either<Exception, T>> deserialize(MultipartFile multipartFile, Class<T> type) { // 使用Flux.using管理资源,确保流关闭 return Flux.using( // 资源提供者:创建BufferedReader () -> new BufferedReader(new InputStreamReader(multipartFile.getInputStream())), // 流处理逻辑 reader -> { MappingStrategy<T> ms = new HeaderColumnNameMappingStrategy<>(); ms.setType(type); CsvToBean<T> csvToBean = new CsvToBeanBuilder<T>(reader) .withType(type) .withMappingStrategy(ms) .withThrowExceptions(false) // 关键:禁止自动抛出解析异常,改为暂存 .build(); // 流式遍历解析结果,包装为Either return Flux.fromStream(csvToBean.stream()) .handle((record, sink) -> { // 检查当前记录是否有解析异常 if (!csvToBean.getCapturedExceptions().isEmpty()) { // 将异常包装为Either.Left并发送 csvToBean.getCapturedExceptions().forEach(ex -> sink.next(new Either.Left<>(ex))); csvToBean.getCapturedExceptions().clear(); // 清理异常列表,避免影响下一条记录 } else { // 正常记录包装为Either.Right并发送 sink.next(new Either.Right<>(record)); } }); }, // 资源清理:关闭BufferedReader BufferedReader::close ) // 捕获流创建阶段的IO异常,包装为Either.Left .onErrorMap(IOException.class, Either.Left::new); }
关键细节说明
- 大文件友好:依然依赖
csvToBean.stream()的流式特性,不会一次性加载整个CSV文件到内存,内存占用保持恒定。 - 异常处理全覆盖:既捕获了输入流创建时的
IOException,也捕获了每条记录的解析异常(如字段类型不匹配、缺失必填字段等)。 - 资源安全:
Flux.using会自动管理BufferedReader的生命周期,无论流正常完成还是异常终止,都会执行关闭操作。
如果使用Vavr的Either,只需将自定义的Either.Left/Either.Right替换为Either.left()/Either.right()即可,逻辑完全一致。
内容的提问来源于stack exchange,提问作者asikorski
相关产品推荐
相关产品推荐

