如何用Spring WebFlux流式处理CSV并生成分块NDJSON响应?
基于Spring WebFlux的响应式CSV流式处理方案
你的核心问题是当前代码会将文件块直接当作完整行处理,且未处理跨buffer的换行场景,无法真正实现逐行流式处理。以下是可行的解决方案:
关键思路
- 利用Spring WebFlux内置的
LineDecoder,将流式的DataBuffer拆分为完整的文本行,它会自动处理跨多个buffer的换行,且不会一次性加载整个文件到内存。 - 配合成熟的CSV解析库(如Apache Commons CSV)对每行进行解析,保持响应式流的特性,每行解析完成后立即输出为NDJSON块。
修正后的代码示例
import org.springframework.core.io.buffer.DataBufferUtils; import org.springframework.core.io.buffer.LineDecoder; import org.springframework.http.HttpStatus; import org.springframework.http.MediaType; import org.springframework.http.codec.multipart.FilePartEvent; import org.springframework.http.codec.multipart.FormPartEvent; import org.springframework.http.codec.multipart.PartEvent; import org.springframework.web.bind.annotation.PostMapping; import org.springframework.web.bind.annotation.RequestBody; import org.springframework.web.bind.annotation.ResponseStatus; import org.springframework.web.bind.annotation.RestController; import org.apache.commons.csv.CSVFormat; import org.apache.commons.csv.CSVParser; import org.apache.commons.csv.CSVRecord; import reactor.core.publisher.Flux; import java.io.IOException; import java.nio.charset.StandardCharsets; @RestController public class CsvProcessingController { @PostMapping(path = "process", consumes = MediaType.MULTIPART_FORM_DATA_VALUE, produces = MediaType.APPLICATION_NDJSON_VALUE) @ResponseStatus(HttpStatus.OK) public Flux<Data> process(@RequestBody Flux<PartEvent> allPartsEvents) { // 创建UTF-8编码的行解码器,自动处理跨buffer换行 LineDecoder lineDecoder = LineDecoder.allocate(StandardCharsets.UTF_8); return allPartsEvents.windowUntil(PartEvent::isLast) .concatMap(window -> window.switchOnFirst((signal, partEvents) -> { if (!signal.hasValue()) { return Flux.error(new RuntimeException("空的分片事件")); } PartEvent event = signal.get(); if (event instanceof FormPartEvent formEvent) { return Flux.just(new Data(formEvent.value())); } else if (event instanceof FilePartEvent) { // 将分片事件转换为DataBuffer流,再解码为完整行 Flux<String> csvLines = partEvents .map(PartEvent::content) .transform(lineDecoder::decode) .doOnDiscard(DataBuffer.class, DataBufferUtils::release); // 确保释放所有缓冲区,避免内存泄漏 // 处理CSV:提取表头,解析数据行 return csvLines.switchOnFirst((lineSignal, lineFlux) -> { if (!lineSignal.hasValue()) { return Flux.error(new RuntimeException("空的CSV文件")); } // 解析表头行,构建CSV格式 String headerLine = lineSignal.get(); CSVFormat csvFormat = CSVFormat.DEFAULT.withHeader(headerLine.split(",")) .withTrim(); // 可选:去除字段前后空格 // 跳过表头行,解析剩余数据行 return lineFlux.skip(1) .map(line -> { try (CSVParser parser = CSVParser.parse(line, csvFormat)) { CSVRecord record = parser.getRecords().get(0); // 根据CSV字段构造Data对象,此处需根据你的Data类结构调整 return new Data(record.get("列1"), record.get("列2")); } catch (IOException e) { throw new RuntimeException("解析CSV行失败: " + line, e); } }); }); } else { return Flux.error(new RuntimeException("不支持的事件类型: " + event.getClass().getName())); } })); } // 示例Data类,需根据实际需求调整 public static class Data { private String field1; private String field2; public Data(String field1, String field2) { this.field1 = field1; this.field2 = field2; } // 需提供getter,以便Spring将其序列化为JSON public String getField1() { return field1; } public String getField2() { return field2; } } }
核心细节说明
- LineDecoder的作用:它会逐个读取DataBuffer中的字节,累计到换行符时输出完整行,完美处理行内容跨多个DataBuffer的场景,实现真正的逐流式处理。
- 内存泄漏防护:通过
doOnDiscard确保所有DataBuffer在流结束或出错时被释放,避免内存泄漏。 - CSV解析的响应式适配:使用Apache Commons CSV对单行进行解析,保持流的响应式特性,每行解析完成后立即生成Data对象,Spring会自动将
Flux<Data>序列化为分块的NDJSON响应。
内容的提问来源于stack exchange,提问作者mfudi
相关产品推荐
相关产品推荐

