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

如何用Spring WebFlux流式处理CSV并生成分块NDJSON响应?

基于Spring WebFlux的响应式CSV流式处理方案

你的核心问题是当前代码会将文件块直接当作完整行处理,且未处理跨buffer的换行场景,无法真正实现逐行流式处理。以下是可行的解决方案:

关键思路

  1. 利用Spring WebFlux内置的LineDecoder,将流式的DataBuffer拆分为完整的文本行,它会自动处理跨多个buffer的换行,且不会一次性加载整个文件到内存。
  2. 配合成熟的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 12:53:10