Commit 4348c14c by zhangxingmin

push

parent 635a68aa
......@@ -14,12 +14,15 @@ 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.redisson.api.RLock;
import org.redisson.api.RedissonClient;
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.*;
import java.util.concurrent.TimeUnit;
@Slf4j
@Service
......@@ -37,39 +40,58 @@ public class ApiChunkedUploadContextServiceImpl implements ApiChunkedUploadConte
@Resource
private IOssProviderService ossProviderService;
@Resource
private RedissonClient redissonClient;
private static final String LOCK_KEY_PREFIX = "chunked:upload:lock:";
/**
* 上传单个分片
* @param taskId
* @param chunkIndex
* @param chunk
* @param projectBizId
* @param source
* 上传单个分片(使用分布式锁保证上下文唯一性)
* @param taskId 录制任务ID
* @param chunkIndex 分片索引(从0开始)
* @param chunk 分片文件
* @param projectBizId 项目业务ID(为空时使用默认OSS)
* @param source 来源标识(如 'recording')
*/
@Override
@Transactional(rollbackFor = Exception.class)
public void uploadChunk(String taskId, Integer chunkIndex, MultipartFile chunk, String projectBizId, String source) {
// ========== 入口日志 ==========
public void uploadChunk(String taskId, Integer chunkIndex, MultipartFile chunk,
String projectBizId, String source) {
// ---------- 前置准备 ----------
log.info("【uploadChunk】开始处理, taskId={}, chunkIndex={}, projectBizId={}, source={}, fileSize={}",
taskId, chunkIndex, projectBizId, source, chunk != null ? chunk.getSize() : 0);
// 获取 OSS 服务商(不依赖上下文,提前获取)
OssProvider provider = ossProviderService.getProviderByProjectId(projectBizId);
if (provider == null) {
log.error("【uploadChunk】未找到OSS服务商, projectBizId={}", projectBizId);
throw new BusinessException("未找到对应的OSS服务商配置");
}
log.info("【uploadChunk】获取到OSS服务商: name={}, bucket={}", provider.getName(), provider.getBucketName());
// 创建 OSS 客户端(提前创建,供后续使用)
OSS ossClient = ossClientFactory.createOssClient(provider);
log.info("【uploadChunk】OSS客户端创建成功");
// ---------- 分布式锁(保护整个方法,确保同一taskId操作串行) ----------
String lockKey = LOCK_KEY_PREFIX + taskId;
RLock lock = redissonClient.getLock(lockKey);
boolean locked = false;
try {
// 1. 获取 OSS 服务商
OssProvider provider = ossProviderService.getProviderByProjectId(projectBizId);
if (provider == null) {
log.error("【uploadChunk】未找到OSS服务商, projectBizId={}", projectBizId);
throw new BusinessException("未找到对应的OSS服务商配置");
// 尝试获取锁,等待5秒,锁持有时间30秒(足够处理上传+DB更新)
locked = lock.tryLock(5, 30, TimeUnit.SECONDS);
if (!locked) {
log.error("【uploadChunk】获取分布式锁失败, taskId={}", taskId);
throw new BusinessException("系统繁忙,请稍后重试");
}
log.info("【uploadChunk】获取到OSS服务商: name={}, bucket={}, endpoint={}",
provider.getName(), provider.getBucketName(), provider.getEndpoint());
// 2. 创建 OSS 客户端
OSS ossClient = ossClientFactory.createOssClient(provider);
log.info("【uploadChunk】OSS客户端创建成功");
log.debug("【uploadChunk】获取锁成功, taskId={}", taskId);
// 3. 查询或创建上传上下文
// ---------- 1. 查询或创建上传上下文(原子操作) ----------
ChunkedUploadContext context = contextMapper.selectOne(
new LambdaQueryWrapper<ChunkedUploadContext>()
.eq(ChunkedUploadContext::getTaskId, taskId)
.last(" limit 1 ")
);
if (context == null) {
......@@ -84,13 +106,14 @@ public class ApiChunkedUploadContextServiceImpl implements ApiChunkedUploadConte
log.info("【uploadChunk】生成 objectKey: {}", objectKey);
// 调用 OSS 初始化分片上传
InitiateMultipartUploadRequest initRequest = new InitiateMultipartUploadRequest(provider.getBucketName(), objectKey);
InitiateMultipartUploadRequest initRequest =
new InitiateMultipartUploadRequest(provider.getBucketName(), objectKey);
InitiateMultipartUploadResult initResult = ossClient.initiateMultipartUpload(initRequest);
context.setUploadId(initResult.getUploadId());
context.setPartEtagsJson("[]");
context.setStatus(1); // 上传中
// 插入数据库
// 插入数据库(此时其他线程被锁阻挡,不会重复插入)
int insertResult = contextMapper.insert(context);
log.info("【uploadChunk】初始化分片上传并入库, taskId={}, uploadId={}, objectKey={}, insertResult={}",
taskId, context.getUploadId(), objectKey, insertResult);
......@@ -100,7 +123,7 @@ public class ApiChunkedUploadContextServiceImpl implements ApiChunkedUploadConte
context.getPartEtagsJson() != null ? context.getPartEtagsJson().length() : 0);
}
// 4. 上传当前分片
// ---------- 2. 上传当前分片 ----------
int partNumber = chunkIndex + 1;
log.info("【uploadChunk】开始上传分片, taskId={}, partNumber={}, chunkSize={}",
taskId, partNumber, chunk.getSize());
......@@ -114,7 +137,7 @@ public class ApiChunkedUploadContextServiceImpl implements ApiChunkedUploadConte
uploadPartRequest.setPartSize(chunk.getSize());
UploadPartResult uploadResult = ossClient.uploadPart(uploadPartRequest);
// 5. 更新 ETag 列表
// ---------- 3. 更新 ETag 列表(在锁保护下,不会并发覆盖) ----------
PartETag partETag = new PartETag(uploadResult.getPartNumber(), uploadResult.getETag());
String json = context.getPartEtagsJson();
List<PartETag> partETags = objectMapper.readValue(json, new TypeReference<List<PartETag>>() {});
......@@ -127,10 +150,19 @@ public class ApiChunkedUploadContextServiceImpl implements ApiChunkedUploadConte
log.info("【uploadChunk】分片上传成功, taskId={}, partNumber={}, ETag={}, 当前总分片数={}",
taskId, partNumber, uploadResult.getETag(), partETags.size());
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
log.error("【uploadChunk】获取锁被中断, taskId={}", taskId, e);
throw new BusinessException("系统中断,请重试");
} catch (Exception e) {
log.error("【uploadChunk】分片上传失败, taskId={}, chunkIndex={}, projectBizId={}",
taskId, chunkIndex, projectBizId, e);
log.error("【uploadChunk】处理失败, taskId={}, chunkIndex={}", taskId, chunkIndex, e);
throw new BusinessException("分片上传失败: " + e.getMessage());
} finally {
// 释放锁(仅当当前线程持有锁时)
if (locked && lock.isHeldByCurrentThread()) {
lock.unlock();
log.debug("【uploadChunk】释放分布式锁, taskId={}", taskId);
}
}
}
......@@ -143,6 +175,8 @@ public class ApiChunkedUploadContextServiceImpl implements ApiChunkedUploadConte
@Override
@Transactional(rollbackFor = Exception.class)
public Result<Map<String, Object>> finishChunks(String taskId, String projectBizId) {
// ... 原有逻辑不变(无需加锁,因为合并时上下文已存在且不再变化)
// 但为了安全,也可以加锁确保合并时上下文不会被更新,但此时上传已完成,不需要锁。
log.info("【finishChunks】开始完成分片上传, taskId={}, projectBizId={}", taskId, projectBizId);
// 1. 查询上下文
......
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