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

如何将CSV反序列化为Either<实体类, Exception>类型的Flux流?

改造方案:返回Either<实体类, Exception>的流式处理实现

完全可以改造为返回Flux<Either<T, Exception>>的形式,同时保持大文件友好的流式处理特性,不需要改用不适合大文件的parse()方法。核心思路是在流式解析每条CSV记录时捕获异常,并将正常记录和异常分别包装到Either的左右分支中,同时确保资源安全和内存高效。

步骤说明

  1. 自定义/引入Either类型:如果不想引入第三方库(如Vavr),可以自己实现一个极简的Either密封接口;若使用Vavr则直接用io.vavr.control.Either。
  2. 关闭CsvToBean自动抛异常:通过withThrowExceptions(false)让解析器将异常暂存,而非直接中断流。
  3. 流式捕获并包装结果:遍历解析流时,检查每条记录对应的异常,分别包装为Either的Right(正常实体)或Left(异常)。
  4. 安全管理资源:用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 16:05:25