头像视频合成

This commit is contained in:
spllzh
2025-10-08 08:02:21 +08:00
parent 76dedb657a
commit 45d061bad7
28 changed files with 689 additions and 13 deletions

View File

@@ -103,6 +103,9 @@ public class AudioStatisticsScheduler {

View File

@@ -0,0 +1,358 @@
package com.rj.scheduler;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONObject;
import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper;
import com.rj.config.AliyunConfig;
import com.rj.entity.VideoSynthesisLog;
import com.rj.mapper.VideoSynthesisLogMapper;
import com.rj.service.MinIOService;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.http.HttpEntity;
import org.springframework.http.HttpHeaders;
import org.springframework.http.HttpMethod;
import org.springframework.http.ResponseEntity;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
import org.springframework.web.client.RestTemplate;
import java.io.ByteArrayInputStream;
import java.io.ByteArrayOutputStream;
import java.io.InputStream;
import java.net.HttpURLConnection;
import java.net.URL;
import java.time.LocalDateTime;
import java.util.List;
/**
* 视频合成任务状态检查定时器
* 每5分钟检查一次PENDING状态的任务,调用阿里云API查询任务状态并更新数据库
*
* @author rj
* @date 2025-01-30
*/
@Slf4j
@Component
public class VideoSynthesisStatusScheduler {
@Autowired
private VideoSynthesisLogMapper videoSynthesisLogMapper;
@Autowired
private AliyunConfig aliyunConfig;
@Autowired
private RestTemplate restTemplate;
@Autowired
private MinIOService minIOService;
/**
* 每5分钟执行一次任务状态检查
*/
@Scheduled(fixedRate = 5 * 60 * 1000) // 5分钟 = 5 * 60 * 1000毫秒
public void checkVideoSynthesisStatus() {
try {
log.info("开始执行视频合成任务状态检查...");
// 查询所有PENDING状态的任务
List<VideoSynthesisLog> pendingTasks = getPendingTasks();
if (pendingTasks.isEmpty()) {
log.info("没有找到PENDING状态的任务");
return;
}
log.info("找到{}个PENDING状态的任务,开始检查状态", pendingTasks.size());
// 遍历每个任务,检查状态
for (VideoSynthesisLog task : pendingTasks) {
try {
checkAndUpdateTaskStatus(task);
} catch (Exception e) {
log.error("检查任务状态失败,任务ID: {}, 错误: {}", task.getTaskId(), e.getMessage(), e);
}
}
log.info("视频合成任务状态检查完成");
} catch (Exception e) {
log.error("视频合成任务状态检查执行失败", e);
}
}
/**
* 查询所有PENDING状态的任务
*/
private List<VideoSynthesisLog> getPendingTasks() {
QueryWrapper<VideoSynthesisLog> queryWrapper = new QueryWrapper<>();
queryWrapper.eq("task_status", "PENDING");
return videoSynthesisLogMapper.selectList(queryWrapper);
}
/**
* 检查并更新单个任务状态
*/
private void checkAndUpdateTaskStatus(VideoSynthesisLog task) {
try {
log.info("检查任务状态,任务ID: {}", task.getTaskId());
if (task == null || task.getTaskId() == null || task.getTaskId().isEmpty()) return;
// 调用阿里云API检查任务状态
String apiUrl = "https://dashscope.aliyuncs.com/api/v1/tasks/" + task.getTaskId();
HttpHeaders headers = new HttpHeaders();
headers.set("Authorization", "Bearer " + aliyunConfig.getApiKey());
headers.set("Content-Type", "application/json");
HttpEntity<String> entity = new HttpEntity<>(headers);
ResponseEntity<String> response = restTemplate.exchange(
apiUrl,
HttpMethod.GET,
entity,
String.class
);
if (response.getStatusCode().is2xxSuccessful()) {
String responseBody = response.getBody();
log.info("API响应: {}", responseBody);
// 解析响应并更新数据库
updateTaskFromResponse(task, responseBody);
} else {
log.error("API调用失败,状态码: {}, 任务ID: {}", response.getStatusCode(), task.getTaskId());
}
} catch (Exception e) {
log.error("检查任务状态异常,任务ID: {}, 错误: {}", task.getTaskId(), e.getMessage(), e);
}
}
/**
* 根据API响应更新任务状态
*/
private void updateTaskFromResponse(VideoSynthesisLog task, String responseBody) {
try {
JSONObject responseJson = JSON.parseObject(responseBody);
if (responseJson.containsKey("output")) {
JSONObject output = responseJson.getJSONObject("output");
// 更新任务状态
String taskStatus = output.getString("task_status");
task.setTaskStatus(taskStatus);
task.setUpdatedAt(LocalDateTime.now());
// 如果任务成功完成,处理视频文件
if ("SUCCEEDED".equals(taskStatus)) {
JSONObject results = output.getJSONObject("results");
if (results != null && results.containsKey("video_url")) {
String originalVideoUrl = results.getString("video_url");
log.info("任务成功完成,原始视频URL: {}", originalVideoUrl);
// 从阿里云OSS下载视频并上传到MinIO
String minioVideoUrl = downloadAndUploadToMinIO(originalVideoUrl, task.getTaskId());
if (minioVideoUrl != null) {
task.setVideoUrl(minioVideoUrl);
log.info("视频已成功转存到MinIO: {}", minioVideoUrl);
} else {
// 如果转存失败,保留原始URL
task.setVideoUrl(originalVideoUrl);
log.warn("视频转存到MinIO失败,保留原始URL: {}", originalVideoUrl);
}
task.setSuccess(true);
task.setResponseTime(LocalDateTime.now());
}
} else if ("FAILED".equals(taskStatus)) {
task.setSuccess(false);
task.setErrorMessage("任务执行失败");
task.setResponseTime(LocalDateTime.now());
}
// 更新usage信息
if (responseJson.containsKey("usage")) {
JSONObject usage = responseJson.getJSONObject("usage");
if (usage.containsKey("video_duration")) {
task.setVideoDuration(usage.getBigDecimal("video_duration"));
}
if (usage.containsKey("video_ratio")) {
task.setVideoRatio(usage.getString("video_ratio"));
}
}
// 保存到数据库
int updateResult = videoSynthesisLogMapper.updateById(task);
if (updateResult > 0) {
log.info("任务状态更新成功,任务ID: {}, 新状态: {}", task.getTaskId(), taskStatus);
} else {
log.error("任务状态更新失败,任务ID: {}", task.getTaskId());
}
} else {
log.error("API响应格式异常,任务ID: {}, 响应: {}", task.getTaskId(), responseBody);
}
} catch (Exception e) {
log.error("解析API响应失败,任务ID: {}, 响应: {}, 错误: {}", task.getTaskId(), responseBody, e.getMessage(), e);
}
}
/**
* 从阿里云OSS下载视频并上传到MinIO
*
* @param ossVideoUrl 阿里云OSS视频URL
* @param taskId 任务ID,用于生成文件名
* @return MinIO中的视频URL,失败时返回null
*/
private String downloadAndUploadToMinIO(String ossVideoUrl, String taskId) {
try {
log.info("开始从阿里云OSS下载视频: {}", ossVideoUrl);
// 从阿里云OSS下载视频
byte[] videoData = downloadVideoFromOSS(ossVideoUrl);
if (videoData == null || videoData.length == 0) {
log.error("从阿里云OSS下载视频失败,数据为空");
return null;
}
log.info("视频下载成功,大小: {} bytes", videoData.length);
// 生成MinIO文件名
String fileName = "video_synthesis/" + taskId + ".mp4";
// 上传到MinIO
String minioUrl = minIOService.uploadFile(
new ByteArrayInputStream(videoData),
fileName,
"video/mp4"
);
log.info("视频已成功上传到MinIO: {}", minioUrl);
return minioUrl;
} catch (Exception e) {
log.error("下载并上传视频到MinIO失败: {}", e.getMessage(), e);
return null;
}
}
/**
* 从阿里云OSS下载视频文件
*
* @param videoUrl 视频URL
* @return 视频字节数组
*/
private byte[] downloadVideoFromOSS(String videoUrl) {
try {
log.info("开始下载视频: {}", videoUrl);
// 方法1: 使用HttpURLConnection下载视频,避免RestTemplate的URL编码问题
byte[] videoData = downloadWithHttpURLConnection(videoUrl);
if (videoData != null) {
return videoData;
}
// 方法2: 如果HttpURLConnection失败,尝试使用RestTemplate但禁用URL编码
log.warn("HttpURLConnection下载失败,尝试使用RestTemplate备用方案");
return downloadWithRestTemplate(videoUrl);
} catch (Exception e) {
log.error("下载视频异常: {}", e.getMessage(), e);
return null;
}
}
/**
* 使用HttpURLConnection下载视频
*/
private byte[] downloadWithHttpURLConnection(String videoUrl) {
try {
URL url = new URL(videoUrl);
HttpURLConnection connection = (HttpURLConnection) url.openConnection();
// 设置请求头
connection.setRequestMethod("GET");
connection.setConnectTimeout(30000); // 30秒连接超时
connection.setReadTimeout(300000); // 5分钟读取超时
connection.setRequestProperty("User-Agent", "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36");
connection.setRequestProperty("Accept", "*/*");
connection.setRequestProperty("Accept-Encoding", "identity"); // 禁用压缩
int responseCode = connection.getResponseCode();
if (responseCode == 200) {
try (InputStream inputStream = connection.getInputStream();
ByteArrayOutputStream outputStream = new ByteArrayOutputStream()) {
byte[] buffer = new byte[8192];
int bytesRead;
while ((bytesRead = inputStream.read(buffer)) != -1) {
outputStream.write(buffer, 0, bytesRead);
}
byte[] videoData = outputStream.toByteArray();
log.info("HttpURLConnection视频下载成功,大小: {} bytes", videoData.length);
return videoData;
}
} else {
log.error("HttpURLConnection下载失败,HTTP状态码: {}", responseCode);
// 读取错误信息
try (InputStream errorStream = connection.getErrorStream()) {
if (errorStream != null) {
byte[] errorBytes = new byte[1024];
int errorBytesRead = errorStream.read(errorBytes);
if (errorBytesRead > 0) {
String errorMessage = new String(errorBytes, 0, errorBytesRead);
log.error("错误响应内容: {}", errorMessage);
}
}
}
return null;
}
} catch (Exception e) {
log.error("HttpURLConnection下载异常: {}", e.getMessage(), e);
return null;
}
}
/**
* 使用RestTemplate下载视频(备用方案)
*/
private byte[] downloadWithRestTemplate(String videoUrl) {
try {
// 创建自定义的RestTemplate,禁用URL编码
RestTemplate customRestTemplate = new RestTemplate();
// 设置请求头
HttpHeaders headers = new HttpHeaders();
headers.set("User-Agent", "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36");
headers.set("Accept", "*/*");
headers.set("Accept-Encoding", "identity");
HttpEntity<String> entity = new HttpEntity<>(headers);
ResponseEntity<byte[]> response = customRestTemplate.exchange(
videoUrl,
HttpMethod.GET,
entity,
byte[].class
);
if (response.getStatusCode().is2xxSuccessful() && response.getBody() != null) {
log.info("RestTemplate视频下载成功,大小: {} bytes", response.getBody().length);
return response.getBody();
} else {
log.error("RestTemplate下载失败,状态码: {}", response.getStatusCode());
return null;
}
} catch (Exception e) {
log.error("RestTemplate下载异常: {}", e.getMessage(), e);
return null;
}
}
}