上传接口与文档主记录创建
这篇文档要讲的是:用户上传一份文档后,后端到底做了哪些事情?从 Controller 接到请求开始,一直到最后把消息丢进 Kafka 队列,整条链路我们一步步拆开来看。
先上一张总览流程图,有个整体印象之后再逐段看源码。
上传链路总览
Controller 层:接收 multipart 请求
上传接口的入口在 DocumentManageController,它是一个标准的 Spring MVC 控制器。这个接口比较特殊——不是普通的 JSON 请求,而是 multipart/form-data,同时包含二进制文件和文档元信息两部分。
/**
* 上传文档并投递后续解析任务。
* <p>
* 请求体采用 multipart/form-data,通常包含两部分:
* 1. {@code file}:真实上传的文档二进制内容;
* 2. {@code meta}:可选的文档元信息,例如展示名称、知识范围、业务分类、操作人等。
* </p>
* <p>
* 控制器本身不直接处理文件落盘、对象存储、文档入库或任务投递,
* 这里只负责把 multipart 中的两部分正确绑定出来,然后交给服务层统一完成:
* 文件校验、文件类型识别、对象存储上传、文档主表写入、任务表写入、日志记录和 Kafka 投递。
* </p>
*
* @param file 上传的文档文件,必须通过 multipart 的 {@code file} 字段传入
* @param dto 上传附带的文档元信息,来自 multipart 的 {@code meta} 字段,可为空
* @return 上传结果,包含文档 ID、任务 ID 以及当前解析/策略/索引状态
*/
@Operation(summary = "上传文档并投递解析任务")
@PostMapping(value = "/upload", consumes = MediaType.MULTIPART_FORM_DATA_VALUE)
public ApiResponse<DocumentUploadVo> upload(@RequestPart("file") MultipartFile file,
@Valid @RequestPart(value = "meta", required = false) DocumentUploadDto dto) {
// meta 允许前端不传;这里统一兜底成一个空 DTO,避免服务层重复判空,让后续取字段时更稳定。
return ApiResponse.ok(documentManageService.upload(file, dto == null ? new DocumentUploadDto() : dto));
}
这段代码的要点:
- 用
@RequestPart("file")接收上传的文件二进制内容,用@RequestPart("meta")接收可选的元信息 JSON meta标记了required = false,前端可以不传。Controller 这里做了一个兜底:如果dto为 null,就 new 一个空的DocumentUploadDto,这样服务层就不用到处判空了- Controller 本身不做任何业务逻辑,纯粹是参数绑定 + 转发给 Service
关于 DocumentUploadDto
DocumentUploadDto是上传时可选的元信息,包含以下字段:
documentName:文档展示名称(不传则用原始文件名)operatorId:操作人 ID(不传则视为系统触发)knowledgeScopeCode/knowledgeScopeName:知识范围编码和名称businessCategory:业务分类documentTags:文档标签
Service 层:upload 方法全流程
接下来就是核心了——DocumentManageServiceImpl.upload() 方法。这个方法没有在方法级别加 @Transactional 注解,而是通过 TransactionTemplate 手动控制事务边界:数据库写入操作(文档主表、任务表、任务日志)放在 transactionTemplate.execute() 内部,事务提交成功后再发送 Kafka 消息。这样做的目的是避免"消息已发出但事务还没提交"的时间窗口问题。
我们按执行顺序,一段一段来看。
第一步:文件校验
// 第一层防线:上传文件对象为空,或者 multipart 中虽然有字段但没有实际内容,都视为非法上传请求。
if (file == null || file.isEmpty()) {
throw new SuperAgentFrameException(DocumentManageCode.EMPTY_FILE_CONTENT.getCode(),
DocumentManageCode.EMPTY_FILE_CONTENT.getMsg());
}
// 原始文件名不仅用于展示,更用于后续识别文件类型,因此缺失时不能继续往下走。
String originalFileName = file.getOriginalFilename();
if (StrUtil.isBlank(originalFileName)) {
throw new SuperAgentFrameException(DocumentManageCode.UNSUPPORTED_FILE_TYPE.getCode(),
"上传文件缺少原始文件名,无法识别文件类型。");
}
// 根据文件名后缀识别文档类型;如果无法识别,就不允许进入后续存储和解析流程。
DocumentFileTypeEnum fileType = DocumentFileTypeEnum.fromFileName(originalFileName);
if (fileType == null) {
throw new SuperAgentFrameException(DocumentManageCode.UNSUPPORTED_FILE_TYPE.getCode(),
DocumentManageCode.UNSUPPORTED_FILE_TYPE.getMsg());
}
这里做了三层校验:
- 文件本身不能为空——
file == null || file.isEmpty()能拦住"有字段但没内容"的情况 - 原始文件名必须存在——后面要靠文件名后缀来判断文件类型
- 文件类型必须是系统支持的——通过
DocumentFileTypeEnum.fromFileName()根据后缀匹配,匹配不上就直接拒绝
第二步:读取文件字节 + 生成文档 ID
// 先把 MultipartFile 一次性读取成字节数组,后续存储、文件大小计算都复用这份内存快照。
byte[] fileBytes = getFileBytes(file);
// 文档主键在上传阶段就提前生成,后面对象存储路径、数据库主表、任务表都要围绕这个 documentId 建关联。
Long documentId = uidGenerator.getUid();
getFileBytes() 是一个私有方法,把 MultipartFile 的输入流一次性读成 byte[]:
/**
* 从 Spring 的 {@link MultipartFile} 中提取原始字节数组。
* <p>
* 上传链路后续要做对象存储上传、文件大小统计、任务日志记录等操作,
* 因此这里先把输入流稳定地读取成字节数组,避免在多个地方重复读取 multipart 流。
* </p>
*/
private byte[] getFileBytes(MultipartFile file) {
try {
return file.getBytes();
}
catch (IOException exception) {
// 读取文件体失败属于上传链路的底层 IO 异常,统一包装成文档存储失败,便于前端和日志统一识别。
throw new SuperAgentFrameException(DocumentManageCode.DOCUMENT_STORAGE_FAILED.getCode(),
"读取上传文件内容失败: " + exception.getMessage(), exception);
}
}
为什么要先读成 byte[]?因为后面对象存储上传要用、文件大小计算也要用,一次性读出来复用比多次读流更稳定。
documentId 通过百度的 UidGenerator 提前生成,这个 ID 后面会贯穿整条链路——对象存储路径里有它,数据库主表的主键是它,任务表也要关联它。
第三步:上传文件到 MinIO 对象存储
// 原始文件先上传到对象存储,拿到 bucket/object/url 等定位信息后,再写入文档主表。
StoredObjectInfo storedObjectInfo = storageService.uploadOriginalFile(
documentId, originalFileName, fileBytes, file.getContentType());
这一步调用了 MinioDocumentStorageService.uploadOriginalFile(),我们深入看看这个方法的实现:
@Override
public StoredObjectInfo uploadOriginalFile(Long documentId, String originalFileName,
byte[] bytes, String contentType) {
// 原始文件路径按"前缀/文档ID/时间戳-原始文件名"组织,便于同一文档的多次上传版本区分与追踪。
String objectName = properties.getMinio().getObjectPrefix() + "/"
+ documentId + "/" + System.currentTimeMillis() + "-" + originalFileName;
upload(objectName, bytes, contentType);
// 返回对象定位信息给上层,供文档主表持久化保存。
return new StoredObjectInfo(properties.getMinio().getBucketName(), objectName,
buildObjectUrl(objectName));
}
这里做了两件事:
- 拼接对象存储路径:格式是
前缀/documentId/时间戳-原始文件名,比如doc-original/123456/1714000000000-报告.pdf。加时间戳是为了同一文档多次上传时不会互相覆盖 - 调用
upload()执行真正的上传,然后把 bucket、objectName、objectUrl 封装成StoredObjectInfo返回
再往下看 upload() 方法:
/**
* 执行一次真正的 MinIO 对象上传。
* <p>
* 这个方法是上传链路在存储层的核心入口,负责两件事:
* 1. 确保 bucket 已存在;
* 2. 把字节数组以指定 contentType 写入 MinIO。
* </p>
*/
private void upload(String objectName, byte[] bytes, String contentType) {
try {
// 先保证桶存在,再写对象;否则首次部署或新环境下第一次上传会直接失败。
ensureBucketExists();
minioClient.putObject(
PutObjectArgs.builder()
.bucket(properties.getMinio().getBucketName())
.object(objectName)
// 如果上层没能提供明确的 MIME 类型,就回退到通用二进制类型。
.contentType(StrUtil.isNotBlank(contentType) ? contentType : "application/octet-stream")
.stream(new ByteArrayInputStream(bytes), bytes.length, -1)
.build()
);
}
catch (Exception exception) {
// MinIO SDK 抛出的各种底层异常统一包装成业务异常,方便上层按"存储失败"统一处理。
throw new SuperAgentFrameException(DocumentManageCode.DOCUMENT_STORAGE_FAILED.getCode(),
"上传 MinIO 文件失败: " + exception.getMessage(), exception);
}
}
ensureBucketExists() 会先检查桶是否存在,不存在就自动创建:
private void ensureBucketExists() throws Exception {
if (!bucketExists()) {
// 只有在桶不存在时才创建,避免每次上传都重复发起创建请求。
minioClient.makeBucket(MakeBucketArgs.builder()
.bucket(properties.getMinio().getBucketName()).build());
}
}
private boolean bucketExists() throws Exception {
String bucketName = properties.getMinio().getBucketName();
return minioClient.bucketExists(BucketExistsArgs.builder().bucket(bucketName).build());
}
最后 buildObjectUrl() 拼出文件的完整访问地址:
private String buildObjectUrl(String objectName) {
String endpoint = properties.getMinio().getEndpoint();
if (endpoint.endsWith("/")) {
// 去掉结尾的"/",避免下面手工拼接 bucket/objectName 时出现重复分隔符。
endpoint = endpoint.substring(0, endpoint.length() - 1);
}
return endpoint + "/" + properties.getMinio().getBucketName() + "/" + objectName;
}
返回值 StoredObjectInfo
StoredObjectInfo是一个简单的数据载体,包含三个字段:
bucketName:MinIO 桶名objectName:对象在桶内的完整路径objectUrl:拼接好的完整访问 URL
第四步:构建文档主记录
文件上传到 MinIO 之后,接下来要在内存中组装文档主记录对象。注意这一步只是构建对象、设置字段,还没有真正写入数据库——入库操作统一放在后面的手动事务中执行。
// 文档主表记录的是"这份文档当前是什么状态、文件存在哪里、后续应该进入什么处理阶段"。
SuperAgentDocument document = new SuperAgentDocument();
document.setId(documentId);
// 如果前端传了 documentName,就优先使用业务展示名;否则回退到上传时的原始文件名。
document.setDocumentName(StrUtil.isNotBlank(dto.getDocumentName())
? dto.getDocumentName() : originalFileName);
document.setOriginalFileName(originalFileName);
document.setFileType(fileType.getCode());
document.setMimeType(file.getContentType());
document.setFileSize((long) fileBytes.length);
document.setStorageType(DocumentStorageTypeEnum.MINIO.getCode());
document.setBucketName(storedObjectInfo.getBucketName());
document.setObjectName(storedObjectInfo.getObjectName());
document.setObjectUrl(storedObjectInfo.getObjectUrl());
// 上传成功后并不代表解析完成,因此解析状态先标记为 PARSING,表示"已入队等待后续处理链继续推进"。
document.setParseStatus(DocumentParseStatusEnum.PARSING.getCode());
// 策略和索引阶段都还没有真正执行,所以这里先放到"等待推荐 / 等待构建"的初始态。
document.setStrategyStatus(DocumentStrategyStatusEnum.WAIT_RECOMMEND.getCode());
document.setIndexStatus(DocumentIndexStatusEnum.WAIT_BUILD.getCode());
// 字符数和 token 数需要等解析文本产出后再统计,因此上传阶段先置 0。
document.setCharCount(0);
document.setTokenCount(0);
// 这些字段属于上传时就能确定的业务标签,用于后续检索范围、文档分类、运营筛选等场景。
document.setKnowledgeScopeCode(StrUtil.trimToNull(dto.getKnowledgeScopeCode()));
document.setKnowledgeScopeName(StrUtil.trimToNull(dto.getKnowledgeScopeName()));
document.setBusinessCategory(StrUtil.trimToNull(dto.getBusinessCategory()));
document.setDocumentTags(StrUtil.trimToNull(dto.getDocumentTags()));
document.setStatus(BusinessStatus.YES.getCode());
这段代码信息量比较大,我们拆开来看几个关键点:
文档名称的取值逻辑:优先用前端传的 documentName(业务展示名),没传就用原始文件名。这样前端可以给文档起一个更友好的名字,比如"2024年Q1财报",而不是"report_2024q1_final_v3.pdf"。
三个状态字段的初始值:
- parseStatus = PARSING:表示"已经进入解析队列,等待异步处理"
- strategyStatus = WAIT_RECOMMEND:策略推荐还没开始
- indexStatus = WAIT_BUILD:索引构建还没开始
这三个状态会随着后续异步链路的推进逐步更新。
charCount 和 tokenCount 先置 0:因为这两个值要等文档解析完、拿到纯文本之后才能统计,上传阶段还不知道。
第五步:构建任务记录
文档主记录组装好之后,紧接着要构建一条"解析路由任务"记录。和文档主记录一样,这一步也只是在内存中组装对象,真正的入库操作放在后面的手动事务中统一执行。
// 上传完成后,不会在当前请求线程里直接做重解析,而是创建一条任务记录交给异步链路处理。
Long taskId = uidGenerator.getUid();
SuperAgentDocumentTask task = new SuperAgentDocumentTask();
task.setId(taskId);
task.setDocumentId(documentId);
// 上传之后的第一类任务是"解析路由",它负责决定后续解析和策略推荐该怎么走。
task.setTaskType(DocumentTaskTypeEnum.PARSE_ROUTE.getCode());
task.setTaskStatus(DocumentTaskStatusEnum.NEW.getCode());
// 当前阶段先记为 FILE_UPLOAD,表示这条任务刚从上传阶段移交出来。
task.setCurrentStage(DocumentTaskStageEnum.FILE_UPLOAD.getCode());
// operatorId 允许为空;为空时表示系统触发,否则表示由具体用户触发。
Long operatorId = parseOptionalLong(dto.getOperatorId());
task.setTriggerSource(resolveTriggerSource(operatorId));
task.setRetryCount(0);
task.setStatus(BusinessStatus.YES.getCode());
几个关键字段说明:
taskType = PARSE_ROUTE:这是上传后的第一类任务,叫"解析路由"。它的职责是决定后续的解析和策略推荐该怎么走taskStatus = NEW:刚创建,还没开始执行currentStage = FILE_UPLOAD:标记这条任务是从"文件上传"阶段移交出来的triggerSource:通过resolveTriggerSource()判断——有 operatorId 就是用户触发,没有就是系统触发
resolveTriggerSource() 和 parseOptionalLong() 这两个辅助方法也值得看一眼:
/**
* 根据操作人 ID 判断任务触发来源。
*/
private Integer resolveTriggerSource(Long operatorId) {
return operatorId == null
? DocumentTriggerSourceEnum.SYSTEM.getCode()
: DocumentTriggerSourceEnum.USER.getCode();
}
/**
* 解析可选的字符串 ID。
* <p>
* 上传接口中的一些字段允许不传,因此这里采取"宽松解析"策略:
* 空白串、非数字、非正数都统一视为 null,而不是直接抛异常。
* </p>
*/
private Long parseOptionalLong(String rawValue) {
if (StrUtil.isBlank(rawValue)) {
return null;
}
try {
// 只有正整数才被视为有效 ID;0 或负数在业务上不具备合法主键语义。
Long value = Long.valueOf(rawValue.trim());
return value > 0 ? value : null;
}
catch (NumberFormatException exception) {
// 可选字段解析失败时不打断主流程,直接按未传处理。
return null;
}
}
parseOptionalLong() 的设计思路是"宽松解析":空白串、非数字、0 或负数都当 null 处理,不会因为前端传了个奇怪的值就把整个上传请求打挂。
第六步:手动事务提交——入库 + 任务日志
前面两步只是在内存中组装好了文档主记录和任务记录,还没有真正写入数据库。接下来通过 TransactionTemplate 手动开启事务,把文档主表、任务表、任务日志的写入放在同一个事务里执行:
// 使用 TransactionTemplate 手动控制事务边界,而不是在方法上加 @Transactional。
// 这样做的核心目的是:确保事务提交之后再发送 Kafka 消息。
// 如果用 @Transactional,Kafka 发送和数据库操作在同一个事务上下文中,
// 存在"消息已发出但事务还没提交"的时间窗口——消费方查数据库可能查不到数据。
DocumentUploadVo uploadVo = transactionTemplate.execute(status -> {
// 文档主表入库:这是整条上传链路的核心写入,后续任务和日志都要引用 documentId
documentMapper.insert(document);
// 任务表入库:创建一条"解析路由"任务,交给异步消费方处理
taskMapper.insert(task);
// 任务日志记录:把"上传已完成并已投递异步链路"这件事显式记录下来,
// 方便后续排查问题时回看时间线
taskLogService.saveLog(taskId, documentId,
DocumentTaskStageEnum.FILE_UPLOAD.getCode(),
DocumentTaskEventTypeEnum.COMPLETE.getCode(),
DocumentLogLevelEnum.INFO.getCode(),
resolveOperatorType(operatorId),
operatorId,
"文件上传完成,已进入解析与策略推荐队列。",
Map.of("originalFileName", originalFileName, "fileSize", fileBytes.length));
// 在事务内部就把返回值组装好,确保返回给前端的数据和数据库中的数据完全一致
return new DocumentUploadVo(documentId, taskId, document.getDocumentName(),
document.getParseStatus(), document.getStrategyStatus(), document.getIndexStatus());
});
这里最关键的设计决策是:用 TransactionTemplate 替代了方法级的 @Transactional 注解。
为什么要这样改?因为原来的写法是在 upload() 方法上加 @Transactional,整个方法体都在事务范围内。这意味着 Kafka 消息发送也在事务提交之前执行——如果消息发得快、消费方处理得也快,消费方去查数据库时事务可能还没提交,就会查不到刚插入的文档和任务记录。
用 TransactionTemplate 之后,事务的边界被精确控制在 execute() 的 lambda 内部。lambda 执行完毕、execute() 方法返回时,事务就已经提交了。后面再发 Kafka 消息,消费方查数据库时数据一定已经落盘。
为什么不用 @Transactional + TransactionSynchronizationManager?
另一种常见做法是保留
@Transactional,然后用TransactionSynchronizationManager.registerSynchronization()在事务提交后回调发送 Kafka。这种方式也能解决问题,但代码更绕,而且回调里如果抛异常不容易处理。直接用TransactionTemplate把事务边界收窄,代码更直观。
第七步:事务提交后发送 Kafka 消息
事务提交成功后,再把任务投递到 Kafka,交给异步消费方去做后续的解析和策略推荐:
// 走到这里说明 transactionTemplate.execute() 已经正常返回,
// 即文档主表、任务表、任务日志都已经事务提交成功。
// 此时再发 Kafka 消息,消费方查数据库一定能查到完整的上下文数据。
kafkaProducer.sendParseRoute(new DocumentParseRouteMessage(documentId, taskId));
Kafka 发送放在事务外的权衡
这样做保证了"数据库先落盘,消息再发出",但也引入了一个新的边界问题:如果事务提交成功但 Kafka 发送失败,数据库里有记录但异步链路不会启动。不过这种情况可以通过补偿机制(定时扫描状态为 PARSING 但长时间没有被消费的文档)来兜底,比"消息发出但数据库没提交"要好处理得多。
DocumentParseRouteMessage 很简单,就两个字段:
@Data
@NoArgsConstructor
@AllArgsConstructor
public class DocumentParseRouteMessage {
private Long documentId;
private Long taskId;
}
DocumentKafkaProducer.sendParseRoute() 的实现:
/**
* 发送"解析路由"消息。
* <p>
* 这条消息由文档上传链路触发,代表"原始文件已经上传并入库,可以进入异步解析与策略推荐阶段"。
* </p>
*/
public void sendParseRoute(DocumentParseRouteMessage message) {
send(SpringUtil.getPrefixDistinctionName() + "-"
+ properties.getKafka().getParseTopic(),
String.valueOf(message.getDocumentId()), message);
}
topic 名称是通过 SpringUtil.getPrefixDistinctionName() 加上配置的 topic 名拼接而成的,这样不同环境(dev、test、prod)的 topic 不会互相干扰。消息的 key 用的是 documentId,这样同一文档的消息会落到同一个 partition,保证顺序性。
底层的 send() 方法:
/**
* 执行一次底层 Kafka 消息发送。
* <p>
* 这里统一完成对象序列化和同步发送确认;如果发送失败,则直接抛出业务异常,
* 让上层知道消息并没有真正进入异步链路。
* </p>
*/
private void send(String topic, String key, Object message) {
try {
// 先把消息对象序列化为 JSON,消费方再按对应消息类型反序列化。
String payload = objectMapper.writeValueAsString(message);
// 这里使用 get() 等待发送结果,确保上传接口返回成功前,消息已经被 broker 接收确认。
kafkaTemplate.send(topic, key, payload).get();
} catch (Exception exception) {
throw new SuperAgentFrameException(DocumentManageCode.KAFKA_SEND_FAILED.getCode(),
"Kafka 消息发送失败: " + exception.getMessage(), exception);
}
}
注意这里用了 .get() 做同步等待——kafkaTemplate.send() 本身是异步的,返回一个 Future,调用 .get() 会阻塞当前线程直到 broker 确认收到消息。这样做的好处是:如果 Kafka 发送失败,上传接口会直接返回错误,而不是"假装成功"但消息其实丢了。
第八步:返回上传结果
// uploadVo 在事务内部就已经组装好了,这里直接返回给 Controller 层
return uploadVo;
返回给前端的 DocumentUploadVo 包含:
@Data
@NoArgsConstructor
@AllArgsConstructor
public class DocumentUploadVo {
private Long documentId;
private Long taskId;
private String documentName;
private Integer parseStatus;
private Integer strategyStatus;
private Integer indexStatus;
}
只返回了文档 ID、任务 ID、文档名称和三个状态码,不会暴露 MinIO 的存储路径、bucket 名称这些内部细节。前端拿到 documentId 和 taskId 之后,就可以用它们去轮询文档的处理进度了。
数据流转结构图
最后用一张图来总结上传链路中各个数据对象之间的关系:
小结
整个上传链路可以用一句话概括:校验文件 → 存到 MinIO → 组装文档和任务对象 → TransactionTemplate 手动事务提交(入库 + 日志)→ 事务提交后发 Kafka 消息 → 返回结果。
上传接口本身是同步的,但它做的事情其实是"为异步链路做好准备"——文件存好了、数据库记录建好了、任务也创建好了,最后通过 Kafka 消息通知异步消费方:"这份文档准备好了,你可以开始解析了"。这里特别值得注意的是事务边界的设计:通过 TransactionTemplate 把数据库操作和 Kafka 发送分开,确保消费方查数据库时数据一定已经落盘。
至于 Kafka 消费方拿到消息之后怎么做解析、怎么做策略推荐,那就是下一篇文档要讲的内容了。
企业级项目导航:⬅️ 02-一定得用Kafka吗 | 03-上传接口与文档主记录创建 | ➡️ 04-异步解析入口
💬 评论