Commit 9b48c583 by zhangxingmin

push

parent 1cd4541f
package com.yd.oss.api.controller;
import com.yd.common.result.Result;
import com.yd.oss.api.service.ApiChunkedUploadContextService;
import com.yd.oss.feign.client.ApiChunkedUploadContextFeignClient;
import lombok.extern.slf4j.Slf4j;
import org.springframework.validation.annotation.Validated;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import org.springframework.web.multipart.MultipartFile;
import javax.annotation.Resource;
import java.util.Map;
/**
* 分片上传上下文信息
* (一个任务上传文件信息拆分成多个ETag上传存储,任务执行完毕,阿里云合并ETag列表为完整的文件信息)
* @author zxm
* @since 2026-08-05
*/
@Slf4j
@RestController
@RequestMapping("/chunkedUploadContext")
@Validated
public class ApiChunkedUploadContextController implements ApiChunkedUploadContextFeignClient {
@Resource
private ApiChunkedUploadContextService apiChunkedUploadContextService;
/**
* 上传单个分片
*/
public Result<Void> uploadChunk(String taskId,Integer chunkIndex,
MultipartFile chunk,String projectBizId,
String source) {
apiChunkedUploadContextService.uploadChunk(taskId, chunkIndex, chunk,projectBizId,source);
return Result.success();
}
/**
* 完成分片上传(合并)
*/
public Result<Map<String, Object>> finishChunks(String taskId,
String projectBizId) {
return apiChunkedUploadContextService.finishChunks(taskId,projectBizId);
}
}
package com.yd.oss.api.service;
import com.yd.common.result.Result;
import org.springframework.web.multipart.MultipartFile;
import java.util.Map;
public interface ApiChunkedUploadContextService {
void uploadChunk(String taskId, Integer chunkIndex, MultipartFile chunk,String projectBizId,String source);
Result<Map<String, Object>> finishChunks(String taskId,String projectBizId);
}
package com.yd.oss.api.service.impl;
import com.aliyun.oss.OSS;
import com.aliyun.oss.model.*;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.yd.common.exception.BusinessException;
import com.yd.common.result.Result;
import com.yd.oss.api.service.ApiChunkedUploadContextService;
import com.yd.oss.service.config.OssClientFactory;
import com.yd.oss.service.dao.ChunkedUploadContextMapper;
import com.yd.oss.service.model.ChunkedUploadContext;
import com.yd.oss.service.model.OssProvider;
import com.yd.oss.service.service.IOssProviderService;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import org.springframework.web.multipart.MultipartFile;
import javax.annotation.Resource;
import java.time.LocalDateTime;
import java.util.*;
@Slf4j
@Service
public class ApiChunkedUploadContextServiceImpl implements ApiChunkedUploadContextService {
@Resource
private ChunkedUploadContextMapper contextMapper;
@Resource
private ObjectMapper objectMapper;
@Resource
private OssClientFactory ossClientFactory;
@Resource
private IOssProviderService ossProviderService;
/**
* 上传单个分片
* @param taskId
* @param chunkIndex
* @param chunk
* @param projectBizId
* @param source
*/
@Override
@Transactional(rollbackFor = Exception.class)
public void uploadChunk(String taskId, Integer chunkIndex, MultipartFile chunk, String projectBizId,String source) {
try {
log.debug("上传分片: taskId={}, chunkIndex={}, size={}", taskId, chunkIndex, chunk.getSize());
// 根据项目ID获取服务商信息
OssProvider provider = ossProviderService.getProviderByProjectId(projectBizId);
if (provider == null) {
log.error("未找到项目对应的OSS服务商,projectBizId={}", projectBizId);
throw new BusinessException("未找到对应的OSS服务商配置");
}
//创建OSS客户端
OSS ossClient = ossClientFactory.createOssClient(provider);
// 1. 从数据库查询上下文,若不存在则新建
ChunkedUploadContext context = contextMapper.selectOne(
new LambdaQueryWrapper<ChunkedUploadContext>()
.eq(ChunkedUploadContext::getTaskId, taskId)
);
if (context == null) {
//首次上传,创建新上下文(初始化 OSS 分片上传)
context = new ChunkedUploadContext();
context.setTaskId(taskId);
// 生成 OSS 对象 Key 结构:sharding+分片来源+年+月+日+xxx.webm
String objectKey = String.format("sharding/"+source+"/%tY/%tm/%s_%d.webm",
new Date(), new Date(), taskId, System.currentTimeMillis());
context.setObjectKey(objectKey);
// 调用 OSS 初始化分片上传
InitiateMultipartUploadRequest initRequest = new InitiateMultipartUploadRequest(provider.getBucketName(), objectKey);
InitiateMultipartUploadResult initResult = ossClient.initiateMultipartUpload(initRequest);
//OSS 分片上传 ID(由 OSS 初始化时返回)
context.setUploadId(initResult.getUploadId());
// 初始化分片列表为空
context.setPartEtagsJson("[]");
context.setStatus(1); // 状态:上传中
// 插入数据库
contextMapper.insert(context);
log.info("初始化分片上传并入库: taskId={}, uploadId={}, objectKey={}",
taskId, context.getUploadId(), objectKey);
}
// 2. 当前分片序号转 OSS partNumber(从 1 开始)
int partNumber = chunkIndex + 1;
// 3. 上传分片到 OSS
UploadPartRequest uploadPartRequest = new UploadPartRequest();
uploadPartRequest.setBucketName(provider.getBucketName());
uploadPartRequest.setKey(context.getObjectKey());
uploadPartRequest.setUploadId(context.getUploadId());
uploadPartRequest.setPartNumber(partNumber);
uploadPartRequest.setInputStream(chunk.getInputStream());
uploadPartRequest.setPartSize(chunk.getSize());
UploadPartResult uploadResult = ossClient.uploadPart(uploadPartRequest);
// 4. 构建新的 PartETag
PartETag partETag = new PartETag(uploadResult.getPartNumber(), uploadResult.getETag());
// 5. 从数据库读取现有的 ETag 列表(JSON -> List)
String json = context.getPartEtagsJson();
List<PartETag> partETags = objectMapper.readValue(json,
new TypeReference<List<PartETag>>() {});
// 6. 添加新 ETag(按 partNumber 升序排序)
partETags.add(partETag);
partETags.sort(Comparator.comparingInt(PartETag::getPartNumber));
// 7. 序列化为 JSON 并更新数据库
context.setPartEtagsJson(objectMapper.writeValueAsString(partETags));
context.setUpdateTime(LocalDateTime.now());
contextMapper.updateById(context);
log.debug("分片上传成功: taskId={}, partNumber={}, ETag={}", taskId, partNumber, uploadResult.getETag());
} catch (Exception e) {
log.error("分片上传失败", e);
throw new BusinessException("分片上传失败: " + e.getMessage());
}
}
/**
* 完成分片上传(合并)
* @param taskId
* @param projectBizId
* @return
*/
@Override
@Transactional(rollbackFor = Exception.class)
public Result<Map<String, Object>> finishChunks(String taskId,String projectBizId) {
log.info("完成分片上传并合并文件: taskId={}", taskId);
// 1. 从数据库查询上下文
ChunkedUploadContext context = contextMapper.selectOne(
new LambdaQueryWrapper<ChunkedUploadContext>()
.eq(ChunkedUploadContext::getTaskId, taskId)
);
if (context == null) {
throw new BusinessException("未找到分片上传上下文,请检查 taskId 是否正确");
}
// 2. 解析分片列表
String json = context.getPartEtagsJson();
List<PartETag> partETags;
try {
partETags = objectMapper.readValue(json, new TypeReference<List<PartETag>>() {});
} catch (Exception e) {
log.error("解析分片 ETag 列表失败", e);
throw new BusinessException("数据异常,无法合并");
}
if (partETags.isEmpty()) {
throw new BusinessException("没有分片数据,无法合并");
}
// 根据项目ID获取服务商信息
OssProvider provider = ossProviderService.getProviderByProjectId(projectBizId);
if (provider == null) {
log.error("未找到项目对应的OSS服务商,projectBizId={}", projectBizId);
throw new BusinessException("未找到对应的OSS服务商配置");
}
//创建OSS客户端
OSS ossClient = ossClientFactory.createOssClient(provider);
// 3. 调用 OSS 完成分片合并
CompleteMultipartUploadRequest completeRequest = new CompleteMultipartUploadRequest(
provider.getBucketName(),
context.getObjectKey(),
context.getUploadId(),
partETags
);
CompleteMultipartUploadResult completeResult = ossClient.completeMultipartUpload(completeRequest);
// 4. 构建文件访问 URL
String fileUrl = String.format("https://%s.%s/%s",
provider.getBucketName(),
provider.getEndpoint().replace("https://", ""),
context.getObjectKey()
);
// 5. 获取文件大小
ObjectMetadata metadata = ossClient.getObjectMetadata(provider.getBucketName(), context.getObjectKey());
long fileSize = metadata.getContentLength();
// 6. 更新录制任务表(状态变为已录制)
// LambdaQueryWrapper<RecordingTask> wrapper = new LambdaQueryWrapper<>();
// wrapper.eq(RecordingTask::getTaskId, taskId);
// RecordingTask task = iRecordingTaskService.getOne(wrapper);
// if (task == null) {
// throw new BusinessException("录制任务不存在,taskId=" + taskId);
// }
// task.setFileUrl(fileUrl);
// task.setFileSize(fileSize);
// task.setStatus("3"); // 已录制
// task.setFileFormat("webm");
// task.setStopTime(LocalDateTime.now());
// iRecordingTaskService.updateById(task);
// 7. 更新上下文状态为“已完成”
context.setStatus(2); // 已完成
contextMapper.updateById(context);
// 8. 返回结果
Map<String, Object> result = new HashMap<>();
result.put("taskId", taskId);
result.put("fileUrl", fileUrl);
result.put("fileKey", context.getObjectKey());
result.put("fileSize", fileSize);
return Result.success(result);
}
}
package com.yd.oss.feign.client;
import com.yd.common.result.Result;
import com.yd.oss.feign.fallback.ApiChunkedUploadContextFeignFallbackFactory;
import org.springframework.cloud.openfeign.FeignClient;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.multipart.MultipartFile;
import java.util.Map;
/**
* 分片上传上下文信息Feign客户端
*/
@FeignClient(name = "yd-oss-api",path = "/oss/api/chunkedUploadContext",fallbackFactory = ApiChunkedUploadContextFeignFallbackFactory.class)
public interface ApiChunkedUploadContextFeignClient {
/**
* 上传单个分片
* @param taskId
* @param chunkIndex
* @param chunk
* @param projectBizId
* @param source
* @return
*/
@PostMapping("/uploadChunk")
Result<Void> uploadChunk(@RequestParam(value = "taskId",required = false) String taskId,
@RequestParam(value = "chunkIndex",required = false) Integer chunkIndex,
@RequestParam(value = "chunk",required = false) MultipartFile chunk,
@RequestParam(value = "projectBizId",required = false) String projectBizId,
@RequestParam(value = "source",required = false) String source);
/**
* 完成分片上传(合并)
* @param taskId
* @param projectBizId
* @return
*/
@PostMapping("/finishChunks")
Result<Map<String, Object>> finishChunks(@RequestParam(value = "taskId",required = false) String taskId,
@RequestParam(value = "projectBizId",required = false) String projectBizId);
}
package com.yd.oss.feign.fallback;
import com.yd.common.result.Result;
import com.yd.oss.feign.client.ApiChunkedUploadContextFeignClient;
import lombok.extern.slf4j.Slf4j;
import org.springframework.cloud.openfeign.FallbackFactory;
import org.springframework.stereotype.Component;
import org.springframework.web.multipart.MultipartFile;
import java.util.Map;
/**
* 分片上传上下文信息Feign降级处理
*/
@Slf4j
@Component
public class ApiChunkedUploadContextFeignFallbackFactory implements FallbackFactory<ApiChunkedUploadContextFeignClient> {
@Override
public ApiChunkedUploadContextFeignClient create(Throwable cause) {
return new ApiChunkedUploadContextFeignClient() {
@Override
public Result<Void> uploadChunk(String taskId, Integer chunkIndex, MultipartFile chunk, String projectBizId, String source) {
return null;
}
@Override
public Result<Map<String, Object>> finishChunks(String taskId, String projectBizId) {
return null;
}
};
}
}
......@@ -4,7 +4,6 @@ import com.yd.common.result.Result;
import com.yd.oss.feign.client.ApiExcelFeignClient;
import com.yd.oss.feign.dto.ExportResult;
import com.yd.oss.feign.request.ApiExportRequest;
import com.yd.oss.feign.request.ApiOssExcelParseRequest;
import com.yd.oss.feign.request.ApiOssExportAppointmentExcelRequest;
import com.yd.oss.feign.request.MultiSheetExportRequest;
import com.yd.oss.feign.response.ApiOssExcelParseResponse;
......
package com.yd.oss.service.dao;
import com.yd.oss.service.model.ChunkedUploadContext;
import com.baomidou.mybatisplus.core.mapper.BaseMapper;
/**
* <p>
* 分片上传上下文表(一个任务上传文件信息拆分成多个ETag上传存储,任务执行完毕,阿里云合并ETag列表为完整的文件信息) Mapper 接口
* </p>
*
* @author zxm
* @since 2026-08-05
*/
public interface ChunkedUploadContextMapper extends BaseMapper<ChunkedUploadContext> {
}
package com.yd.oss.service.model;
import com.baomidou.mybatisplus.annotation.IdType;
import com.baomidou.mybatisplus.annotation.TableField;
import com.baomidou.mybatisplus.annotation.TableId;
import com.baomidou.mybatisplus.annotation.TableName;
import java.io.Serializable;
import java.time.LocalDateTime;
import lombok.Getter;
import lombok.Setter;
/**
* <p>
* 分片上传上下文表(一个任务上传文件信息拆分成多个ETag上传存储,任务执行完毕,阿里云合并ETag列表为完整的文件信息)
* </p>
*
* @author zxm
* @since 2026-08-05
*/
@Getter
@Setter
@TableName("chunked_upload_context")
public class ChunkedUploadContext implements Serializable {
private static final long serialVersionUID = 1L;
/**
* 主键
*/
@TableId(value = "id", type = IdType.AUTO)
private Long id;
/**
* 任务ID
*/
@TableField("task_id")
private String taskId;
/**
* OSS 分片上传 ID(由 OSS 初始化时返回)
*/
@TableField("upload_id")
private String uploadId;
/**
* 文件在 OSS 中的完整路径(不含 Bucket 名称)
*/
@TableField("object_key")
private String objectKey;
/**
* 已上传分片的 ETag 列表(JSON 数组格式,如 [{"partNumber":1,"eTag":"xxx"}])
*/
@TableField("part_etags_json")
private String partEtagsJson;
/**
* 状态:0-初始化,1-上传中,2-已完成,3-已取消
*/
@TableField("status")
private Integer status;
/**
* 通用备注
*/
@TableField("remark")
private String remark;
/**
* 删除标识: 0-正常, 1-删除
*/
@TableField("is_deleted")
private Integer isDeleted;
/**
* 创建人ID
*/
@TableField("creator_id")
private String creatorId;
/**
* 更新人ID
*/
@TableField("updater_id")
private String updaterId;
/**
* 创建时间
*/
@TableField("create_time")
private LocalDateTime createTime;
/**
* 更新时间
*/
@TableField("update_time")
private LocalDateTime updateTime;
}
package com.yd.oss.service.service;
import com.yd.oss.service.model.ChunkedUploadContext;
import com.baomidou.mybatisplus.extension.service.IService;
/**
* <p>
* 分片上传上下文表(一个任务上传文件信息拆分成多个ETag上传存储,任务执行完毕,阿里云合并ETag列表为完整的文件信息) 服务类
* </p>
*
* @author zxm
* @since 2026-08-05
*/
public interface IChunkedUploadContextService extends IService<ChunkedUploadContext> {
}
package com.yd.oss.service.service.impl;
import com.yd.oss.service.model.ChunkedUploadContext;
import com.yd.oss.service.dao.ChunkedUploadContextMapper;
import com.yd.oss.service.service.IChunkedUploadContextService;
import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl;
import org.springframework.stereotype.Service;
/**
* <p>
* 分片上传上下文表(一个任务上传文件信息拆分成多个ETag上传存储,任务执行完毕,阿里云合并ETag列表为完整的文件信息) 服务实现类
* </p>
*
* @author zxm
* @since 2026-08-05
*/
@Service
public class ChunkedUploadContextServiceImpl extends ServiceImpl<ChunkedUploadContextMapper, ChunkedUploadContext> implements IChunkedUploadContextService {
}
......@@ -21,7 +21,7 @@ public class MyBatisPlusCodeGenerator {
})
.strategyConfig(builder -> {
builder.addInclude(
"material","rel_object_material"
"chunked_upload_context"
)
.entityBuilder()
.enableLombok()
......
<?xml version="1.0" encoding="UTF-8"?>
<!DOCTYPE mapper PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN" "http://mybatis.org/dtd/mybatis-3-mapper.dtd">
<mapper namespace="com.yd.oss.service.dao.ChunkedUploadContextMapper">
</mapper>
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment