大文件分片上传、断点续传和秒传
第一次写技术文章,有很多不足,希望大家多提提意见。
需求分析
开发分片上传、断点续传和秒传功能是为了解决大文件传输场景中的效率、可靠性与资源浪费问题。分片上传通过将大文件拆分为小块,降低服务器压力并支持并行传输,避免单次请求超时或内存溢出;断点续传允许网络中断后从中断点继续上传,减少重复传输带来的时间和流量损耗;秒传通过文件哈希值校验识别重复内容,直接复用已有文件,消除冗余存储并实现“瞬时”完成上传。
通过 AWS S3 的 SDK 访问了 MinIO 服务,为什么使用 AWS S3 SDK 而不是 MinIO 原生 SDK 来访问 MinIO 服务呢?AWS S3 是最广泛使用的对象存储服务,可移植性强,而 MinIO 原生 SDK 是 S3 SDK 的轻量级封装,功能上基本一致。并且 AWS S3 生态、稳定性和性能优化更好。
概念
- 分片上传:将要上传的文件按照一定的大小分为多个小分片(数据块),逐个上传到服务器,最后将所有的小分片合并成原始的完整文件。
- 断点续传:上传中断后,再次上传相同文件时已经上传的分片不再上传,从中断点继续上传。
- 秒传:根据文件唯一标识判断是否已经上传,已经上传过的文件不会再次上传。
方案设计
分片上传:在前端开始上传文件前,先通过 MD5 算法获取文件的唯一标识,发送请求给后端获取文件的上传任务记录,如果数据库中没有查到该文件的上传任务记录,数据库新增一条上传任务记录。前端将文件分为若干个分片,依次将分片序号和文件标识发送给后端,返回该分片的预签名url给前端,前端上传分片到 MinIO 中。循环上传分片直到所有分片上传完成后发送合并分片请求给后端,后端调用 MinIO 的 api 进行合并分片。
断点续传:开始上传文件前,发送请求给后端根据文件唯一标识获取该文件的上传任务记录,获取到已经上传过的分片,前端基于中断点继续上传分片。
秒传:开始上传文件前,发送请求给后端根据文件唯一标识获取该文件的上传任务记录,如果上传任务已完成,则不需要重新上传。
后端开发
建表
▼sql复制代码-- 创建库 create database if not exists minioUpload; -- 切换库 use minioUpload; SET NAMES utf8mb4; SET FOREIGN_KEY_CHECKS = 0; DROP TABLE IF EXISTS `sys_upload_task`; CREATE TABLE `sys_upload_task` ( `id` bigint NOT NULL, `upload_id` varchar(255) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci NOT NULL COMMENT '分片上传的uploadId', `file_identifier` varchar(500) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci NOT NULL COMMENT '文件唯一标识(md5)', `file_name` varchar(500) COLLATE utf8mb4_general_ci NOT NULL COMMENT '文件名', `bucket_name` varchar(255) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci NOT NULL COMMENT '所属桶名', `object_key` varchar(500) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci NOT NULL COMMENT '文件的key', `total_size` bigint NOT NULL COMMENT '文件大小(byte)', `chunk_size` bigint NOT NULL COMMENT '每个分片大小(byte)', `chunk_num` int NOT NULL COMMENT '分片数量', PRIMARY KEY (`id`), UNIQUE KEY `uq_file_identifier` (`file_identifier`) USING BTREE, UNIQUE KEY `uq_upload_id` (`upload_id`) USING BTREE ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_general_ci COMMENT='分片上传-分片任务记录'; SET FOREIGN_KEY_CHECKS = 1;
利用 Mybatis-Plus 生成实体类、mapper 文件、xml 文件。实体类如下:
▼java复制代码@Data @Accessors(chain = true) @TableName("sys_upload_task") @NoArgsConstructor @AllArgsConstructor @Builder public class SysUploadTask implements Serializable { private Long id; //分片上传的uploadId private String uploadId; //文件唯一标识(md5) private String fileIdentifier; //文件名 private String fileName; //所属桶名 private String bucketName; //文件的key private String objectKey; //文件大小(byte) private Long totalSize; //每个分片大小(byte) private Long chunkSize; //分片数量 private Integer chunkNum; @TableField(exist = false) private URL url; }
引入AWS S3(引入依赖,添加配置,配置类)
在项目 pom.xml 文件中,引入一下依赖配置:
▼xml复制代码<dependency> <groupId>com.amazonaws</groupId> <artifactId>aws-java-sdk-s3</artifactId> <version>1.12.263</version> </dependency>
在 application.yml 文件中添加配置
▼yaml复制代码minio: access-key: admin access-secret: admin123456 endpoint: http://127.0.0.1:9005 bucket: minio-upload
读取配置文件中的MinIO相关配置
▼java复制代码@Data @Component @ConfigurationProperties(prefix = "minio") public class MinioProperties { @Value("endpoint") private String endpoint; @Value("access-key") private String accessKey; @Value("access-secret") private String accessSecret; @Value("bucket") private String bucket; }
编写 AmazonS3Config 配置类
▼java复制代码@Configuration public class AmazonS3Config { @Resource private MinioProperties minioProperties; @Bean(name = "amazonS3Client") public AmazonS3 amazonS3Client () { //设置连接时的参数 ClientConfiguration config = new ClientConfiguration(); //设置连接方式为HTTP,可选参数为HTTP和HTTPS config.setProtocol(Protocol.HTTP); //设置网络访问超时时间 config.setConnectionTimeout(5000); config.setUseExpectContinue(true); AWSCredentials credentials = new BasicAWSCredentials(minioProperties.getAccessKey(), minioProperties.getAccessSecret()); //设置Endpoint AwsClientBuilder.EndpointConfiguration end_point = new AwsClientBuilder.EndpointConfiguration(minioProperties.getEndpoint(), Regions.US_EAST_1.name()); AmazonS3 amazonS3 = AmazonS3ClientBuilder.standard() .withClientConfiguration(config) .withCredentials(new AWSStaticCredentialsProvider(credentials)) .withEndpointConfiguration(end_point) .withPathStyleAccessEnabled(true).build(); return amazonS3; } }
关键点:
- Endpoint:通过
minioProperties.getEndpoint()指定了MinIO服务的地址(如http://127.0.0.1:9005),而不是AWS S3的默认地址。 - Path Style Access:通过
withPathStyleAccessEnabled(true)启用了路径风格的访问(MinIO通常需要这种方式)。 - 认证信息:使用MinIO的
access-key和access-secret替代AWS的认证信息。
Minio 相关常量类
▼java复制代码public interface MinioConstant { // 分块大小 int DEFAULT_CHUNK_SIZE = 10 * 1024 * 1024; // 预签名url过期时间(ms) Long PRE_SIGN_URL_EXPIRE = 60 * 10 * 1000L; }
获取上传进度
前端根据 MD5 算法计算出文件的哈希值,根据文件哈希值获取文件上传进度
▼java复制代码// Controller 层 /** * 获取上传进度 * @param identifier 文件 md5 * @return */ @GetMapping("/{identifier}") public Result<TaskInfoDTO> taskInfo (@PathVariable("identifier") String identifier) { return Result.ok(sysUploadTaskService.getTaskInfo(identifier)); } // Server 层 public TaskInfoDTO getTaskInfo(String identifier) { SysUploadTask task = getByIdentifier(identifier); if (task == null) { return null; } TaskInfoDTO result = new TaskInfoDTO().setFinished(true).setTaskRecord(TaskRecordDTO.convertFromEntity(task)).setPath(getPath(task.getBucketName(), task.getObjectKey())); boolean doesObjectExist = amazonS3.doesObjectExist(task.getBucketName(), task.getObjectKey()); if (!doesObjectExist) { // 未上传完,返回已上传的分片 ListPartsRequest listPartsRequest = new ListPartsRequest(task.getBucketName(), task.getObjectKey(), task.getUploadId()); PartListing partListing = amazonS3.listParts(listPartsRequest); result.setFinished(false).getTaskRecord().setExitPartList(partListing.getParts()); } return result; } public SysUploadTask getByIdentifier(String identifier) { return sysUploadTaskMapper.selectOne(new QueryWrapper<SysUploadTask>().lambda().eq(SysUploadTask::getFileIdentifier, identifier)); }
返回结果 TaskInfoDTO 类
▼java复制代码@Data @ToString @NoArgsConstructor @AllArgsConstructor @Accessors(chain = true) public class TaskInfoDTO { /** * 是否完成上传(是否已经合并分片) */ private boolean finished; /** * 文件地址 */ private String path; /** * 上传记录 */ private TaskRecordDTO taskRecord; }
TaskRecordDTO 类
▼java复制代码@Data @ToString @Accessors(chain = true) public class TaskRecordDTO extends SysUploadTask { /** * 已上传完的分片 */ private List<PartSummary> exitPartList; public static TaskRecordDTO convertFromEntity (SysUploadTask task) { TaskRecordDTO dto = new TaskRecordDTO(); BeanUtil.copyProperties(task, dto); return dto; } }
初始化上传任务
若数据库中不存在该文件的上传任务,新建一个上传任务。
▼java复制代码// Controller /** * 创建一个上传任务 * @return */ @PostMapping public Result<TaskInfoDTO> initTask (@Valid @RequestBody InitTaskParam param, BindingResult bindingResult) { if (bindingResult.hasErrors()) { return Result.error(bindingResult.getFieldError().getDefaultMessage()); } return Result.ok(sysUploadTaskService.initTask(param)); } // Server层 @Override public TaskInfoDTO initTask(InitTaskParam param) { Date currentDate = new Date(); String bucketName = minioProperties.getBucket(); String fileName = param.getFileName(); String suffix = fileName.substring(fileName.lastIndexOf(".")+1, fileName.length()); String key = StrUtil.format("{}/{}.{}", DateUtil.format(currentDate, "YYYY-MM-dd"), IdUtil.randomUUID(), suffix); String contentType = MediaTypeFactory.getMediaType(key).orElse(MediaType.APPLICATION_OCTET_STREAM).toString(); ObjectMetadata objectMetadata = new ObjectMetadata(); objectMetadata.setContentType(contentType); InitiateMultipartUploadResult initiateMultipartUploadResult = amazonS3 .initiateMultipartUpload(new InitiateMultipartUploadRequest(bucketName, key).withObjectMetadata(objectMetadata)); String uploadId = initiateMultipartUploadResult.getUploadId(); URL url = amazonS3.getUrl(bucketName, key); System.out.println("返回的url的地址为"+ url); int chunkNum = (int) Math.ceil(param.getTotalSize() * 1.0 / param.getChunkSize()); SysUploadTask task = SysUploadTask.builder() .bucketName(minioProperties.getBucket()) .chunkNum(chunkNum) .chunkSize(param.getChunkSize()) .totalSize(param.getTotalSize()) .fileIdentifier(param.getIdentifier()) .fileName(fileName) .objectKey(key) .uploadId(uploadId) .build(); sysUploadTaskMapper.insert(task); return new TaskInfoDTO().setFinished(false).setTaskRecord(TaskRecordDTO.convertFromEntity(task)).setPath(getPath(bucketName, key)); } @Override public String getPath(String bucket, String objectKey) { return StrUtil.format("{}/{}/{}", minioProperties.getEndpoint(), bucket, objectKey); }
初始化上传任务请求的请求参数接收类 InitTaskParam
▼java复制代码@Data @ToString @Accessors(chain = true) public class InitTaskParam { /** * 文件唯一标识(MD5) */ @NotBlank(message = "文件标识不能为空") private String identifier; /** * 文件大小(byte) */ @NotNull(message = "文件大小不能为空") private Long totalSize; /** * 分片大小(byte) */ @NotNull(message = "分片大小不能为空") private Long chunkSize; /** * 文件名称 */ @NotBlank(message = "文件名称不能为空") private String fileName; }
分片上传
前端获取每个分片的预签名上传地址后依次上传分片
▼java复制代码// Controller /** * 获取每个分片的预签名上传地址 * @param identifier * @param partNumber * @return */ @GetMapping("/{identifier}/{partNumber}") public Result preSignUploadUrl (@PathVariable("identifier") String identifier, @PathVariable("partNumber") Integer partNumber) { SysUploadTask task = sysUploadTaskService.getByIdentifier(identifier); if (task == null) { return Result.error("分片任务不存在"); } Map<String, String> params = new HashMap<>(); params.put("partNumber", partNumber.toString()); params.put("uploadId", task.getUploadId()); return Result.ok(sysUploadTaskService.genPreSignUploadUrl(task.getBucketName(), task.getObjectKey(), params)); } @Override public String genPreSignUploadUrl(String bucket, String objectKey, Map<String, String> params) { Date currentDate = new Date(); Date expireDate = DateUtil.offsetMillisecond(currentDate, MinioConstant.PRE_SIGN_URL_EXPIRE.intValue()); GeneratePresignedUrlRequest request = new GeneratePresignedUrlRequest(bucket, objectKey) .withExpiration(expireDate).withMethod(HttpMethod.PUT); if (params != null) { params.forEach((key, val) -> request.addRequestParameter(key, val)); } URL preSignedUrl = amazonS3.generatePresignedUrl(request); return preSignedUrl.toString(); } public SysUploadTask getByIdentifier(String identifier) { return sysUploadTaskMapper.selectOne(new QueryWrapper<SysUploadTask>().lambda().eq(SysUploadTask::getFileIdentifier, identifier)); }
合并分片
所有分片上传完毕后,发送合并分片请求给后端
▼java复制代码// Controller /** * 合并分片 * @param identifier * @return */ @PostMapping("/merge/{identifier}") public Result merge (@PathVariable("identifier") String identifier) { sysUploadTaskService.merge(identifier); return Result.ok(); } // Server @Override public void merge(String identifier) { SysUploadTask task = getByIdentifier(identifier); if (task == null) { throw new RuntimeException("分片任务不存在"); } ListPartsRequest listPartsRequest = new ListPartsRequest(task.getBucketName(), task.getObjectKey(), task.getUploadId()); PartListing partListing = amazonS3.listParts(listPartsRequest); List<PartSummary> parts = partListing.getParts(); if (!task.getChunkNum().equals(parts.size())) { // 已上传分块数量与记录中的数量不对应,不能合并分块 throw new RuntimeException("分片缺失,请重新上传"); } CompleteMultipartUploadRequest completeMultipartUploadRequest = new CompleteMultipartUploadRequest() .withUploadId(task.getUploadId()) .withKey(task.getObjectKey()) .withBucketName(task.getBucketName()) .withPartETags(parts.stream().map(partSummary -> new PartETag(partSummary.getPartNumber(), partSummary.getETag())).collect(Collectors.toList())); CompleteMultipartUploadResult result = amazonS3.completeMultipartUpload(completeMultipartUploadRequest); }
