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

Reactive WebFlux读取CSV仅获取单行问题求助

解决Flux.generate仅读取CSV第一行的问题

问题根源

你的readLinesFromCsv方法中,Flux.generate的实现逻辑存在错误:

  • 状态初始化器仅在第一次调用时执行一次,读取第一行后就不再继续读取新行
  • 生成器函数每次返回null,导致后续状态变为null,无法触发下一次读取操作

修复后的readLinesFromCsv方法

修改Flux.generate的实现,让它持续读取文件直到结束:

private Flux<String> readLinesFromCsv() {
    return Flux.using(
        () -> new BufferedReader(new InputStreamReader(resourceLoader.getResource("classpath:mycsvFile.csv").getInputStream())),
        reader -> Flux.generate(
            () -> reader, // 将BufferedReader作为状态维护,供每次读取复用
            (bufferedReader, sink) -> {
                try {
                    String line = bufferedReader.readLine();
                    if (line != null) {
                        sink.next(line);
                    } else {
                        sink.complete();
                    }
                } catch (IOException e) {
                    sink.error(e);
                }
                return bufferedReader; // 保持状态,用于下一次生成调用
            }
        ),
        reader -> {
            try {
                reader.close();
            } catch (IOException e) {
                Exceptions.propagate(e);
            }
        }
    );
}

关键修改点:

  • 将BufferedReader作为Flux.generate的状态对象,每次生成元素时都用它读取下一行
  • 移除多余的.map(buffer -> buffer.toString()),因为readLine()返回的已经是String类型
  • 读取异常时通过sink.error()传递错误,便于后续错误处理

额外优化:修复数据处理的潜在问题

原代码中buffer(100)会把CSV表头和数据行混在同一个批次中,导致后续非第一个批次的lines.get(0)取到的是数据行而非表头,引发转换错误。建议先分离表头和数据行:

public Mono<ResponseEntity<Void>> readFromFileAndStoreData() {
    return readLinesFromCsv()
        .take(1) // 单独读取表头行
        .flatMap(headerLine -> {
            String[] headers = headerLine.split(",");
            log.info("headers to read: {}", Arrays.toString(headers));
            
            // 处理剩余的数据行,按100条分批存储
            return readLinesFromCsv()
                .skip(1) // 跳过表头行
                .buffer(100)
                .flatMap(lines -> {
                    log.info("total records to store: {}", lines.size());
                    List<MyPojoData> pojos = lines.stream()
                        .map(line -> convert(headers, line.split(",")))
                        .collect(Collectors.toList());
                    return Mono.fromCallable(() -> myDataRepository.saveAll(pojos))
                        .onErrorResume(e -> {
                            log.error("Failed to save batch", e);
                            return Mono.empty();
                        });
                })
                .then(); // 等待所有批次处理完成
        })
        .map(unused -> ResponseEntity.ok().<Void>build())
        .onErrorResume(e -> {
            log.error("Error processing CSV", e);
            return Mono.just(ResponseEntity.status(HttpStatus.INTERNAL_SERVER_ERROR).build());
        });
}

优化说明:

  • 先单独读取表头,再处理后续数据行,避免表头被混入数据批次
  • 移除直接调用subscribe()的方式,改为返回Mono<ResponseEntity>,让Spring WebFlux统一控制执行流程
  • 添加错误日志,方便排查存储失败的问题

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 08:45:18