欢迎光临
我们一直在努力

Spring WebFlux + Apache Tika 文件上传与解析功能深度解析

博主最近想使用 Spring Boot + Spring WebFlux + Spring AI 做一个 RAG 知识库系统,结果遇到了一个非常有趣的技术挑战——如何在响应式编程环境下处理文件上传和解析!

【先贴一下核心功能上来~】

✅ 支持多种文件格式:PDF、Word、Excel、PPT、TXT、Markdown
✅ 自动文本提取:使用 Apache Tika 智能识别并提取文档内容
✅ 响应式处理:基于 Spring WebFlux 的非阻塞文件处理
✅ 向量化存储:文档切分后生成 Embedding 并存入 pgvector
✅ 完整错误处理:上传失败时自动清理文件和数据库记录

下面会分享我个人的一些思路和踩坑经历 ~~


一、文件上传与解析

文件上传与解析是 RAG 知识库系统的核心入口,它决定了后续检索和生成的质量。整个流程就像是一个"文档加工厂":

  • 接收原材料(文件上传)
  • 识别文件类型(格式验证)
  • 提取文本内容(文档解析)
  • 保存元信息(数据库记录)
  • 加工处理(切分、向量化、存储)
  • 没错,我们熟知的 Spring WebFlux、Apache Tika、pgvector 等等底层都有相关的实现。Spring 的响应式编程也是大量用到 Mono、Flux 等响应式类型。

    我认为文件上传与解析有两大特点:异步非阻塞和流式处理。这使得系统在处理大文件时仍能保持良好的性能以及避免内存溢出。

    下面我来简单说说文件上传与解析开发时的核心概念。


    二、核心概念

    1. Spring WebFlux 响应式编程

    WebFlux 是 Spring 5 引入的响应式 Web 框架

    • 基于 Reactor 库(Mono、Flux)
    • 非阻塞 I/O,提高并发性能
    • 适合处理流式数据(如文件上传、SSE 推送)

    关键类型:

    • Mono<T>:表示 0 或 1 个元素的异步序列
    • Flux<T>:表示 0 到 N 个元素的异步序列
    • FilePart:WebFlux 中的文件上传对象(非阻塞)

    2. Apache Tika 文档解析

    Tika 是 Apache 基金会的文档解析库

    • 自动识别文件类型(MIME 类型检测)
    • 统一接口提取文本(parseToString())
    • 支持 100+ 种文件格式(PDF、Office、图片等)

    核心 API:

    Tika tika = new Tika();
    String text = tika.parseToString(inputStream); // 一行代码搞定!

    3. 文档向量化流程

    RAG 系统的核心:将文本转换为向量

  • 文档切分:将长文档切分为固定大小的片段(chunk)
  • 生成 Embedding:使用 AI 模型将文本转换为向量
  • 向量存储:将向量存入 pgvector 数据库
  • 相似度检索:查询时通过向量相似度找到相关片段

  • 三、问题出现的完整流程

    下面是博主实现文件上传与解析功能的完整流程

    概括起来其实主要分为七步:
    在这里插入图片描述

    步骤 1:Controller 接收文件上传请求

    @PostMapping(value = "/upload", consumes = MediaType.MULTIPART_FORM_DATA_VALUE)
    public Mono<Result<Document>> uploadDocument(
    @RequestPart("file") FilePart filePart,
    @RequestPart(value = "description", required = false) FormFieldPart descriptionPart) {

    String description = descriptionPart != null ? descriptionPart.value() : null;

    return documentService.uploadDocument(filePart, description)
    .map(document -> Result.success("文档上传成功", document))
    .onErrorResume(e -> {
    log.error("上传文档失败", e);
    return Mono.just(Result.error("上传文档失败: " + e.getMessage()));
    });
    }

    关键点:

    • @RequestPart("file"):接收文件部分
    • @RequestPart(value = "description", required = false):接收可选的描述字段
    • Mono<Result<Document>>:返回响应式的单值结果
    • onErrorResume:错误处理,确保始终返回 Mono

    步骤 2:Service 层验证文件类型

    @Override
    @Transactional
    public Mono<Document> uploadDocument(FilePart filePart, String description) {
    String filename = filePart.filename();
    if (filename.isEmpty()) {
    return Mono.error(new IllegalArgumentException("文件名不能为空"));
    }

    // 验证文件类型
    if (!documentParserService.isSupportedFileType(filename)) {
    return Mono.error(new IllegalArgumentException(
    "不支持的文件类型,支持:PDF、Word、Excel、PPT、TXT、Markdown"));
    }
    // …
    }

    支持的文件类型:

    private static final Set<String> SUPPORTED_EXTENSIONS = Set.of(
    "pdf", "doc", "docx", "xls", "xlsx", "ppt", "pptx", "txt", "md"
    );

    验证逻辑:

    @Override
    public boolean isSupportedFileType(String filename) {
    if(filename == null){
    return false;
    }
    String extension = getFileExtension(filename).toLowerCase();
    return SUPPORTED_EXTENSIONS.contains(extension);
    }

    步骤 3:保存文件到本地磁盘

    // 保存文件到本地
    String savedFilename = System.currentTimeMillis() + "_" + filename;
    Path filePath = Paths.get(uploadDir, savedFilename);

    // 使用 DataBufferUtils 非阻塞写入文件
    return DataBufferUtils.write(filePart.content(), filePath,
    StandardOpenOption.CREATE,
    StandardOpenOption.WRITE,
    StandardOpenOption.TRUNCATE_EXISTING)
    .then(documentParserService.parseDocumentReactive(filePart))
    // …

    关键点:

    • System.currentTimeMillis() + "_" + filename:添加时间戳避免文件名冲突
    • DataBufferUtils.write():响应式写入文件(非阻塞)
    • .then():文件写入完成后,继续执行解析操作

    目录初始化:

    @PostConstruct
    public void init() {
    // 创建上传目录(在字段注入完成后执行)
    try {
    Path uploadPath = Paths.get(uploadDir);
    Files.createDirectories(uploadPath);
    log.info("文档上传目录已创建: {}", uploadPath.toAbsolutePath());
    } catch (Exception e) {
    log.error("创建上传目录失败: {}", uploadDir, e);
    }
    }

    ⚠️ 重要提示:

    • 使用 @PostConstruct 而非构造函数初始化目录
    • 因为 @Value 注入在构造函数执行时还未完成,会导致 NullPointerException

    步骤 4:使用 Tika 解析文件内容

    @Override
    public Mono<String> parseDocumentReactive(FilePart filePart) {
    if (filePart == null) {
    return Mono.error(new IllegalArgumentException("文件不能为空"));
    }

    String filename = filePart.filename();
    if (filename.isEmpty()) {
    return Mono.error(new IllegalArgumentException("文件名不能为空"));
    }

    // 使用 DataBufferUtils 读取文件内容为字节数组
    return DataBufferUtils.join(filePart.content())
    .map(dataBuffer -> {
    byte[] bytes = new byte[dataBuffer.readableByteCount()];
    dataBuffer.read(bytes);
    DataBufferUtils.release(dataBuffer); // 释放内存
    return bytes;
    })
    .map(bytes -> {
    try (InputStream inputStream = new ByteArrayInputStream(bytes)) {
    String text = tika.parseToString(inputStream);
    log.info("成功解析文档:{},提取文本长度:{}", filename, text.length());
    return text;
    } catch (IOException | TikaException e) {
    log.error("解析文档失败:{}", filename, e);
    throw new RuntimeException("文档解析失败:" + e.getMessage(), e);
    }
    });
    }

    解析流程:

  • DataBufferUtils.join(filePart.content()):将流式数据合并为单个 DataBuffer
  • .map():将 DataBuffer 转换为字节数组
  • DataBufferUtils.release(dataBuffer):**重要!**释放内存,避免内存泄漏
  • tika.parseToString(inputStream):Tika 自动识别文件类型并提取文本
  • Tika 的魔法:

    • 自动识别:根据文件内容(而非扩展名)识别文件类型
    • 统一接口:无论 PDF、Word、Excel,都使用 parseToString() 方法
    • 智能提取:自动处理表格、图片、格式等,提取纯文本

    步骤 5:保存文档元信息到数据库

    .flatMap(content -> {
    // 2. 保存文档元信息
    Document document = new Document();
    document.setFilename(filename);
    document.setFileType(documentParserService.getFileType(filename));
    try {
    document.setFileSize(Files.size(filePath));
    } catch (Exception e) {
    log.warn("获取文件大小失败", e);
    document.setFileSize(0L);
    }
    document.setDescription(description);
    document.setCreatedBy("system");

    Document savedDocument = documentMetaRepository.save(document);
    log.info("文档元信息已保存: ID={}, filename={}", savedDocument.getId(), filename);
    // …
    })

    Document 实体:

    @Entity
    @Table(name = "document")
    @Data
    @NoArgsConstructor
    public class Document {
    @Id
    @GeneratedValue(strategy = GenerationType.IDENTITY)
    private Long id;

    @Column(nullable = false)
    private String filename; // 原始文件名

    @Column(nullable = false)
    private String fileType; // 文件类型:pdf, docx, xlsx 等

    @Column
    private Long fileSize; // 文件大小(字节)

    @Column(columnDefinition = "text")
    private String description; // 文档描述

    @Column(name = "created_at")
    private LocalDateTime createdAt = LocalDateTime.now();

    @Column(name = "created_by")
    private String createdBy; // 上传用户
    }

    关键点:

    • .flatMap():将 Mono<String>(解析内容)转换为 Mono<Document>(保存结果)
    • Files.size(filePath):获取文件大小
    • documentMetaRepository.save(document):保存到数据库,自动生成 ID

    步骤 6:文档切分与向量化

    // 3. 文档入库(切分 + embedding + 存储)
    try {
    ragService.ingestDocument(content, savedDocument.getId());
    log.info("文档已入库: ID={}, filename={}", savedDocument.getId(), filename);
    } catch (Exception e) {
    log.error("文档入库失败: ID={}, filename={}", savedDocument.getId(), filename, e);
    // 回滚:删除已保存的文档元信息
    documentMetaRepository.delete(savedDocument);
    // 删除已保存的文件
    try {
    Files.deleteIfExists(filePath);
    } catch (Exception ex) {
    log.warn("删除文件失败", ex);
    }
    return Mono.<Document>error(new RuntimeException("文档入库失败: " + e.getMessage(), e));
    }

    文档入库流程(RagServiceImpl.ingestDocument):

    @Override
    public void ingestDocument(String documentContent, Long docId) {
    //1、文档切分 – 按500字切
    List<String> chunks = splitText(documentContent, 500);

    //2、批量生成embedding
    List<float[]> embeddings = embeddingService.embedTexts(chunks);

    //3、存入数据库
    for (int i = 0; i < chunks.size(); i++) {
    String embeddingStr = embeddingService.toVectorString(embeddings.get(i));
    documentRepository.insertWithVector(
    chunks.get(i),
    docId,
    embeddingStr,
    LocalDateTime.now()
    );
    }
    }

    文档切分策略:

    private List<String> splitText(String text, int chunkSize) {
    List<String> chunks = new ArrayList<>();
    int length = text.length();

    for (int i = 0; i < length; i += chunkSize) {
    int end = Math.min(i + chunkSize, length);
    chunks.add(text.substring(i, end));
    }

    return chunks;
    }

    向量存储(pgvector):

    @Modifying
    @Query(value = "INSERT INTO document_chunk (content, doc_id, embedding, created_at) " +
    "VALUES (:content, :docId, CAST(:embedding AS vector), :createdAt)",
    nativeQuery = true)
    @Transactional
    void insertWithVector(@Param("content") String content,
    @Param("docId") Long docId,
    @Param("embedding") String embedding,
    @Param("createdAt") LocalDateTime createdAt);

    ⚠️ 关键点:

    • CAST(:embedding AS vector):必须显式转换,否则会报类型错误
    • @Modifying:标记为修改操作,不能使用 RETURNING 子句

    步骤 7:错误处理与资源清理

    .doOnError(error -> {
    // 出错时删除文件
    try {
    Files.deleteIfExists(filePath);
    } catch (Exception e) {
    log.warn("删除文件失败", e);
    }
    });

    错误处理策略:

  • 文件保存失败:直接返回错误,不创建数据库记录
  • 解析失败:删除已保存的文件,返回错误
  • 向量化失败:删除已保存的文件和数据库记录,返回错误
  • 确保数据一致性:

    • 使用 @Transactional 确保数据库操作的原子性
    • 使用 doOnError 确保文件清理
    • 使用 try-catch 处理异常,避免资源泄漏

    四、为什么使用响应式编程?

    响应式编程的设计初衷

    原本的设计场景是高并发、非阻塞的场景,如代码所示👇

    // 传统阻塞式(Servlet)
    @PostMapping("/upload")
    public Result<Document> uploadDocument(@RequestParam MultipartFile file) {
    // 阻塞:等待文件上传完成
    // 阻塞:等待文件解析完成
    // 阻塞:等待数据库保存完成
    return Result.success(document);
    } // ← 线程被占用,无法处理其他请求

    // 响应式(WebFlux)
    @PostMapping("/upload")
    public Mono<Result<Document>> uploadDocument(@RequestPart FilePart filePart) {
    return documentService.uploadDocument(filePart, description)
    .map(Result::success); // 非阻塞,线程可以处理其他请求
    }

    响应式编程的优势

    1. 非阻塞 I/O

    • 传统方式:一个线程处理一个请求,线程被阻塞
    • 响应式方式:一个线程可以处理多个请求,线程不阻塞

    2. 流式处理

    • 文件上传时,可以边上传边处理,不需要等待整个文件上传完成
    • 适合处理大文件,避免内存溢出

    3. 背压(Backpressure)

    • 当处理速度跟不上数据产生速度时,可以自动调节
    • 避免系统过载

    简单来说

    响应式编程让系统在高并发场景下仍能保持良好的性能,特别是在处理文件上传、流式数据、实时推送等场景时,优势明显。


    五、关键技术点深度解析

    1. FilePart vs MultipartFile

    MultipartFile(传统 Servlet):

    @PostMapping("/upload")
    public Result upload(@RequestParam MultipartFile file) {
    byte[] bytes = file.getBytes(); // 阻塞:读取整个文件到内存
    // …
    }

    FilePart(WebFlux):

    @PostMapping("/upload")
    public Mono<Result> upload(@RequestPart FilePart filePart) {
    return DataBufferUtils.join(filePart.content()) // 流式读取
    .map(dataBuffer -> {
    // 处理数据
    });
    }

    区别:

    • MultipartFile:一次性读取整个文件到内存(阻塞)
    • FilePart:流式读取,可以边读边处理(非阻塞)

    2. DataBufferUtils 的使用

    DataBufferUtils.join():

    DataBufferUtils.join(filePart.content())
    .map(dataBuffer -> {
    byte[] bytes = new byte[dataBuffer.readableByteCount()];
    dataBuffer.read(bytes);
    DataBufferUtils.release(dataBuffer); // ⚠️ 重要:释放内存
    return bytes;
    });

    关键点:

    • join():将多个 DataBuffer 合并为一个
    • readableByteCount():获取可读字节数
    • release():必须调用,释放内存,避免内存泄漏

    3. Apache Tika 的魔法

    自动识别文件类型:

    Tika tika = new Tika();
    String mimeType = tika.detect(inputStream); // 自动检测 MIME 类型
    String text = tika.parseToString(inputStream); // 提取文本

    支持的格式:

    • 文档:PDF、Word、Excel、PPT、RTF、ODT
    • 文本:TXT、Markdown、HTML、XML
    • 图片:JPEG、PNG、GIF(OCR 提取文本)
    • 压缩:ZIP、TAR、GZIP

    工作原理:

  • 文件签名检测:读取文件头部字节,识别文件类型
  • 解析器选择:根据文件类型选择对应的解析器
  • 文本提取:使用解析器提取文本内容
  • 4. 文档切分策略

    固定大小切分:

    private List<String> splitText(String text, int chunkSize) {
    List<String> chunks = new ArrayList<>();
    int length = text.length();

    for (int i = 0; i < length; i += chunkSize) {
    int end = Math.min(i + chunkSize, length);
    chunks.add(text.substring(i, end));
    }

    return chunks;
    }

    优化建议:

    • 重叠切分:相邻 chunk 之间保留部分重叠,避免语义断裂
    • 按段落切分:优先在段落边界切分,保持语义完整
    • 智能切分:使用 NLP 技术识别句子边界

    示例(重叠切分):

    private List<String> splitTextWithOverlap(String text, int chunkSize, int overlap) {
    List<String> chunks = new ArrayList<>();
    int length = text.length();
    int i = 0;

    while (i < length) {
    int end = Math.min(i + chunkSize, length);
    chunks.add(text.substring(i, end));
    i += (chunkSize overlap); // 重叠部分
    }

    return chunks;
    }

    5. pgvector 向量存储

    向量类型转换:

    @Query(value = "INSERT INTO document_chunk (content, doc_id, embedding, created_at) " +
    "VALUES (:content, :docId, CAST(:embedding AS vector), :createdAt)",
    nativeQuery = true)
    void insertWithVector(@Param("content") String content,
    @Param("docId") Long docId,
    @Param("embedding") String embedding, // 格式:[0.1, 0.2, 0.3, …]
    @Param("createdAt") LocalDateTime createdAt);

    向量格式:

    • 输入:"[0.1, 0.2, 0.3, …]"(字符串)
    • 转换:CAST(:embedding AS vector)(PostgreSQL 类型)
    • 存储:vector(1536)(1536 维向量)

    相似度检索:

    @Query(value = "SELECT * FROM document_chunk " +
    "ORDER BY embedding <-> CAST(:queryEmbedding AS vector) " +
    "LIMIT :limit", nativeQuery = true)
    List<DocumentChunk> findSimilarChunks(@Param("queryEmbedding") String queryEmbedding,
    @Param("limit") int limit);

    运算符说明:

    • <->:欧氏距离(L2 距离)
    • <#>:余弦距离
    • <=>:内积距离

    六、踩坑经历

    坑 1:FilePart 与 MultipartFile 的混淆

    问题:

    // ❌ 错误:WebFlux 中不能使用 MultipartFile
    @PostMapping("/upload")
    public Mono<Result> upload(@RequestParam MultipartFile file) {
    // …
    }

    错误信息:

    org.springframework.web.server.ServerWebInputException:
    Required request part 'file' is not present

    解决方案:

    // ✅ 正确:使用 FilePart
    @PostMapping("/upload")
    public Mono<Result> upload(@RequestPart("file") FilePart filePart) {
    // …
    }

    原因:

    • MultipartFile 是 Servlet API,用于传统阻塞式应用
    • FilePart 是 WebFlux API,用于响应式应用

    坑 2:DataBuffer 内存泄漏

    问题:

    // ❌ 错误:未释放 DataBuffer
    return DataBufferUtils.join(filePart.content())
    .map(dataBuffer -> {
    byte[] bytes = new byte[dataBuffer.readableByteCount()];
    dataBuffer.read(bytes);
    // 忘记释放 dataBuffer!
    return bytes;
    });

    现象:

    • 内存占用持续增长
    • 最终导致 OutOfMemoryError

    解决方案:

    // ✅ 正确:必须释放 DataBuffer
    return DataBufferUtils.join(filePart.content())
    .map(dataBuffer -> {
    byte[] bytes = new byte[dataBuffer.readableByteCount()];
    dataBuffer.read(bytes);
    DataBufferUtils.release(dataBuffer); // ⚠️ 重要!
    return bytes;
    });

    原因:

    • DataBuffer 使用引用计数管理内存
    • 必须调用 release() 释放内存,否则会导致内存泄漏

    坑 3:pgvector 类型转换错误

    问题:

    // ❌ 错误:直接插入字符串,类型不匹配
    @Query(value = "INSERT INTO document_chunk (content, embedding) " +
    "VALUES (:content, :embedding)", nativeQuery = true)
    void insertWithVector(@Param("embedding") String embedding);

    错误信息:

    ERROR: column "embedding" is of type vector but expression is of type character varying

    解决方案:

    // ✅ 正确:使用 CAST 显式转换
    @Query(value = "INSERT INTO document_chunk (content, doc_id, embedding, created_at) " +
    "VALUES (:content, :docId, CAST(:embedding AS vector), :createdAt)",
    nativeQuery = true)
    void insertWithVector(@Param("embedding") String embedding);

    原因:

    • pgvector 的 vector 类型是 PostgreSQL 的自定义类型
    • 必须使用 CAST 显式转换,不能直接插入字符串

    坑 4:@Modifying 查询使用 RETURNING

    问题:

    // ❌ 错误:@Modifying 方法不能使用 RETURNING
    @Modifying
    @Query(value = "INSERT INTO document_chunk (…) VALUES (…) RETURNING id",
    nativeQuery = true)
    Long insertWithVector(...); // 期望返回 Long,但实际返回 int

    错误信息:

    org.springframework.orm.jpa.JpaSystemException:
    JDBC exception executing SQL […] [传回预期之外的结果。]

    解决方案:

    // ✅ 正确:移除 RETURNING,返回受影响行数
    @Modifying
    @Query(value = "INSERT INTO document_chunk (…) VALUES (…)",
    nativeQuery = true)
    @Transactional
    void insertWithVector(...); // 返回 void 或 int(受影响行数)

    原因:

    • @Modifying 方法期望返回 int(受影响行数)或 void
    • 不能使用 RETURNING 子句获取生成的值
    • 如果需要获取 ID,应该先保存实体,再获取 ID

    坑 5:@PostConstruct 与构造函数初始化

    问题:

    // ❌ 错误:在构造函数中使用 @Value 注入的字段
    public DocumentServiceImpl(...) {
    Path uploadPath = Paths.get(uploadDir); // uploadDir 还是 null!
    Files.createDirectories(uploadPath); // NullPointerException
    }

    错误信息:

    java.lang.NullPointerException
    at DocumentServiceImpl.<init>(DocumentServiceImpl.java:49)

    解决方案:

    // ✅ 正确:使用 @PostConstruct
    @PostConstruct
    public void init() {
    Path uploadPath = Paths.get(uploadDir); // 此时 uploadDir 已注入
    Files.createDirectories(uploadPath);
    }

    原因:

    • @Value 注入在构造函数执行之后才完成
    • 构造函数中访问 @Value 字段会得到 null
    • @PostConstruct 方法在依赖注入完成后执行,此时字段已注入

    七、总结

    1. 文件上传与解析的本质:

    • 响应式处理:使用 Mono、Flux 实现非阻塞文件处理
    • 流式读取:使用 DataBufferUtils 流式读取文件内容
    • 自动解析:使用 Apache Tika 自动识别并提取文本

    2. 问题的根源:

    • 响应式编程:需要理解 Mono、Flux 的使用方式
    • 内存管理:必须释放 DataBuffer,避免内存泄漏
    • 类型转换:pgvector 需要显式类型转换

    3. 解决思路:

    • 使用 FilePart:WebFlux 中必须使用 FilePart 而非 MultipartFile
    • 释放资源:使用 DataBufferUtils.release() 释放内存
    • 显式转换:使用 CAST(:embedding AS vector) 转换类型
    • 初始化时机:使用 @PostConstruct 而非构造函数初始化

    4. 实践场景:

    • 文件上传:使用 FilePart + DataBufferUtils 流式处理
    • 文档解析:使用 Apache Tika 统一接口提取文本
    • 向量存储:使用 pgvector 存储和检索向量
    • 错误处理:使用 doOnError 确保资源清理

    5. 性能优化建议:

    • 批量处理:批量生成 Embedding,减少 API 调用
    • 异步处理:文档入库可以异步处理,不阻塞上传响应
    • 缓存机制:解析结果可以缓存,避免重复解析
    • 分片上传:大文件可以分片上传,提高上传速度

    以上就是博主了解的文件上传与解析功能的所有原理和流程啦。核心其实就是:响应式编程 + 流式处理 + 自动解析 + 向量存储,实现高效、可靠的文档处理流程。

    好了,博主的分享到此就结束啦!

    如果博文对您有帮助的话,可以给博主一键三连哟 !记得关注博主哟,会继续改进学习创作更多优秀的文章!!!

    如果博文内容有误,也欢迎各位佬在评论区批评指正!!!

    赞(0)
    未经允许不得转载:171主机测评 » Spring WebFlux + Apache Tika 文件上传与解析功能深度解析
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址