上传接口与文档主记录创建

这篇文档要讲的是:用户上传一份文档后,后端到底做了哪些事情?从 Controller 接到请求开始,一直到最后把消息丢进 Kafka 队列,整条链路我们一步步拆开来看。

先上一张总览流程图,有个整体印象之后再逐段看源码。

上传链路总览

FrDssNgoaGj3xM5amhckSNBLD8SS-5428612a

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 名称这些内部细节。前端拿到 documentIdtaskId 之后,就可以用它们去轮询文档的处理进度了。

数据流转结构图

最后用一张图来总结上传链路中各个数据对象之间的关系:

FsKIkYywPG4G_4H6n2a6HpuTVv9p-3f0972d7

小结

整个上传链路可以用一句话概括:校验文件 → 存到 MinIO → 组装文档和任务对象 → TransactionTemplate 手动事务提交(入库 + 日志)→ 事务提交后发 Kafka 消息 → 返回结果

上传接口本身是同步的,但它做的事情其实是"为异步链路做好准备"——文件存好了、数据库记录建好了、任务也创建好了,最后通过 Kafka 消息通知异步消费方:"这份文档准备好了,你可以开始解析了"。这里特别值得注意的是事务边界的设计:通过 TransactionTemplate 把数据库操作和 Kafka 发送分开,确保消费方查数据库时数据一定已经落盘。

至于 Kafka 消费方拿到消息之后怎么做解析、怎么做策略推荐,那就是下一篇文档要讲的内容了。


企业级项目导航:⬅️ 02-一定得用Kafka吗 | 03-上传接口与文档主记录创建 | ➡️ 04-异步解析入口