Skip to content
Projects
Groups
Snippets
Help
This project
Loading...
Sign in / Register
Toggle navigation
Y
yd-oss
Overview
Overview
Details
Activity
Cycle Analytics
Repository
Repository
Files
Commits
Branches
Tags
Contributors
Graph
Compare
Charts
Issues
0
Issues
0
List
Board
Labels
Milestones
Merge Requests
0
Merge Requests
0
CI / CD
CI / CD
Pipelines
Jobs
Schedules
Charts
Wiki
Wiki
Snippets
Snippets
Members
Collapse sidebar
Close sidebar
Activity
Graph
Charts
Create a new issue
Jobs
Commits
Issue Boards
Open sidebar
xingmin
yd-oss
Commits
aecbbbff
Commit
aecbbbff
authored
Aug 05, 2026
by
zhangxingmin
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
push
parent
ecc6b626
Hide whitespace changes
Inline
Side-by-side
Showing
1 changed file
with
56 additions
and
59 deletions
+56
-59
yd-oss-api/src/main/java/com/yd/oss/api/service/impl/ApiChunkedUploadContextServiceImpl.java
+56
-59
No files found.
yd-oss-api/src/main/java/com/yd/oss/api/service/impl/ApiChunkedUploadContextServiceImpl.java
View file @
aecbbbff
...
@@ -21,7 +21,6 @@ import javax.annotation.Resource;
...
@@ -21,7 +21,6 @@ import javax.annotation.Resource;
import
java.time.LocalDateTime
;
import
java.time.LocalDateTime
;
import
java.util.*
;
import
java.util.*
;
@Slf4j
@Slf4j
@Service
@Service
public
class
ApiChunkedUploadContextServiceImpl
implements
ApiChunkedUploadContextService
{
public
class
ApiChunkedUploadContextServiceImpl
implements
ApiChunkedUploadContextService
{
...
@@ -48,55 +47,64 @@ public class ApiChunkedUploadContextServiceImpl implements ApiChunkedUploadConte
...
@@ -48,55 +47,64 @@ public class ApiChunkedUploadContextServiceImpl implements ApiChunkedUploadConte
*/
*/
@Override
@Override
@Transactional
(
rollbackFor
=
Exception
.
class
)
@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
)
{
try
{
// ========== 入口日志 ==========
log
.
debug
(
"上传分片: taskId={}, chunkIndex={}, size={}"
,
taskId
,
chunkIndex
,
chunk
.
getSize
());
log
.
info
(
"【uploadChunk】开始处理, taskId={}, chunkIndex={}, projectBizId={}, source={}, fileSize={}"
,
taskId
,
chunkIndex
,
projectBizId
,
source
,
chunk
!=
null
?
chunk
.
getSize
()
:
0
);
// 根据项目ID获取服务商信息
try
{
// 1. 获取 OSS 服务商
OssProvider
provider
=
ossProviderService
.
getProviderByProjectId
(
projectBizId
);
OssProvider
provider
=
ossProviderService
.
getProviderByProjectId
(
projectBizId
);
if
(
provider
==
null
)
{
if
(
provider
==
null
)
{
log
.
error
(
"
未找到项目对应的OSS服务商,
projectBizId={}"
,
projectBizId
);
log
.
error
(
"
【uploadChunk】未找到OSS服务商,
projectBizId={}"
,
projectBizId
);
throw
new
BusinessException
(
"未找到对应的OSS服务商配置"
);
throw
new
BusinessException
(
"未找到对应的OSS服务商配置"
);
}
}
//创建OSS客户端
log
.
info
(
"【uploadChunk】获取到OSS服务商: name={}, bucket={}, endpoint={}"
,
provider
.
getName
(),
provider
.
getBucketName
(),
provider
.
getEndpoint
());
// 2. 创建 OSS 客户端
OSS
ossClient
=
ossClientFactory
.
createOssClient
(
provider
);
OSS
ossClient
=
ossClientFactory
.
createOssClient
(
provider
);
log
.
info
(
"【uploadChunk】OSS客户端创建成功"
);
//
1. 从数据库查询上下文,若不存在则新建
//
3. 查询或创建上传上下文
ChunkedUploadContext
context
=
contextMapper
.
selectOne
(
ChunkedUploadContext
context
=
contextMapper
.
selectOne
(
new
LambdaQueryWrapper
<
ChunkedUploadContext
>()
new
LambdaQueryWrapper
<
ChunkedUploadContext
>()
.
eq
(
ChunkedUploadContext:
:
getTaskId
,
taskId
)
.
eq
(
ChunkedUploadContext:
:
getTaskId
,
taskId
)
);
);
if
(
context
==
null
)
{
if
(
context
==
null
)
{
//首次上传,创建新上下文(初始化 OSS 分片上传)
log
.
info
(
"【uploadChunk】上下文不存在,将初始化新的分片上传, taskId={}"
,
taskId
);
context
=
new
ChunkedUploadContext
();
context
=
new
ChunkedUploadContext
();
context
.
setTaskId
(
taskId
);
context
.
setTaskId
(
taskId
);
// 生成 OSS 对象 Key
结构:sharding+分片来源+年+月+日+xxx.webm
// 生成 OSS 对象 Key
String
objectKey
=
String
.
format
(
"sharding/"
+
source
+
"/%tY/%tm/%s_%d.webm"
,
String
objectKey
=
String
.
format
(
"sharding/"
+
source
+
"/%tY/%tm/%s_%d.webm"
,
new
Date
(),
new
Date
(),
taskId
,
System
.
currentTimeMillis
());
new
Date
(),
new
Date
(),
taskId
,
System
.
currentTimeMillis
());
context
.
setObjectKey
(
objectKey
);
context
.
setObjectKey
(
objectKey
);
log
.
info
(
"【uploadChunk】生成 objectKey: {}"
,
objectKey
);
// 调用 OSS 初始化分片上传
// 调用 OSS 初始化分片上传
InitiateMultipartUploadRequest
initRequest
=
new
InitiateMultipartUploadRequest
(
provider
.
getBucketName
(),
objectKey
);
InitiateMultipartUploadRequest
initRequest
=
new
InitiateMultipartUploadRequest
(
provider
.
getBucketName
(),
objectKey
);
InitiateMultipartUploadResult
initResult
=
ossClient
.
initiateMultipartUpload
(
initRequest
);
InitiateMultipartUploadResult
initResult
=
ossClient
.
initiateMultipartUpload
(
initRequest
);
//OSS 分片上传 ID(由 OSS 初始化时返回)
context
.
setUploadId
(
initResult
.
getUploadId
());
context
.
setUploadId
(
initResult
.
getUploadId
());
// 初始化分片列表为空
context
.
setPartEtagsJson
(
"[]"
);
context
.
setPartEtagsJson
(
"[]"
);
context
.
setStatus
(
1
);
//
状态:
上传中
context
.
setStatus
(
1
);
// 上传中
// 插入数据库
// 插入数据库
contextMapper
.
insert
(
context
);
int
insertResult
=
contextMapper
.
insert
(
context
);
log
.
info
(
"初始化分片上传并入库: taskId={}, uploadId={}, objectKey={}"
,
log
.
info
(
"【uploadChunk】初始化分片上传并入库, taskId={}, uploadId={}, objectKey={}, insertResult={}"
,
taskId
,
context
.
getUploadId
(),
objectKey
);
taskId
,
context
.
getUploadId
(),
objectKey
,
insertResult
);
}
else
{
log
.
info
(
"【uploadChunk】找到已存在的上下文, taskId={}, uploadId={}, status={}, partCount={}"
,
taskId
,
context
.
getUploadId
(),
context
.
getStatus
(),
context
.
getPartEtagsJson
()
!=
null
?
context
.
getPartEtagsJson
().
length
()
:
0
);
}
}
//
2. 当前分片序号转 OSS partNumber(从 1 开始)
//
4. 上传当前分片
int
partNumber
=
chunkIndex
+
1
;
int
partNumber
=
chunkIndex
+
1
;
log
.
info
(
"【uploadChunk】开始上传分片, taskId={}, partNumber={}, chunkSize={}"
,
taskId
,
partNumber
,
chunk
.
getSize
());
// 3. 上传分片到 OSS
UploadPartRequest
uploadPartRequest
=
new
UploadPartRequest
();
UploadPartRequest
uploadPartRequest
=
new
UploadPartRequest
();
uploadPartRequest
.
setBucketName
(
provider
.
getBucketName
());
uploadPartRequest
.
setBucketName
(
provider
.
getBucketName
());
uploadPartRequest
.
setKey
(
context
.
getObjectKey
());
uploadPartRequest
.
setKey
(
context
.
getObjectKey
());
...
@@ -106,27 +114,22 @@ public class ApiChunkedUploadContextServiceImpl implements ApiChunkedUploadConte
...
@@ -106,27 +114,22 @@ public class ApiChunkedUploadContextServiceImpl implements ApiChunkedUploadConte
uploadPartRequest
.
setPartSize
(
chunk
.
getSize
());
uploadPartRequest
.
setPartSize
(
chunk
.
getSize
());
UploadPartResult
uploadResult
=
ossClient
.
uploadPart
(
uploadPartRequest
);
UploadPartResult
uploadResult
=
ossClient
.
uploadPart
(
uploadPartRequest
);
//
4. 构建新的 PartETag
//
5. 更新 ETag 列表
PartETag
partETag
=
new
PartETag
(
uploadResult
.
getPartNumber
(),
uploadResult
.
getETag
());
PartETag
partETag
=
new
PartETag
(
uploadResult
.
getPartNumber
(),
uploadResult
.
getETag
());
// 5. 从数据库读取现有的 ETag 列表(JSON -> List)
String
json
=
context
.
getPartEtagsJson
();
String
json
=
context
.
getPartEtagsJson
();
List
<
PartETag
>
partETags
=
objectMapper
.
readValue
(
json
,
List
<
PartETag
>
partETags
=
objectMapper
.
readValue
(
json
,
new
TypeReference
<
List
<
PartETag
>>()
{});
new
TypeReference
<
List
<
PartETag
>>()
{});
// 6. 添加新 ETag(按 partNumber 升序排序)
partETags
.
add
(
partETag
);
partETags
.
add
(
partETag
);
partETags
.
sort
(
Comparator
.
comparingInt
(
PartETag:
:
getPartNumber
));
partETags
.
sort
(
Comparator
.
comparingInt
(
PartETag:
:
getPartNumber
));
// 7. 序列化为 JSON 并更新数据库
context
.
setPartEtagsJson
(
objectMapper
.
writeValueAsString
(
partETags
));
context
.
setPartEtagsJson
(
objectMapper
.
writeValueAsString
(
partETags
));
context
.
setUpdateTime
(
LocalDateTime
.
now
());
context
.
setUpdateTime
(
LocalDateTime
.
now
());
contextMapper
.
updateById
(
context
);
contextMapper
.
updateById
(
context
);
log
.
debug
(
"分片上传成功: taskId={}, partNumber={}, ETag={}"
,
taskId
,
partNumber
,
uploadResult
.
getETag
());
log
.
info
(
"【uploadChunk】分片上传成功, taskId={}, partNumber={}, ETag={}, 当前总分片数={}"
,
taskId
,
partNumber
,
uploadResult
.
getETag
(),
partETags
.
size
());
}
catch
(
Exception
e
)
{
}
catch
(
Exception
e
)
{
log
.
error
(
"分片上传失败"
,
e
);
log
.
error
(
"【uploadChunk】分片上传失败, taskId={}, chunkIndex={}, projectBizId={}"
,
taskId
,
chunkIndex
,
projectBizId
,
e
);
throw
new
BusinessException
(
"分片上传失败: "
+
e
.
getMessage
());
throw
new
BusinessException
(
"分片上传失败: "
+
e
.
getMessage
());
}
}
}
}
...
@@ -139,17 +142,20 @@ public class ApiChunkedUploadContextServiceImpl implements ApiChunkedUploadConte
...
@@ -139,17 +142,20 @@ public class ApiChunkedUploadContextServiceImpl implements ApiChunkedUploadConte
*/
*/
@Override
@Override
@Transactional
(
rollbackFor
=
Exception
.
class
)
@Transactional
(
rollbackFor
=
Exception
.
class
)
public
Result
<
Map
<
String
,
Object
>>
finishChunks
(
String
taskId
,
String
projectBizId
)
{
public
Result
<
Map
<
String
,
Object
>>
finishChunks
(
String
taskId
,
String
projectBizId
)
{
log
.
info
(
"
完成分片上传并合并文件: taskId={}"
,
task
Id
);
log
.
info
(
"
【finishChunks】开始完成分片上传, taskId={}, projectBizId={}"
,
taskId
,
projectBiz
Id
);
// 1.
从数据库
查询上下文
// 1. 查询上下文
ChunkedUploadContext
context
=
contextMapper
.
selectOne
(
ChunkedUploadContext
context
=
contextMapper
.
selectOne
(
new
LambdaQueryWrapper
<
ChunkedUploadContext
>()
new
LambdaQueryWrapper
<
ChunkedUploadContext
>()
.
eq
(
ChunkedUploadContext:
:
getTaskId
,
taskId
)
.
eq
(
ChunkedUploadContext:
:
getTaskId
,
taskId
)
);
);
if
(
context
==
null
)
{
if
(
context
==
null
)
{
log
.
error
(
"【finishChunks】未找到分片上传上下文, taskId={}"
,
taskId
);
throw
new
BusinessException
(
"未找到分片上传上下文,请检查 taskId 是否正确"
);
throw
new
BusinessException
(
"未找到分片上传上下文,请检查 taskId 是否正确"
);
}
}
log
.
info
(
"【finishChunks】找到上下文, taskId={}, uploadId={}, status={}, objectKey={}"
,
taskId
,
context
.
getUploadId
(),
context
.
getStatus
(),
context
.
getObjectKey
());
// 2. 解析分片列表
// 2. 解析分片列表
String
json
=
context
.
getPartEtagsJson
();
String
json
=
context
.
getPartEtagsJson
();
...
@@ -157,23 +163,26 @@ public class ApiChunkedUploadContextServiceImpl implements ApiChunkedUploadConte
...
@@ -157,23 +163,26 @@ public class ApiChunkedUploadContextServiceImpl implements ApiChunkedUploadConte
try
{
try
{
partETags
=
objectMapper
.
readValue
(
json
,
new
TypeReference
<
List
<
PartETag
>>()
{});
partETags
=
objectMapper
.
readValue
(
json
,
new
TypeReference
<
List
<
PartETag
>>()
{});
}
catch
(
Exception
e
)
{
}
catch
(
Exception
e
)
{
log
.
error
(
"
解析分片 ETag 列表失败"
,
e
);
log
.
error
(
"
【finishChunks】解析分片 ETag 列表失败, taskId={}, json={}"
,
taskId
,
json
,
e
);
throw
new
BusinessException
(
"数据异常,无法合并"
);
throw
new
BusinessException
(
"数据异常,无法合并"
);
}
}
if
(
partETags
.
isEmpty
())
{
if
(
partETags
.
isEmpty
())
{
log
.
error
(
"【finishChunks】没有分片数据, taskId={}"
,
taskId
);
throw
new
BusinessException
(
"没有分片数据,无法合并"
);
throw
new
BusinessException
(
"没有分片数据,无法合并"
);
}
}
log
.
info
(
"【finishChunks】解析到 {} 个分片, taskId={}"
,
partETags
.
size
(),
taskId
);
//
根据项目ID获取服务商信息
//
3. 获取 OSS 服务商
OssProvider
provider
=
ossProviderService
.
getProviderByProjectId
(
projectBizId
);
OssProvider
provider
=
ossProviderService
.
getProviderByProjectId
(
projectBizId
);
if
(
provider
==
null
)
{
if
(
provider
==
null
)
{
log
.
error
(
"
未找到项目对应的OSS服务商,
projectBizId={}"
,
projectBizId
);
log
.
error
(
"
【finishChunks】未找到OSS服务商,
projectBizId={}"
,
projectBizId
);
throw
new
BusinessException
(
"未找到对应的OSS服务商配置"
);
throw
new
BusinessException
(
"未找到对应的OSS服务商配置"
);
}
}
//创建OSS客户端
log
.
info
(
"【finishChunks】获取到OSS服务商: name={}, bucket={}"
,
provider
.
getName
(),
provider
.
getBucketName
());
OSS
ossClient
=
ossClientFactory
.
createOssClient
(
provider
);
OSS
ossClient
=
ossClientFactory
.
createOssClient
(
provider
);
//
3. 调用 OSS 完成分片
合并
//
4. 执行
合并
CompleteMultipartUploadRequest
completeRequest
=
new
CompleteMultipartUploadRequest
(
CompleteMultipartUploadRequest
completeRequest
=
new
CompleteMultipartUploadRequest
(
provider
.
getBucketName
(),
provider
.
getBucketName
(),
context
.
getObjectKey
(),
context
.
getObjectKey
(),
...
@@ -181,43 +190,30 @@ public class ApiChunkedUploadContextServiceImpl implements ApiChunkedUploadConte
...
@@ -181,43 +190,30 @@ public class ApiChunkedUploadContextServiceImpl implements ApiChunkedUploadConte
partETags
partETags
);
);
CompleteMultipartUploadResult
completeResult
=
ossClient
.
completeMultipartUpload
(
completeRequest
);
CompleteMultipartUploadResult
completeResult
=
ossClient
.
completeMultipartUpload
(
completeRequest
);
log
.
info
(
"【finishChunks】合并成功, taskId={}, location={}"
,
taskId
,
completeResult
.
getLocation
());
//
4. 构建文件访问 URL
//
5. 构建文件 URL 和大小
String
fileUrl
=
String
.
format
(
"https://%s.%s/%s"
,
String
fileUrl
=
String
.
format
(
"https://%s.%s/%s"
,
provider
.
getBucketName
(),
provider
.
getBucketName
(),
provider
.
getEndpoint
().
replace
(
"https://"
,
""
),
provider
.
getEndpoint
().
replace
(
"https://"
,
""
),
context
.
getObjectKey
()
context
.
getObjectKey
()
);
);
// 5. 获取文件大小
ObjectMetadata
metadata
=
ossClient
.
getObjectMetadata
(
provider
.
getBucketName
(),
context
.
getObjectKey
());
ObjectMetadata
metadata
=
ossClient
.
getObjectMetadata
(
provider
.
getBucketName
(),
context
.
getObjectKey
());
long
fileSize
=
metadata
.
getContentLength
();
long
fileSize
=
metadata
.
getContentLength
();
// 6. 更新录制任务表(状态变为已录制)
// 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
);
// 已完成
context
.
setStatus
(
2
);
// 已完成
contextMapper
.
updateById
(
context
);
contextMapper
.
updateById
(
context
);
log
.
info
(
"【finishChunks】上下文状态更新为已完成, taskId={}"
,
taskId
);
//
8
. 返回结果
//
7
. 返回结果
Map
<
String
,
Object
>
result
=
new
HashMap
<>();
Map
<
String
,
Object
>
result
=
new
HashMap
<>();
result
.
put
(
"taskId"
,
taskId
);
result
.
put
(
"taskId"
,
taskId
);
result
.
put
(
"fileUrl"
,
fileUrl
);
result
.
put
(
"fileUrl"
,
fileUrl
);
result
.
put
(
"fileKey"
,
context
.
getObjectKey
());
result
.
put
(
"fileKey"
,
context
.
getObjectKey
());
result
.
put
(
"fileSize"
,
fileSize
);
result
.
put
(
"fileSize"
,
fileSize
);
log
.
info
(
"【finishChunks】完成, taskId={}, fileUrl={}"
,
taskId
,
fileUrl
);
return
Result
.
success
(
result
);
return
Result
.
success
(
result
);
}
}
}
}
\ No newline at end of file
Write
Preview
Markdown
is supported
0%
Try again
or
attach a new file
Attach a file
Cancel
You are about to add
0
people
to the discussion. Proceed with caution.
Finish editing this message first!
Cancel
Please
register
or
sign in
to comment