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

