博主最近想使用 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 系统的核心:将文本转换为向量
三、问题出现的完整流程
下面是博主实现文件上传与解析功能的完整流程
概括起来其实主要分为七步:

步骤 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);
}
});
}
解析流程:
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 调用
- 异步处理:文档入库可以异步处理,不阻塞上传响应
- 缓存机制:解析结果可以缓存,避免重复解析
- 分片上传:大文件可以分片上传,提高上传速度
以上就是博主了解的文件上传与解析功能的所有原理和流程啦。核心其实就是:响应式编程 + 流式处理 + 自动解析 + 向量存储,实现高效、可靠的文档处理流程。
好了,博主的分享到此就结束啦!
如果博文对您有帮助的话,可以给博主一键三连哟 !记得关注博主哟,会继续改进学习创作更多优秀的文章!!!
如果博文内容有误,也欢迎各位佬在评论区批评指正!!!


