别把文件上传写成一个 save:可靠上传接口的设计与实现
开发中的痛点
在头像、附件、合同、简历等业务中,文件上传接口看起来并不复杂:接收 MultipartFile,把文件写入对象存储,再保存一条数据库记录。
最直接的实现通常是:
@Transactional
public Long upload(MultipartFile file, Long userId) {
// 前置其他效验
........
// 效验文件类型
storage.verity(file);
// 先上传对象,再保存业务记录
String objectKey = storage.put(file);
// 如果业务有解析文本需求
parserService.parser(file);
FileRecord record = FileRecord.active(userId, objectKey);
fileRecordMapper.insert(record);
return record.getId();
}
这段代码能够通过正常流程测试,但它隐藏了几个关键问题:
- 事务边界错误:数据库事务不能自动回滚已经写入
OSS、S3或文件系统的对象。 - 并发配额失效:如果业务有要求,如简历最多创建5份,两个请求都可能在
count < 5时通过检查,最终插入第 6 条记录。 - 失败状态缺失:服务在上传后、保存前宕机时,既没有完整记录,也没有可靠的清理依据。
- 长事务占用资源:把网络上传、文本解析放进事务,会长时间占用数据库连接和锁。
- 文件校验不足:客户端文件名和
Content-Type都可以伪造,不能直接作为可信依据。
日常开发中,常见的处理方式有三种。
方案一:先上传文件,再保存数据库
文件成功写入对象存储后,再插入业务记录。
优点:
- 代码简单。
- 数据库中通常不会出现指向未上传对象的有效记录。
缺点:
- 数据库写入失败时会留下孤儿文件。
- 服务在两步之间宕机后,缺少可追踪的补偿依据。
方案二:先保存数据库,再上传文件
先创建记录,再执行对象存储上传。
优点:
- 上传开始前已经存在可追踪记录。
- 可以记录处理进度和失败原因。
缺点:
- 如果记录直接标记为可用,业务方可能读取到尚不存在的文件。
- 上传失败后必须修复状态,否则会形成脏数据。
方案三:状态预占、短事务与失败补偿
先用短事务创建 UPLOADING 记录并占用配额;事务提交后上传文件;成功后再用短事务将状态推进为 ACTIVE。与此同时,提前保存一条延迟执行的清理任务,负责兜底删除未完成上传遗留的对象。
这套方案增加了状态和补偿逻辑,却能够明确回答三个问题:
- 当前上传进行到了哪一步?
- 任意一步失败后由谁恢复?
- 服务在任意位置宕机后,系统如何自行收敛?
真正的问题不是“文件能不能传上去”,而是:
如何让数据库记录、对象存储、并发配额和失败补偿最终保持一致?
阅读须知
- 适用场景:后端代理上传中小型文件,并需要保存业务记录、限制数量或执行后续处理。
- 前置知识:了解
Spring Boot、数据库事务、对象存储和基本并发问题。- 语言版本:
Java 17+。- 技术环境:
Spring Boot 3.x、关系型数据库、任意兼容的对象存储。- 方案边界:超大文件、断点续传和高带宽场景更适合客户端预签名直传,文末会单独说明。
本文优先讲清楚 为什么需要 → 如何拆分事务 → 如何控制并发 → 如何补偿失败 → 如何验证边界。
本文要解决什么
本文最终希望实现:
- 对文件大小、扩展名、服务端检测的媒体类型和文件名进行分层校验。
- 同一用户不能重复上传相同文件。
- 用户最多保留指定数量的有效文件,并发上传也不能突破限制。
- 上传期间不持有数据库长事务。
- 上传失败、服务宕机和补偿失败后均可恢复。
- 文件记录具有明确且可校验的状态转换。
最终调用方式保持简单:
// Controller 只调用统一上传入口,不参与事务和补偿编排
UploadResult result = uploadService.upload(userId, file, metadata);
复杂度应封装在服务内部,而不是转嫁给 Controller。
整体设计
本文方案由以下组件组成:
| 组件 | 主要职责 |
|---|---|
FileValidator |
校验大小、扩展名、文件签名和内容类型 |
FileHashService |
流式计算 SHA-256,用于业务查重 |
FileStorage |
生成对象键、上传、检测和删除文件 |
UploadTxService |
执行预占、确认和失败标记三个短事务 |
CleanupIntentRepository |
在外部上传前保存延迟清理意图 |
CleanupWorker / MQ Consumer |
通过任务租约或消息确认机制执行幂等删除 |
UploadRepairJob |
修复超时状态和补偿任务 |
完整链路如下:
文件表只保存上传业务状态,不保存后台任务的执行状态:
UPLOADING 不是一个展示文案,而是并发控制和故障恢复都依赖的业务状态。清理任务是否正在执行、能否重试,则由数据库任务表或消息队列单独管理。
实现前需要理解的关键机制
1. MultipartFile 中的数据并不天然可信
getOriginalFilename() 返回的是客户端提供的文件名。它可能为空,也可能包含 ..、路径片段或特殊字符。Spring 官方文档明确建议不要直接使用该值作为存储文件名,而应生成新的唯一名称。MultipartFile 官方文档
同样,Content-Type 来自请求,也可以被伪造。可靠校验至少应包含:
- 业务允许的文件大小。
- 扩展名白名单。
- 服务端检测的媒体类型。
- 文件头或文件签名。
- 必要时进行病毒扫描、沙箱分析或
CDR处理。
这些校验需要组合使用,不能把其中任意一项当成绝对可信证据。OWASP 的文件上传指南还建议使用应用生成的文件名、限制上传权限,并将文件放在站点根目录之外或独立存储服务中。OWASP File Upload Cheat Sheet
2. 数据库事务不能覆盖普通对象存储
常见的 @Transactional 管理的是数据库连接所参与的事务。对象存储的 putObject 和 deleteObject 通常不是这个本地事务的资源,因此数据库回滚不会撤销已经完成的上传。
这意味着下面的代码即使抛出异常,文件也可能仍然存在:
@Transactional
public void wrongUpload(MultipartFile file) {
// 对象存储不会跟随数据库事务回滚
storage.put(file);
fileRecordMapper.insert(buildRecord());
throw new IllegalStateException("rollback database only");
}
解决方向不是把事务包得更大,而是缩短数据库事务,并为外部副作用设计补偿和对账。
3. “查询后插入”不是原子操作
假设业务要求最多创建 5 个记录(比如,招聘平台一般简历不会允许你创建很多份),当前已有 4 个。两个请求可能同时执行:
@Transactional
public Long upload(MultipartFile file, Long userId) {
// 前置其他效验
........
Integer count = lambdaQuery
.eq(Resume::getUserId,userId)
.count();
// 注意这里仅为演示,所以使用了魔法值。实际开发中不要使用魔法值
if(count >= 5 ){
throw new businessException(ErrorCode.RESUME_MOST_IS_FIVE);
}
// 其他操作
......
}
请求 A:count = 4,通过
请求 B:count = 4,通过
请求 A:insert,数量变为 5
请求 B:insert,数量变为 6
把 count 和 insert 写在同一事务里仍不够。普通隔离级别下,两个事务依然可能读取到相同数量。必须让同一用户的“检查 + 预占”串行执行,或者将配额改造成数据库可原子更新的计数器。
4. 状态转换必须检查受影响行数
上传成功后的确认不能是不带前置状态的覆盖更新:
-- 只允许仍处于上传中的记录完成确认
UPDATE file_upload
SET upload_status = 'ACTIVE'
WHERE id = ?
AND upload_status = 'UPLOADING';
只有受影响行数等于 1,才说明当前请求仍然拥有推进状态的资格。如果记录已被超时任务改为 UPLOAD_FAILED,确认操作必须失败,不能重新把它改回 ACTIVE。
本文示例中的 UploadException 是非受检异常。Spring 默认只对 RuntimeException 及其子类和 Error 回滚;如果业务异常继承受检异常,则必须显式配置 rollbackFor。Spring 回滚规则
public class UploadException extends RuntimeException {
// 业务异常需要保留原始原因,便于日志定位失败阶段
public UploadException(String message) {
super(message);
}
public UploadException(String message, Throwable cause) {
super(message, cause);
}
}
从设计到实现
第一步:定义状态与数据契约
上传状态应与文件解析、审核等后续业务状态分开,避免一个枚举同时承担两套生命周期。
public enum UploadStatus {
UPLOADING, // 已预占,文件尚未确认
ACTIVE, // 文件已上传并可被业务使用
UPLOAD_FAILED // 上传失败,由独立清理任务处理对象
}
预占阶段需要保存足够的恢复信息:
// 预占事务需要持久化的文件元数据
public record FileMetadata(
String originalName,
String objectKey,
String mediaType,
String sha256,
long size
) {
}
// 预占成功后交给事务外上传流程使用
public record UploadReservation(
Long recordId,
Long cleanupTaskId,
String objectKey
) {
}
数据库记录至少包含:
- 所有者
owner_id。 - 展示用原始文件名
original_name。 - 应用生成的对象键
object_key。 - 文件哈希、大小和服务端识别的媒体类型。
- 上传状态、失败原因和时间字段。
- 上传截止时间
upload_deadline_at。 - 逻辑删除或可见性字段。
对象键应由服务端生成,例如:
resume/1029384756/7a04f625d74f4fd3a2a3c1b3f367e311.pdf
原始文件名只用于展示,不参与对象存储路径拼接。
第二步:完成文件校验与哈希计算
下面的代码只展示核心逻辑。实际项目可以使用 Apache Tika、专用文件解析库或安全扫描服务完成内容识别。
@Component
public class FileValidator {
private static final long MAX_SIZE = 10 * 1024 * 1024L;
private static final Set<String> ALLOWED_EXTENSIONS = Set.of("pdf", "doc", "docx");
private static final Set<String> ALLOWED_MEDIA_TYPES = Set.of(
"application/pdf",
"application/msword",
"application/vnd.openxmlformats-officedocument.wordprocessingml.document"
);
public ValidatedFile validate(MultipartFile file) {
// 第一层:拒绝空文件和超限文件
if (file == null || file.isEmpty() || file.getSize() > MAX_SIZE) {
throw new UploadException("文件为空或大小超限");
}
// 第二层:校验清理后的展示文件名和扩展名
String displayName = sanitizeDisplayName(file.getOriginalFilename());
String extension = extractExtension(displayName);
if (!ALLOWED_EXTENSIONS.contains(extension)) {
throw new UploadException("文件扩展名不受支持");
}
// 第三层:根据文件内容检测服务端媒体类型
String detectedType = detectMediaType(file, displayName);
if (!ALLOWED_MEDIA_TYPES.contains(detectedType)) {
throw new UploadException("文件内容类型不受支持");
}
// 哈希只用于查重,不替代安全扫描
String sha256 = calculateSha256(file);
return new ValidatedFile(displayName, extension, detectedType, sha256, file.getSize());
}
}
还应在容器层设置请求上限,避免超大请求完整进入业务代码:
spring:
servlet:
multipart:
# 在请求进入业务代码前限制体积
max-file-size: 10MB
max-request-size: 12MB
容器限制负责尽早拒绝请求,业务校验负责按文件类型、用户等级等规则给出精确判断。相关配置项可参考 Spring Boot 应用属性文档。
计算哈希时不要直接调用 getBytes() 将整个文件装入堆内存,应使用流式摘要:
public String calculateSha256(MultipartFile file) {
// 使用流式摘要,避免将整个文件装入堆内存
try (InputStream input = file.getInputStream()) {
MessageDigest digest = MessageDigest.getInstance("SHA-256");
byte[] buffer = new byte[8192];
int length;
while ((length = input.read(buffer)) != -1) {
// 只处理本次实际读取到的字节
digest.update(buffer, 0, length);
}
return HexFormat.of().formatHex(digest.digest());
} catch (IOException | NoSuchAlgorithmException e) {
throw new UploadException("计算文件哈希失败", e);
}
}
文件哈希适合用于业务查重,但它不是病毒扫描结果,也不能证明文件来自可信用户。
第三步:用短事务完成加锁、检查与预占
预占事务需要依次完成:
- 锁定当前用户的文件变更过程。
- 重新检查相同哈希是否已存在。
- 检查当前有效文件数量。
- 插入
UPLOADING记录。 - 插入延迟清理任务。
下面使用 PostgreSQL 事务级 Advisory Lock 演示。该锁会在事务结束时自动释放,适合没有天然“用户文件守卫行”的场景。PostgreSQL Advisory Lock 官方文档
@Mapper
public interface FileUploadMapper {
// 第一段锁键划分业务域,第二段锁键标识当前用户
@Select("SELECT pg_advisory_xact_lock(#{namespace}, #{ownerSlot})")
void lockOwner(int namespace, int ownerSlot);
}
锁键生成规则应集中管理,避免上传、删除和定时修复各自复制命名空间常量:
public final class UploadLockKey {
// PostgreSQL 双 int 锁键的第一段固定为文件业务命名空间
public static final int NAMESPACE = 0x46494C45;
private UploadLockKey() {
}
public static int ownerSlot(long ownerId) {
// 将 64 位用户 ID 稳定映射到第二段锁键
return Long.hashCode(ownerId);
}
}
不要用 ownerId ^ NAMESPACE 实现命名空间隔离:异或只是把一个 64 位值映射为另一个 64 位值,不同业务仍可能生成同一个锁键。双键版本显式保留业务域;ownerSlot 即使发生哈希碰撞,也只会让两个用户暂时额外串行,属于性能降级,不会导致并发约束失效,况且哈希碰撞的概率本来就很低,可以接受这个小损失(除非你的用户量极大且对并发敏感,则需对其进行优化升级)。
核心预占逻辑位于独立事务服务中:
@Service
@RequiredArgsConstructor
public class UploadTxService {
private static final int MAX_FILE_COUNT = 5;
private final FileUploadRepository fileRepository;
private final CleanupTaskRepository cleanupRepository;
@Transactional
public UploadReservation reserve(Long ownerId, FileMetadata metadata) {
// 串行化同一用户的“检查 + 预占”过程
fileRepository.lockOwner(
UploadLockKey.NAMESPACE,
UploadLockKey.ownerSlot(ownerId)
);
// 锁内重新查重,事务外预检查只能作为优化
if (fileRepository.existsAliveByHash(ownerId, metadata.sha256())) {
throw new UploadException("相同文件已经存在");
}
// UPLOADING 记录同样占用有效配额
long count = fileRepository.countAliveByOwner(ownerId);
if (count >= MAX_FILE_COUNT) {
throw new UploadException("文件数量已达到上限");
}
// 写入明确截止时间,供定时修复判断是否超时
Instant uploadDeadline = Instant.now().plus(Duration.ofMinutes(15));
FileUpload record = FileUpload.uploading(ownerId, metadata, uploadDeadline);
fileRepository.insert(record);
// 外部上传前先保存补偿意图
CleanupTask cleanupTask = CleanupTask.pending(
record.getId(),
metadata.objectKey(),
Instant.now().plus(Duration.ofMinutes(30))
);
cleanupRepository.insert(cleanupTask);
return new UploadReservation(
record.getId(),
cleanupTask.getId(),
metadata.objectKey()
);
}
}
事务中没有文件上传、文本解析或网络请求,因此锁持有时间只覆盖必要的数据库操作。
UploadLockKey 的目的不是加密,而是统一生成同一业务资源的锁键。所有会改变同一用户文件集合的入口,都必须使用相同的锁规则。否则上传接口加了锁,另一个批量导入接口却绕过锁,配额仍可能失效。
Advisory Lock 能协调连接到同一套 PostgreSQL 的多个服务实例,并不只适用于单机。但如果未来拆分后各服务使用不同数据库,该锁就不能跨库生效,需要改用集中式配额服务、共享 Redis 锁或其他分布式协调方案。
对于重复文件,还可以在数据库增加防御性约束。以 PostgreSQL 为例:
-- 防止同一用户保留两条相同内容的有效记录
CREATE UNIQUE INDEX uk_file_owner_hash_alive
ON file_upload (owner_id, file_hash)
WHERE is_deleted = false;
业务锁负责“最多 5 份”的集合约束,唯一索引负责兜住“同一有效文件只能存在一份”的单行约束,两者职责不同。
第四步:在事务外上传,再用短事务确认
应用服务只负责编排流程:
@Service
@RequiredArgsConstructor
public class FileUploadService {
private final FileValidator validator;
private final FileStorage storage;
private final UploadTxService txService;
public UploadResult upload(Long ownerId, MultipartFile file) {
// 先完成安全校验、哈希计算和对象键生成
ValidatedFile validated = validator.validate(file);
String objectKey = storage.generateObjectKey(ownerId, validated.extension());
FileMetadata metadata = validated.toMetadata(objectKey);
// 短事务预占配额并创建清理任务
UploadReservation reservation = txService.reserve(ownerId, metadata);
try {
// 网络上传位于数据库事务之外
storage.put(
reservation.objectKey(),
file.getInputStream(),
file.getSize(),
validated.mediaType()
);
// 上传成功后再通过条件更新激活记录
txService.confirm(reservation.recordId(), reservation.cleanupTaskId());
return new UploadResult(reservation.recordId(), validated.originalName());
} catch (Exception uploadError) {
// 标记失败异常不能覆盖最初的上传异常
safeMarkFailed(reservation.recordId(), uploadError);
throw new UploadException("文件上传失败", uploadError);
}
}
private void safeMarkFailed(Long recordId, Exception cause) {
try {
// 使用独立事务释放配额并保留失败状态
txService.markFailed(recordId, conciseReason(cause));
} catch (Exception markError) {
log.error("标记上传失败异常,等待修复任务处理。recordId={}", recordId, markError);
}
}
}
成功确认必须和清理任务失效处于同一数据库事务:
@Transactional
public void confirm(Long recordId, Long cleanupTaskId) {
// 只有 UPLOADING 状态能够被当前请求确认
boolean transitioned = fileRepository.transition(
recordId,
UploadStatus.UPLOADING,
UploadStatus.ACTIVE
);
if (!transitioned) {
throw new UploadException("上传记录状态已变化,不能确认");
}
// 与状态确认处于同一事务,避免有效文件被补偿任务删除
cleanupRepository.markIgnored(cleanupTaskId);
}
上面展示的是数据库任务表版本。如果采用 Outbox + MQ,同一位置应把尚未发布的 Outbox 事件标记为 CANCELLED;消息若已经发布,消费者会在看到文件已是 ACTIVE 后直接确认,不执行删除。两种实现都必须保持“激活记录”和“撤销清理意图”位于同一数据库事务。
这里必须检查 transition() 的返回值。更新失败通常意味着记录不存在、已超时、已被清理任务接管,或者发生了重复确认。继续返回成功会让客户端以为文件已经可用,而数据库实际上没有完成状态转换。
失败标记使用独立事务:
@Transactional(propagation = Propagation.REQUIRES_NEW)
public void markFailed(Long recordId, String reason) {
// 条件更新失败时交给超时修复任务继续兜底
fileRepository.markFailedIfUploading(recordId, reason);
}
markFailedIfUploading 应在一条条件更新中完成以下操作:将 UPLOADING 改为 UPLOAD_FAILED、写入失败原因、更新审计字段,并将记录标记为逻辑删除。这样失败记录不会继续占用有效配额,用户也可以重新上传相同内容。
将该方法放在独立的 Spring Bean 中,可以避免同类内部调用绕过事务代理。即使失败标记自身失败,预先创建的清理任务和超时扫描仍能继续兜底。修复任务查询失败记录时,需要绕过逻辑删除的默认过滤条件。
第五步:为兜底清理选择可靠的执行机制
前面的预占事务已经在对象上传之前保存了清理意图,因此即使服务在上传后宕机,系统仍然知道“这个对象可能需要被删除”。
但有清理任务,不等于清理一定能够完成。
假设我们只是写一个普通定时任务:
扫描待清理记录
↓
标记任务处理中
↓
删除对象存储文件
↓
标记任务完成
任何一步都可能再次失败。
例如,工作节点刚领取任务就宕机:
任务已标记处理中
↓
服务宕机
↓
没有人继续执行删除
↓
任务可能永久卡住
或者对象已经删除成功,但服务在更新任务状态之前宕机:
对象删除成功
↓
服务宕机
↓
数据库仍认为任务没有完成
↓
任务之后可能被再次执行
再比如,对象存储暂时不可用:
执行删除
↓
网络异常
↓
删除失败
↓
如果没有重试机制,对象仍会永久残留
因此,真正可靠的补偿机制需要解决的不只是:
“我要删除哪个对象?”
还必须解决:
“谁负责执行?执行者宕机后谁接手?失败后如何重试?重复执行是否安全?”
这也是为什么不能把兜底清理简单理解成一个 @Scheduled 定时任务。
一个可靠的执行机制至少需要具备:
- 可重新领取:执行节点宕机后,任务不能永久卡在“处理中”。
- 失败重试:数据库、网络或对象存储临时异常后能够再次执行。
- 幂等执行:同一个删除任务重复执行不会产生错误结果。
- 执行权控制:多个实例不能无序地同时处理同一个任务。
- 最终可观测:超过重试次数后能够进入失败状态、死信队列或触发告警。
实现这些能力通常有两种思路:
方案一:数据库任务表 + 租约
↓
由数据库管理任务领取、租约过期和重试
方案二:Transactional Outbox + MQ
↓
由消息队列提供确认、重新投递和死信能力
两种方案共享以下约束:
- 文件表只保存
UPLOADING、ACTIVE、UPLOAD_FAILED。 - 删除前必须确认记录不是
ACTIVE,对象键也没有被其他有效记录引用。 deleteIfExists必须幂等,对象不存在也视为成功。- 任务可能重复执行,完成动作必须能够识别旧执行者。
解决了“任务如何可靠执行”之后,还需要解决另一个问题:**清理任务被成功调度,并不代表它拿到对象键后就可以直接删除。**从清理意图创建到真正执行之间,业务状态可能已经发生变化,因此执行前必须重新检查文件状态和引用关系。
清理任务判断“可以删除”
↓
上传请求将记录确认成 ACTIVE
↓
清理任务删除了正在使用的对象
清理命令至少需要携带业务记录和对象键,不能只保存一个 URL:
public record FileDeleteCommand(
Long recordId,
Long ownerId,
String objectKey
) {
// recordId 用于重新检查业务状态,objectKey 用于执行幂等删除
}
接下来分别讨论这数据库任务表 + 租约 和 Outbox + MQ两种方案。
方案一:数据库任务表 + 租约
没有消息队列,或者任务规模不大时,可以直接把清理任务保存在数据库。任务表负责执行状态,文件表不再增加 CLEANING/CLEANED:
public enum CleanupTaskStatus {
PENDING, // 等待首次执行
PROCESSING, // 已被某个工作节点领取
RETRY, // 执行失败,等待退避重试
SUCCESS, // 对象已删除
FAILED, // 超过最大重试次数
IGNORED // 文件已激活或仍被引用,无需删除
}
任务表应包含:
record_id
object_key
status
next_retry_at
retry_count
max_retries
lease_owner
lease_until
lease_version
lease_until 解决工作节点宕机后的重新领取;lease_version 是递增的隔离令牌,用于阻止已经失效的旧工作节点提交结果。
领取任务只锁任务表,不在同一事务中锁文件记录,避免与上传确认形成反向加锁:
@Transactional
public List<CleanupLease> claimDueTasks(String workerId, int limit, Instant now) {
// 查询 PENDING、到期 RETRY,以及租约已经过期的 PROCESSING
List<CleanupTask> tasks = cleanupRepository
.selectClaimableForUpdateSkipLocked(now, limit);
return tasks.stream().map(task -> {
long newVersion = task.getLeaseVersion() + 1;
Instant leaseUntil = now.plus(Duration.ofMinutes(2));
// 领取结果与 leaseVersion 一起返回给事务外执行器
cleanupRepository.markProcessing(
task.getId(), workerId, leaseUntil, newVersion
);
return CleanupLease.from(task, workerId, leaseUntil, newVersion);
}).toList();
}
核心查询可以使用 FOR UPDATE SKIP LOCKED,让多个服务实例并行领取不同任务:
-- 活跃租约不会被重复领取,失效租约可以由其他节点接管
SELECT *
FROM cleanup_task
WHERE (status IN ('PENDING', 'RETRY') AND next_retry_at <= :now)
OR (status = 'PROCESSING' AND lease_until <= :now)
ORDER BY next_retry_at
LIMIT :limit
FOR UPDATE SKIP LOCKED;
取得租约不等于可以立即删除。执行器还要在短事务中按照统一顺序锁定 file_upload → cleanup_task,重新验证租约和业务状态:
@Transactional
public DeleteDecision prepareDelete(CleanupLease lease, Instant now) {
// 所有涉及两张表的流程统一先锁文件,再锁任务
FileUpload record = fileRepository
.selectIncludingDeletedForUpdate(lease.recordId());
CleanupTask task = cleanupRepository.selectForUpdate(lease.taskId());
if (!task.isOwnedBy(lease, now)) {
return DeleteDecision.LEASE_LOST;
}
// 有效文件或仍被引用的对象绝不能删除
if (record.getUploadStatus() == UploadStatus.ACTIVE
|| fileRepository.isObjectKeyReferenced(lease.objectKey())) {
cleanupRepository.markIgnoredIfOwned(lease);
return DeleteDecision.IGNORED;
}
// UPLOADING 是否超时统一由 upload_deadline_at 判断,状态修复交给 UploadRepairJob
if (record.getUploadStatus() == UploadStatus.UPLOADING) {
Instant retryAt = record.getUploadDeadlineAt().isAfter(now)
? record.getUploadDeadlineAt()
: now.plusSeconds(30);
cleanupRepository.scheduleRetryIfOwned(lease, retryAt);
return DeleteDecision.POSTPONED;
}
return record.getUploadStatus() == UploadStatus.UPLOAD_FAILED
? DeleteDecision.ALLOWED
: DeleteDecision.REJECTED;
}
对象删除仍然放在事务外:
public void executeCleanup(CleanupLease lease) {
DeleteDecision decision = cleanupTxService.prepareDelete(lease, Instant.now());
if (decision != DeleteDecision.ALLOWED) {
return;
}
try {
// 重复删除必须产生相同结果
storage.deleteIfExists(lease.objectKey());
// 只有仍持有相同 leaseVersion 的节点才能完成任务
cleanupTxService.markSuccessIfOwned(lease);
} catch (Exception e) {
// 失败更新同样校验 leaseVersion,防止旧节点覆盖新租约
cleanupTxService.scheduleRetryIfOwned(lease, e.getMessage());
}
}
如果节点在取得租约后宕机,它来不及执行重试更新,但 lease_until 到期后任务会重新进入可领取集合。如果对象已经删除、数据库尚未标记成功,新节点再次调用 deleteIfExists 仍会成功,然后补上任务状态。
方案二:Transactional Outbox + MQ
如果系统已经使用可靠消息队列,可以把任务执行权交给消息中间件:
RabbitMQ使用手动确认。消费者断开连接且消息尚未ACK时,消息会自动重新入队。RabbitMQ Consumer AcknowledgementsAmazon SQS使用Visibility Timeout。消费者未在可见性超时内删除消息,消息会重新变为可见。SQS Visibility Timeout
但不能在数据库事务中直接发送 MQ 消息:数据库提交和消息发送属于两个独立资源,任意一方失败都会产生双写不一致。通常先在数据库事务中写入 Outbox,再由独立发布器投递消息。
Outbox 事件还应保存 not_before,其值至少为 upload_deadline_at + cleanup_grace。发布器只投递已经到期的事件,避免文件仍在正常上传时就触发删除。确认成功时,如果事件尚未发布,则在同一事务中将其改为 CANCELLED。
消费者处理逻辑保持幂等,并区分“无需删除”和“暂时不能删除”:
public void onFileDelete(FileDeleteCommand command, MessageAck ack) {
try {
DeleteDecision decision = cleanupService.prepareDelete(command);
switch (decision) {
case ALLOWED -> {
// 对象不存在也视为成功,保证重复投递安全
storage.deleteIfExists(command.objectKey());
cleanupService.markStorageDeleted(command.recordId());
ack.success();
}
case IGNORED -> {
// 文件已经激活或仍被引用,无需再清理
ack.success();
}
case POSTPONED -> {
// 尚未到清理时间,使用延迟重试,不能直接 ACK 丢弃
ack.retryLater(cleanupService.nextRetryAt(command));
}
default -> {
// 非预期状态进入死信并告警,等待人工或对账任务处理
ack.deadLetter("文件状态不允许执行清理");
}
}
} catch (RetryableException e) {
// 临时故障重新入队;达到上限后进入死信队列
ack.retryOrDeadLetter(e.getMessage());
}
}
这里的 MessageAck 是对不同消息中间件确认能力的抽象:success 对应正常确认。RabbitMQ 可先把消息可靠发布到延迟重试队列并等待 Publisher Confirm,再确认原消息;SQS 可调整当前消息的 Visibility Timeout。不要简单使用立即 requeue,否则持续失败时会形成高频重投。若改为重新写入带执行时间的 Outbox,必须先持久化重试事件,再确认原消息;宕机窗口造成的重复由幂等消费吸收。
MQ 解决的是消息执行权和宕机重投,不会自动解决数据库与对象存储的一致性。消费者仍需重新检查业务状态、支持重复投递,并为超过重试上限的消息配置死信队列和告警。发布器在消息已发送但尚未来得及把 Outbox 标记为 SENT 时宕机,也可能造成重复投递,因此幂等不是可选优化。
两种方案如何选择?
| 方案 | 优点 | 成本 | 适用场景 |
|---|---|---|---|
| 数据库任务表 + 租约 | 不新增中间件,可与业务记录同事务写入 | 需要维护租约、重试和扫描逻辑 | 单体、早期微服务、中低任务量 |
Outbox + MQ |
利用确认、重投、积压和死信能力,扩展性更好 | 需要部署和运维消息队列,仍需 Outbox | 已有 MQ、任务量高、服务已拆分 |
如果项目已经有本地消息表,优先完善显式租约,不必只为文件删除引入 MQ。如果企业基础设施已经提供 RabbitMQ、RocketMQ 或 SQS,则优先使用 Outbox + MQ,避免在每个业务模块重复实现任务租约。
无论选择哪种方案,上传记录使用逻辑删除时,清理查询都必须显式包含已删除记录。很多 ORM 会自动附加 is_deleted = false,导致清理执行器看不到最需要处理的失败数据。
第六步:定时修复超时的 UPLOADING 死记录
延迟清理任务负责处理已知的补偿任务,但还需要一个独立的 UploadRepairJob 定期扫描数据库。它解决的是以下异常链路:
上传或确认失败
↓
catch 中 markFailed() 再次失败
↓
记录永久停留在 UPLOADING
↓
定时任务发现 upload_deadline_at 已过期
↓
UPLOADING → UPLOAD_FAILED
↓
逻辑删除并补建清理意图
不要直接根据 created_at 判断死记录。上传时间会受文件大小、网络环境和对象存储状态影响,预占时应写入明确的 upload_deadline_at。该期限必须大于对象存储客户端允许的最大上传时长。
定时任务只负责发现候选记录,每条记录交给独立事务修复:
@Component
@RequiredArgsConstructor
public class UploadRepairJob {
private static final int BATCH_SIZE = 100;
private final FileUploadRepository fileRepository;
private final UploadRepairTxService repairTxService;
@Scheduled(fixedDelayString = "${upload.repair.fixed-delay:60000}")
public void repairExpiredUploading() {
// 每轮只扫描一个有限批次,避免长时间占用调度线程
Instant now = Instant.now();
List<ExpiredUploadCandidate> candidates =
fileRepository.findExpiredUploading(now, BATCH_SIZE);
for (ExpiredUploadCandidate candidate : candidates) {
try {
// 每条死记录使用独立事务,单条失败不影响后续记录
repairTxService.repairOne(candidate, now);
} catch (Exception e) {
log.error("修复上传死记录失败。recordId={}", candidate.recordId(), e);
}
}
}
}
应用需要启用定时任务:
@Configuration
@EnableScheduling
public class SchedulingConfiguration {
// 启用 Spring 定时任务调度
}
查询条件应限定状态、截止时间和批次大小:
-- 只扫描已超过上传截止时间的有效预占记录
SELECT id, owner_id
FROM file_upload
WHERE upload_status = 'UPLOADING'
AND is_deleted = false
AND upload_deadline_at <= :now
ORDER BY upload_deadline_at
LIMIT :batchSize;
真正的状态修改需要再次校验,不能相信定时任务之前读取的快照:
@Service
@RequiredArgsConstructor
public class UploadRepairTxService {
private static final Duration CLEANUP_GRACE = Duration.ofMinutes(5);
private final FileUploadRepository fileRepository;
private final CleanupTaskRepository cleanupRepository;
@Transactional
public boolean repairOne(ExpiredUploadCandidate candidate, Instant now) {
// 与正常上传使用同一用户锁,避免集合状态并发变化
fileRepository.lockOwner(
UploadLockKey.NAMESPACE,
UploadLockKey.ownerSlot(candidate.ownerId())
);
// 获取最新状态,不能直接相信扫描阶段的快照
FileUpload record = fileRepository.selectIncludingDeleted(candidate.recordId());
if (record == null
|| record.getUploadStatus() != UploadStatus.UPLOADING
|| record.getUploadDeadlineAt().isAfter(now)) {
return false;
}
// 通过条件更新标记失败并释放用户配额
boolean repaired = fileRepository.markExpiredAsFailed(
record.getId(),
now,
"上传超时,由定时任务修复"
);
if (!repaired) {
return false;
}
// 补偿任务缺失或失效时重新置为待执行
cleanupRepository.ensurePending(
record.getId(),
record.getObjectKey(),
now.plus(CLEANUP_GRACE)
);
return true;
}
}
markExpiredAsFailed 必须是带原状态和截止时间的条件更新:
-- 只有仍处于 UPLOADING 且确实过期的记录才允许修复
UPDATE file_upload
SET upload_status = 'UPLOAD_FAILED',
is_deleted = true,
failure_reason = :reason,
updated_at = :now
WHERE id = :recordId
AND upload_status = 'UPLOADING'
AND upload_deadline_at <= :now;
这里再次获取当前用户的上传锁,是为了让定时修复和上传、删除等集合变更遵守同一并发协议。即使多个服务实例同时执行定时任务,锁内复查和条件更新也只允许一个实例完成修复。
ensurePending 不能无条件插入重复任务。数据库租约方案应保证本地任务为 PENDING/RETRY;MQ 方案应保证对应 Outbox 事件仍可发布。建议建立 (record_id, task_type) 唯一约束,并让清理意图与状态修复处于同一个数据库事务。
定时修复只修改数据库并确保清理任务存在,不在事务中调用对象存储。实际删除至少延迟一个 CLEANUP_GRACE,并且删除前仍需重新检查状态。这样可以降低“原上传请求刚好仍在结束阶段,修复任务提前删除对象”的竞态风险。
还应满足两个时间约束:
对象存储客户端最大上传时长 < upload_deadline_at - created_at
清理执行时间 > upload_deadline_at + CLEANUP_GRACE
如果使用对象存储分片上传,还需要定期终止未完成的分片上传,或配置存储侧生命周期规则。仅删除最终对象无法清理尚未完成的分片。
验证实现
不要只验证一次正常上传。可靠性问题通常出现在异常窗口和并发时序中。
| 场景 | 条件 | 预期结果 |
|---|---|---|
| 正常上传 | 合法文件,配额充足 | 状态最终为 ACTIVE,清理意图变为 IGNORED/CANCELLED |
| 非法类型 | 扩展名合法但文件签名不匹配 | 上传前拒绝 |
| 重复上传 | 同一用户、相同哈希 | 只保留一条有效记录 |
| 配额并发 | 已有 4 份,同时提交 2 个文件 | 只有一个请求成功预占 |
| 存储失败 | put 抛出异常 |
记录变为 UPLOAD_FAILED,清理任务保留 |
| 失败标记也失败 | markFailed 数据库异常 |
UploadRepairJob 超时后标记失败并补建清理任务 |
| 确认失败 | 文件已上传,数据库更新失败 | 不返回成功,后续补偿删除对象 |
| 进程宕机 | 上传成功后、确认前停止进程 | 超时任务接管并清理 |
| 补偿失败 | 对象存储暂时不可用 | 清理任务退避重试 |
| 租约执行者宕机 | 任务为 PROCESSING,但租约过期 |
其他节点重新领取并继续幂等删除 |
| 旧执行者恢复 | 原节点持有旧 leaseVersion |
不能覆盖新节点提交的任务状态 |
MQ 重复投递 |
同一删除消息被消费多次 | 重复删除安全,最终只保留一个业务结果 |
MQ 消费者宕机 |
删除对象后、确认消息前退出 | 消息重新投递,并通过幂等删除完成收敛 |
并发配额测试的核心断言如下:
@Test
void onlyOneRequestCanOccupyLastSlot() throws Exception {
// given:用户只剩一个可用配额
givenUserAlreadyHasFiles(4);
// when:两个请求同时尝试预占
Future<?> first = pool.submit(() -> upload(fileA));
Future<?> second = pool.submit(() -> upload(fileB));
waitFor(first, second);
// then:最终数量不能突破上限
assertThat(countAliveFiles()).isEqualTo(5);
assertThat(successCount(first, second)).isEqualTo(1);
}
死记录修复至少要验证状态和补偿任务同时收敛:
@Test
void shouldRepairExpiredUploadingWhenMarkFailedWasLost() {
// given:失败标记和原清理任务均未正常生效
Long recordId = givenExpiredUploadingRecordWithoutUsableCleanupTask();
// when:定时修复扫描到过期记录
repairJob.repairExpiredUploading();
// then:记录退出配额,并重新获得待执行清理任务
assertThat(loadIncludingDeleted(recordId).getUploadStatus())
.isEqualTo(UploadStatus.UPLOAD_FAILED);
assertThat(loadIncludingDeleted(recordId).getIsDeleted()).isTrue();
assertThat(findCleanupTask(recordId).getStatus()).isEqualTo(TaskStatus.PENDING);
}
数据库租约方案还要验证过期重领和隔离令牌:
@Test
void shouldReclaimExpiredCleanupLeaseAndRejectStaleWorker() {
// given:节点 A 的租约已经过期
CleanupLease leaseA = givenExpiredProcessingTask();
// when:节点 B 重新领取相同任务
CleanupLease leaseB = cleanupService.claimDueTasks("worker-B", 1, Instant.now())
.get(0);
// then:租约版本递增,节点 A 不能再提交成功
assertThat(leaseB.leaseVersion()).isGreaterThan(leaseA.leaseVersion());
assertThat(cleanupService.markSuccessIfOwned(leaseA)).isFalse();
assertThat(cleanupService.markSuccessIfOwned(leaseB)).isTrue();
}
MQ 方案至少要重复投递同一条删除消息,验证消费者能够安全执行多次:
@Test
void shouldHandleDuplicatedFileDeleteMessageIdempotently() {
FileDeleteCommand command = givenFailedUploadDeleteCommand();
// 同一条消息被重复投递
consumer.onFileDelete(command, firstAck);
consumer.onFileDelete(command, secondAck);
assertThat(storage.exists(command.objectKey())).isFalse();
assertThat(firstAck.isSuccess()).isTrue();
assertThat(secondAck.isSuccess()).isTrue();
}
仅使用内存数据库通常无法验证 PostgreSQL Advisory Lock。这类测试应使用真实 PostgreSQL,可以通过 Testcontainers 启动隔离实例。
为什么这样设计?
方案对比
| 方案 | 优点 | 缺点 | 适用场景 |
|---|---|---|---|
| 先上传后写库 | 实现最简单 | 孤儿文件难追踪 | 可丢弃的临时文件 |
| 长事务包住全流程 | 代码表面集中 | 不能回滚对象存储,长期占用连接和锁 | 不推荐用于网络上传 |
Redis 分布式锁 |
可跨应用实例和数据库 | 引入额外依赖,需要处理超时、续期和锁失效 | 多数据库、多服务协调 |
| 数据库守卫行 | 语义直观,可使用 FOR UPDATE |
需要稳定存在的可锁行 | 数据模型已有用户或配额行 |
| 短事务预占 + 补偿 | 状态可追踪,故障可恢复 | 状态机和任务表增加复杂度 | 需要业务记录和可靠性的上传接口 |
本文选择数据库锁,是因为配额检查和文件记录本来就位于数据库中。锁、计数和预占可以在同一事务中完成,减少跨系统协调。
如果上传服务已经拆分为独立数据库,而多个服务仍能修改同一份配额,数据库锁就不再满足协调范围。此时应重新定义资源所有权:优先让一个服务独占文件写入口;确实需要跨服务写入时,再评估 Redis、etcd 或集中式配额服务。
边界条件与失败场景
1. 是否应该复用逻辑删除文件的对象键?
通过哈希找到已删除记录,如果对象仍然存在,确实可以跳过重复上传。
很多读者可能会想到一个很常见的现象:为什么百度网盘等云存储产品上传某些几十 GB 的大文件时,有时几秒甚至瞬间就显示上传成功?
这类能力通常被称为“秒传”。其核心思路就是内容去重:客户端先计算文件哈希或其他文件指纹,服务端根据这些信息判断存储系统中是否已经存在相同内容。如果已经存在,就不需要再次传输完整文件,而是直接复用已有的物理数据,并为当前用户建立新的逻辑引用。
例如:
用户 A 上传 file.zip
↓
物理存储保存一份文件
↓
content_hash = abc123
用户 B 上传相同 file.zip
↓
计算得到相同哈希 abc123
↓
发现物理文件已经存在
↓
不再上传几十 GB 数据
↓
直接建立新的逻辑引用
从用户角度看,几十 GB 的文件可能“瞬间上传完成”;但从存储系统角度看,并不是几十 GB 数据真的在几秒内完成了网络传输,而是系统识别出了已经存在的相同内容并进行了复用。
不过,这种设计也意味着:一个物理对象可能同时被多条业务记录引用。
此时删除文件就不再是简单的“删除当前记录对应的对象”,而必须额外解决:
- 是否还有其他有效记录引用该对象?
- 新引用创建和旧对象清理是否存在竞态?
- 对象是否保证不可变?
- 引用计数或引用关系更新失败后如何恢复和对账?
例如:
┌→ 用户 A 的文件记录
物理文件 X ──┼→ 用户 B 的文件记录
└→ 用户 C 的文件记录
此时用户 A 删除文件时:
删除 A 的逻辑引用
↓
检查是否仍有 B、C 引用
↓
有引用 → 不能删除物理文件
最后一个引用被删除
↓
才允许真正删除物理文件
这实际上已经从简单的“文件上传”演变成了内容寻址存储 + 引用生命周期管理问题。
因此,通用业务系统通常默认:
每次有效上传使用独立对象键,业务记录与物理对象保持一对一关系。
这样虽然可能牺牲一部分存储去重能力,但删除、补偿和故障恢复都会简单很多。
只有当文件规模较大、重复率较高,存储成本明显值得优化时,再考虑引入类似网盘“秒传”的内容去重方案,例如内容寻址存储、独立对象引用表或引用计数机制。
2. 文件上传成功,但失败状态也更新失败
不要在 catch 中无限重试数据库,也不要让标记失败异常覆盖最初的上传异常。
正确处理是:
- 记录包含
recordId的错误日志和监控指标。 - 保留预先创建的清理任务。
UploadRepairJob定时扫描超过upload_deadline_at的UPLOADING记录。- 修复事务将其推进为
UPLOAD_FAILED,逻辑删除并补建清理任务。
补偿逻辑也会失败,所以可靠方案需要“补偿的重试”和“状态的对账”,而不是只写一个 catch。
3. 大文件不适合经过应用服务器中转
视频、压缩包等大文件会占用应用带宽、连接和临时磁盘。更合适的流程通常是:
客户端申请上传会话
↓
服务端校验权限并返回预签名地址
↓
客户端直传对象存储
↓
客户端或存储事件通知服务端确认
↓
服务端校验对象大小、摘要并激活记录
直传改变的是数据链路,不会消除状态预占、最终确认、超时清理和幂等问题。
4. 上传成功,但客户端没有收到响应
服务可能已经把记录确认成 ACTIVE,但响应在网络中丢失。客户端重试后,相同哈希查重会阻止再次写入,却只能得到“文件已存在”,无法确定上一次请求是否成功。
需要严格请求幂等时,可以让客户端发送 Idempotency-Key,服务端将它与用户、上传记录和最终结果关联。同一用户使用相同幂等键重试时,直接返回第一次处理结果。
文件哈希不能完全替代幂等键:两个不同业务请求可能上传相同内容,而同一次请求也可能因为客户端重新编码产生不同字节。
5. 开启 S3 Versioning 后,普通删除不等于物理删除
本文代码中的 deleteIfExists(objectKey) 默认描述未开启版本控制的对象存储。S3 Bucket 开启 Versioning 后,不携带 versionId 删除对象,只会创建一个 Delete Marker;旧版本仍保留在 Bucket 中,可以被指定版本读取,也会继续产生存储费用。AWS S3 Versioning 工作原理
此时必须重新定义 FileStorage 的删除契约:
- 上传成功后保存对象存储返回的
versionId,清理时删除指定版本。 - 如果进程在上传成功、保存
versionId前宕机,使用“每次上传生成唯一对象键”的前提,按对象键枚举并删除所有版本和Delete Marker。 - 配置 Noncurrent Version Expiration 生命周期规则,作为长期存储费用的最后一道兜底。
- 清理账号需要最小化授予列举版本和删除指定版本的权限,并记录审计日志。
因此,幂等删除接口不能只保证“再次读取返回 404”,还要根据 Bucket 的版本控制策略保证目标版本确实被删除。若业务需要保留历史版本,则不能全量删除,而应把应删除的 versionId 持久化到清理任务中。
常见踩坑
坑一:只在上传前执行一次查重
事务外查重可以作为快速失败优化,但不能作为最终依据。两个并发请求仍可能同时通过。权威查重必须在锁内再次执行,并配合数据库唯一约束。
坑二:只给上传接口加锁
批量导入、恢复已删除文件等入口如果也会改变文件集合,就必须遵守同一锁协议。锁的正确性取决于所有参与者是否协作。
坑三:把文件解析放进预占事务
PDF、DOCX 解析可能耗时且消耗大量内存。它不应延长配额锁和数据库事务。需要解析时,可在上传确认后发布任务,将解析状态作为独立生命周期管理。
坑四:把文件哈希当作安全校验
哈希只说明两份字节内容是否相同,不能判断文件是否恶意,也不能证明扩展名与服务端检测结果匹配。
坑五:清理任务只检查 URL 是否存在
“对象存在”不能说明它应该被删除。清理任务必须结合记录状态、upload_deadline_at 和引用关系判断;执行权则由数据库租约或 MQ 确认机制管理。
最佳实践
在生产环境中,建议:
- 分层限流:网关限制请求速率,容器限制请求大小,业务层限制文件类型和用户配额。
- 生成对象键:原始文件名只用于展示,不直接参与存储路径。
- 使用状态机:禁止任意覆盖状态,通过条件更新保证合法转换。
- 保持短事务:事务内只执行锁、查询和必要写入。
- 提前保存补偿意图:外部副作用发生前,先留下可恢复依据。
- 分离业务状态与任务状态:文件表记录上传结果,任务表或
MQ管理执行权。 - 补偿操作幂等:重复删除、租约重领和消息重复投递都应产生确定结果。
- 防止旧执行者提交:数据库租约使用
leaseVersion,MQ消费使用确认或可见性超时。 - 增加可观测性:监控
UPLOADING超时、租约过期重领、死信消息、孤儿对象和上传失败率。 - 定期对账:数据库记录和对象存储清单应能相互核对。
同时避免:
- 直接信任
Content-Type、扩展名或原始文件名。 - 在数据库事务中执行长时间网络上传。
- 认为写了
@Transactional就解决了跨资源一致性。 - 让失败记录永久占用用户配额。
- 在没有引用模型时复用同一个物理对象。
总结
可靠的文件上传接口不是一次简单的文件复制,而是一段跨数据库和对象存储的状态流转。
整体思路可以概括为:
安全校验与哈希
↓
同一所有者串行检查
↓
短事务预占 UPLOADING
↓
事务外上传文件
↓
条件更新确认 ACTIVE
↓
数据库租约或 Outbox + MQ 执行失败补偿
↓
超时修复与定期对账
这套方案的核心优势是:
- 并发配额不会被简单的“先查后写”突破。
- 数据库事务不会被网络上传长期占用。
- 任意阶段失败后都有状态和任务可追踪。
- 清理执行器宕机后,租约重领或
MQ重投能够继续处理。 - 超时上传和异常任务能够通过定时修复与对账自行收敛。
真正需要记住的是:
数据库无法替对象存储回滚;可靠上传依赖短事务、业务状态与任务状态分离,以及能够在执行者宕机后重新投递的补偿机制。
参考资料
OWASP File Upload Cheat SheetSpring Framework MultipartFileSpring Boot Common Application PropertiesPostgreSQL Advisory LocksSpring Framework Transaction ManagementSpring Framework Rollback RulesRabbitMQ Consumer AcknowledgementsRabbitMQ Dead Letter ExchangesAmazon SQS Visibility TimeoutAmazon S3 Versioning