[fix]修改视频上传

This commit is contained in:
wangqx 2025-09-05 17:58:30 +08:00
parent 2616b5fff8
commit bf728cee9b
4 changed files with 45 additions and 22 deletions

View File

@ -34,4 +34,6 @@ public class RocketMQConfig {
* 视频上传主题 * 视频上传主题
*/ */
public static final String VLOG_UPLOAD_TOPIC = "VLOG_UPLOAD_TOPIC"; public static final String VLOG_UPLOAD_TOPIC = "VLOG_UPLOAD_TOPIC";
public static final String VLOG_UPLOAD_GROUP = "VLOG_UPLOAD_GROUP";
} }

View File

@ -230,6 +230,29 @@ public class OssClient {
downloadFile.completionFuture().join(); downloadFile.completionFuture().join();
return tempFilePath; return tempFilePath;
} }
/**
* 下载文件从 Amazon S3 到临时目录
*
* @param path 文件在 Amazon S3 中的对象键
* @return 下载后的文件在本地的临时路径
* @throws OssException 如果下载失败抛出自定义异常
*/
public Path fileDownload(String path,String tempFile) {
// 构建临时文件
Path tempFilePath = FileUtils.createTempFile(tempFile,true).toPath();
// 使用 S3TransferManager 下载文件
FileDownload downloadFile = transferManager.downloadFile(
x -> x.getObjectRequest(
y -> y.bucket(properties.getBucketName())
.key(removeBaseUrl(path))
.build())
.addTransferListener(LoggingTransferListener.create())
.destination(tempFilePath)
.build());
// 等待文件下载操作完成
downloadFile.completionFuture().join();
return tempFilePath;
}
/** /**
* 下载文件从 Amazon S3 输出流 * 下载文件从 Amazon S3 输出流

View File

@ -34,7 +34,7 @@ import java.nio.file.Path;
@RequiredArgsConstructor @RequiredArgsConstructor
@RocketMQMessageListener( @RocketMQMessageListener(
topic = RocketMQConfig.VLOG_UPLOAD_TOPIC, topic = RocketMQConfig.VLOG_UPLOAD_TOPIC,
consumerGroup = RocketMQConfig.CONSUMER_GROUP_SYS_MSG, consumerGroup = RocketMQConfig.VLOG_UPLOAD_GROUP,
selectorExpression = "upload" selectorExpression = "upload"
) )
public class VlogUploadMessageConsumer implements RocketMQListener<MessageExt> { public class VlogUploadMessageConsumer implements RocketMQListener<MessageExt> {
@ -44,12 +44,14 @@ public class VlogUploadMessageConsumer implements RocketMQListener<MessageExt> {
@Override @Override
public void onMessage(MessageExt messageExt) { // 参数为MessageExt public void onMessage(MessageExt messageExt) { // 参数为MessageExt
String message = new String(messageExt.getBody()); String message = new String(messageExt.getBody());
log.info("接收到RocketMQ消息: {}, msgId: {}", message, messageExt.getMsgId()); log.info("接收到RocketMQ消息: {}, msgId: {}", message, messageExt.getMsgId());
try { try {
MQMessage mqMessage = JsonUtils.parseObject(message, MQMessage.class); MQMessage mqMessage = JsonUtils.parseObject(message, MQMessage.class);
if (mqMessage != null) {
if (mqMessage != null&&messageExt.getBody()!=null) {
boolean result = this.handleMessage(mqMessage); boolean result = this.handleMessage(mqMessage);
if (result) { if (result) {
// 消息处理成功手动确认 // 消息处理成功手动确认
@ -80,6 +82,9 @@ public class VlogUploadMessageConsumer implements RocketMQListener<MessageExt> {
//视频上传消息会包含视频的内容 //视频上传消息会包含视频的内容
Vlog vlog = vlogService.getById((String)mqMessage.getData()); Vlog vlog = vlogService.getById((String)mqMessage.getData());
//获取文件内容 //获取文件内容
if(ObjectUtil.isNull(vlog)){
return false;
}
String fileId=vlog.getFileId(); String fileId=vlog.getFileId();
//检查该文件是否为oss文件如果不是说明已经上传过了 //检查该文件是否为oss文件如果不是说明已经上传过了
if(!vlog.getUrl().contains("#")){ if(!vlog.getUrl().contains("#")){
@ -89,12 +94,14 @@ public class VlogUploadMessageConsumer implements RocketMQListener<MessageExt> {
String storagePath = "/data/vlogdata"; String storagePath = "/data/vlogdata";
//从oss下载文件 //从oss下载文件
SysOssVo sysOss =ossService.getById(Long.getLong(fileId)) ; SysOssVo sysOss =ossService.getById(Long.parseLong(fileId)) ;
if (ObjectUtil.isNull(sysOss)) { if (ObjectUtil.isNull(sysOss)) {
throw new ServiceException("文件数据不存在!"); throw new ServiceException("文件数据不存在!");
} }
OssClient ossClient = OssFactory.instance(sysOss.getService()); OssClient ossClient = OssFactory.instance(sysOss.getService());
Path filePath=ossClient.fileDownload(sysOss.getUrl()); Path filePath=ossClient.fileDownload(sysOss.getUrl(),sysOss.getOriginalName());
//对文件进行重命名
fileId = qcCloud.uploadViaTempFile(filePath); fileId = qcCloud.uploadViaTempFile(filePath);
log.info("视频发布ID:" + fileId); log.info("视频发布ID:" + fileId);

View File

@ -161,24 +161,15 @@ public class VlogServiceImpl extends ServiceImpl<VlogMapper, Vlog> implements Vl
//发出mq消息异步处理上传 //发出mq消息异步处理上传
// MQMessage message = MQMessage.builder() MQMessage message = MQMessage.builder()
// .messageType("json") .messageType("json")
// .data(vlog.getId()) .data(vlog.getId())
// .source("app") .source("app")
// .topic("VLOG_UPLOAD_TOPIC") .topic("VLOG_UPLOAD_TOPIC")
// .tag("upload") .tag("upload")
// .sendTime(LocalDateTime.now()) .sendTime(LocalDateTime.now())
// .build(); .build();
// MQMessage message = MQMessage.builder() MqUtil.sendMessage(RocketMQConfig.VLOG_UPLOAD_TOPIC+":"+"upload", message);
// .topic("VLOG_UPLOAD_TOPIC")
// .tag("upload")
// .messageType("json")
// .data("d")
// .source("vlog_service")
// .sendTime(LocalDateTime.now())
// .build();
MQMessage mqMessage = new MQMessage();
MqUtil.sendMessage(RocketMQConfig.VLOG_UPLOAD_TOPIC+":"+"upload", mqMessage);
} }