feat(aihr): support video upload with asr transcript and keyframe ocr

新增 AihrVideoService:ffmpeg 抽音轨按5分钟分段调 asr 转写([mm:ss]
时间戳),30s间隔(≤40帧)抽关键帧调 vision 提取画面文字(重复画面去重),
合并文本复用既有归类/切片/向量化链路。视频只走异步队列(同步接口
拒绝),≤500MB/≤60分钟,同时只加工1个,卡死判定放宽到120分钟;无
asr/vision 结果按「待处理」落库不报错。前端 accept/筛选/文案与
overview 类型徽标同步。端到端验证:30s培训视频语音逐字转写+画面
四要点完整提取,自动归类投诉处理SOP。
This commit is contained in:
2026-07-04 00:24:24 +08:00
parent 90e96322bb
commit fe1eef3905
6 changed files with 321 additions and 31 deletions
@@ -93,13 +93,17 @@ public class AihrSopSeedService {
private final ISysOssService ossService;
private final String importRootConfig;
private final ScheduledExecutorService scheduledExecutorService;
private final AihrVideoService videoService;
public AihrSopSeedService(ObjectMapper objectMapper, JdbcTemplate jdbcTemplate, ISysOssService ossService, @Value("${aihr.import.root:}") String importRootConfig, ScheduledExecutorService scheduledExecutorService) {
public AihrSopSeedService(ObjectMapper objectMapper, JdbcTemplate jdbcTemplate, ISysOssService ossService,
@Value("${aihr.import.root:}") String importRootConfig,
ScheduledExecutorService scheduledExecutorService, AihrVideoService videoService) {
this.objectMapper = objectMapper;
this.jdbcTemplate = jdbcTemplate;
this.ossService = ossService;
this.importRootConfig = importRootConfig;
this.scheduledExecutorService = scheduledExecutorService;
this.videoService = videoService;
}
public SearchResponse search(SearchRequest request) {
@@ -468,8 +472,8 @@ public class AihrSopSeedService {
try {
content = contentReader.get();
if (content.isBlank()) {
if (imageFile(fileName)) {
// 作战清单 M7:无视觉模型或 OCR 无结果的图片先落库标「待处理」,不算失败
if (imageFile(fileName) || AihrVideoService.videoFile(fileName)) {
// 作战清单 M7:无视觉模型或 OCR/转写无结果的图片、视频先落库标「待处理」,不算失败
return savePendingImage(fileName, category, ossUploader);
}
throw new ServiceException("文件内容不能为空");
@@ -529,19 +533,27 @@ public class AihrSopSeedService {
}
/**
* OCR 拿不到文字的图片:只落 OSS + attach 标「等待解析」(status 0),0 片段返回;
* 配置视觉模型后重新上传同名文件即可替换加工。
* OCR/转写拿不到文字的图片或视频:只落 OSS + attach 标「等待解析」(status 0),0 片段返回;
* 配置视觉/ASR 模型后重新上传同名文件即可替换加工。
*/
private UploadResponse savePendingImage(String fileName, String category, Supplier<SysOssVo> ossUploader) {
String knowledgeName = isAutoCategory(category) || category == null || category.isBlank() ? "图片素材" : cleanCategory(category);
boolean video = AihrVideoService.videoFile(fileName);
String defaultName = video ? "培训视频" : "图片素材";
String knowledgeName = isAutoCategory(category) || category == null || category.isBlank() ? defaultName : cleanCategory(category);
KnowledgeConfig config = knowledgeConfig(knowledgeName);
SysOssVo oss = ossUploader.get();
String docId = upsertAttach(config.knowledgeId(), oss.getOssId(), fileName);
markAttachStatus(config.knowledgeId(), docId, 0, "pending-ocr");
String hint = visionRuntime().isPresent()
? "图片 OCR 未提取到文字(当前模型可能不支持图片输入,或图片中无可识别文字),已入库标记待处理"
: "未配置可用视觉模型,图片已入库标记待处理;在模型管理启用 vision/多模态 chat 模型后重新上传即可解析";
return new UploadResponse(docId, oss.getOssId(), fileName, knowledgeName, 0, hint, List.of("图片", "待OCR"), List.of());
markAttachStatus(config.knowledgeId(), docId, 0, video ? "pending-transcribe" : "pending-ocr");
String hint;
if (video) {
hint = "视频未提取到语音转写或画面文字,已入库标记待处理;确认已启用 asr(语音)或 vision(画面)模型后重新上传即可解析";
} else {
hint = visionRuntime().isPresent()
? "图片 OCR 未提取到文字(当前模型可能不支持图片输入,或图片中无可识别文字),已入库标记待处理"
: "未配置可用视觉模型,图片已入库标记待处理;在模型管理启用 vision/多模态 chat 模型后重新上传即可解析";
}
List<String> tags = video ? List.of("视频", "待转写") : List.of("图片", "待OCR");
return new UploadResponse(docId, oss.getOssId(), fileName, knowledgeName, 0, hint, tags, List.of());
}
private DocumentInsight documentInsight(String fileName, String content, String requestedCategory) {
@@ -1274,6 +1286,8 @@ public class AihrSopSeedService {
case "ppt", "pptx" -> "PPT";
case "pdf" -> "PDF";
case "txt", "md", "markdown" -> "Text";
case "jpg", "jpeg", "png", "gif", "webp", "bmp" -> "图片";
case "mp4", "mov", "avi", "mkv", "webm", "m4v" -> "视频";
default -> type.toUpperCase();
};
}
@@ -1978,6 +1992,10 @@ public class AihrSopSeedService {
}
private String readContent(MultipartFile file, String fileName) {
if (AihrVideoService.videoFile(fileName)) {
// 视频加工耗时数分钟,只允许走异步队列(暂存路径版本),同步接口直接拒绝
throw new ServiceException("视频请使用资料处理中心的批量导入(异步队列)上传");
}
if (textFile(fileName)) {
return normalizeMarkdown(readUtf8(file));
}
@@ -2007,6 +2025,20 @@ public class AihrSopSeedService {
}
private String readContent(Path file, String fileName) {
if (AihrVideoService.videoFile(fileName)) {
return videoService.extractText(file, fileName, (bytes, mimeType) -> {
Optional<ChatRuntime> runtime = visionRuntime();
if (runtime.isEmpty()) {
return Optional.empty();
}
try {
return Optional.of(normalizeExtractedText(callVisionOcr(runtime.get(), bytes, mimeType)));
} catch (Exception e) {
log.warn("video frame ocr failed for {}: {}", fileName, e.getMessage());
return Optional.empty();
}
});
}
if (textFile(fileName)) {
return normalizeMarkdown(readUtf8(file));
}
@@ -2180,6 +2212,7 @@ public class AihrSopSeedService {
String lower = fileName.toLowerCase();
return textFile(lower)
|| imageFile(lower)
|| AihrVideoService.videoFile(lower)
|| lower.endsWith(".pdf")
|| lower.endsWith(".doc")
|| lower.endsWith(".docx")
@@ -37,11 +37,15 @@ public class AihrUploadQueueService {
private static final String TENANT_ID = "000000";
private static final DateTimeFormatter TIME_FORMATTER = DateTimeFormatter.ofPattern("MM-dd HH:mm:ss");
private static final long MAX_FILE_BYTES = 100L * 1024 * 1024;
/** 视频单独放宽:培训视频普遍大于文档;更大的走服务端目录导入。 */
private static final long MAX_VIDEO_BYTES = 500L * 1024 * 1024;
private static final int MAX_ERROR_CHARS = 480;
/** 加工并发上限:链路含 LLM 归类与 embedding 外呼,并发过高会互相争抢配额并拖慢单文件耗时。 */
private static final int WORKER_PERMITS = 2;
/** 处理中超过该时长视为孤儿(服务重启/线程异常丢失),由清扫任务重置回待处理。需大于最慢单文件加工时长。 */
private static final int STUCK_MINUTES = 30;
/** 视频加工链路(ffmpeg+分段ASR+关键帧OCR)耗时更长,卡死判定单独放宽。 */
private static final int VIDEO_STUCK_MINUTES = 120;
/** 失败条目的暂存文件保留时长,超时由清扫回收磁盘(之后重试会提示重新上传)。 */
private static final int STAGING_RETENTION_HOURS = 72;
@@ -73,10 +77,13 @@ public class AihrUploadQueueService {
}
String fileName = sanitizeFileName(file.getOriginalFilename());
if (!AihrSopSeedService.supportedFile(fileName)) {
throw new ServiceException("仅支持 txt/md/markdown/pdf/doc/docx/xls/xlsx/ppt/pptx/图片文件");
throw new ServiceException("仅支持 txt/md/markdown/pdf/doc/docx/xls/xlsx/ppt/pptx/图片/视频文件");
}
if (file.getSize() > MAX_FILE_BYTES) {
throw new ServiceException("文件不能超过 100MB");
long maxBytes = AihrVideoService.videoFile(fileName) ? MAX_VIDEO_BYTES : MAX_FILE_BYTES;
if (file.getSize() > maxBytes) {
throw new ServiceException(AihrVideoService.videoFile(fileName)
? "视频不能超过 500MB,更大的请放到导入根目录走服务端导入"
: "文件不能超过 100MB");
}
ensureTable();
String batch = sanitizeBatchId(batchId);
@@ -177,25 +184,30 @@ public class AihrUploadQueueService {
private Long claimNext() {
try {
// 同名文件不并发领取:saveDocument 对同 knowledge+fileName 是替换语义,并发会撞 attach 唯一键
List<Long> candidates = jdbcTemplate.query("""
select i.id from aihr_knowledge_upload_item i
// 同名文件不并发领取:saveDocument 对同 knowledge+fileName 是替换语义,并发会撞 attach 唯一键。
// 视频最多同时加工 1 个:单个视频要占数分钟的 ffmpeg+ASR+vision,双路并发互相拖慢。
boolean videoBusy = hasProcessingVideo();
List<QueueCandidate> candidates = jdbcTemplate.query("""
select i.id, i.file_name from aihr_knowledge_upload_item i
where i.tenant_id = ? and i.status = 0
and not exists (
select 1 from (select file_name from aihr_knowledge_upload_item where tenant_id = ? and status = 1) p
where p.file_name = i.file_name
)
order by i.id asc
limit 5
""", (rs, rowNum) -> rs.getLong(1), TENANT_ID, TENANT_ID);
for (Long id : candidates) {
limit 10
""", (rs, rowNum) -> new QueueCandidate(rs.getLong(1), rs.getString(2)), TENANT_ID, TENANT_ID);
for (QueueCandidate candidate : candidates) {
if (videoBusy && AihrVideoService.videoFile(candidate.fileName())) {
continue;
}
int updated = jdbcTemplate.update("""
update aihr_knowledge_upload_item
set status = 1, update_time = now()
where tenant_id = ? and id = ? and status = 0
""", TENANT_ID, id);
""", TENANT_ID, candidate.id());
if (updated == 1) {
return id;
return candidate.id();
}
}
return null;
@@ -205,6 +217,16 @@ public class AihrUploadQueueService {
}
}
private boolean hasProcessingVideo() {
List<String> processing = jdbcTemplate.query("""
select file_name from aihr_knowledge_upload_item where tenant_id = ? and status = 1
""", (rs, rowNum) -> rs.getString(1), TENANT_ID);
return processing.stream().anyMatch(AihrVideoService::videoFile);
}
private record QueueCandidate(Long id, String fileName) {
}
private void processItem(Long id) {
ItemRow row = jdbcTemplate.query("""
select id, file_name, category, staging_path, status
@@ -249,8 +271,12 @@ public class AihrUploadQueueService {
int reset = jdbcTemplate.update("""
update aihr_knowledge_upload_item
set status = 0, update_time = now()
where tenant_id = ? and status = 1 and update_time < (now() - interval %d minute)
""".formatted(STUCK_MINUTES), TENANT_ID);
where tenant_id = ? and status = 1
and (
(update_time < (now() - interval %d minute) and file_name not regexp '\\\\.(mp4|mov|avi|mkv|webm|m4v)$')
or update_time < (now() - interval %d minute)
)
""".formatted(STUCK_MINUTES, VIDEO_STUCK_MINUTES), TENANT_ID);
if (reset > 0) {
log.info("upload queue reset {} stuck items", reset);
}
@@ -0,0 +1,223 @@
package org.dromara.aihr.service;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.dromara.common.core.exception.ServiceException;
import org.springframework.stereotype.Service;
import java.io.IOException;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.Comparator;
import java.util.List;
import java.util.Locale;
import java.util.Optional;
import java.util.concurrent.TimeUnit;
import java.util.stream.Stream;
/**
* 视频加工:ffmpeg 抽音轨分段 + ASR 转写(带时间戳),抽关键帧 + 视觉 OCR(PPT 录屏字幕类画面文字)。
* 产物是纯文本,供上层复用既有的归类/切片/向量化链路。视觉 OCR 由调用方注入(避免与 SOP 服务循环依赖)。
*/
@Service
@RequiredArgsConstructor
@Slf4j
public class AihrVideoService {
/** 单段音频时长(秒):SenseVoice 对长音频不友好,分段转写后拼接。 */
private static final int AUDIO_SEGMENT_SECONDS = 300;
/** 关键帧抽取间隔(秒)与上限:控制 vision 调用成本(G6)。 */
private static final int FRAME_INTERVAL_SECONDS = 30;
private static final int MAX_FRAMES = 40;
private static final int MAX_AUDIO_SEGMENTS = 12;
/** 支持的最长视频(秒),超长直接拒绝,引导切分后再上传。 */
private static final int MAX_DURATION_SECONDS = 3600;
private static final int FFMPEG_TIMEOUT_SECONDS = 300;
private final AihrSpeechService speechService;
/** 视觉 OCR 回调:入参 (imageBytes, mimeType),返回识别文本;不可用/失败返回 empty。 */
@FunctionalInterface
public interface FrameOcr {
Optional<String> ocr(byte[] imageBytes, String mimeType);
}
public static boolean videoFile(String fileName) {
String lower = fileName == null ? "" : fileName.toLowerCase(Locale.ROOT);
return lower.endsWith(".mp4") || lower.endsWith(".mov") || lower.endsWith(".avi")
|| lower.endsWith(".mkv") || lower.endsWith(".webm") || lower.endsWith(".m4v");
}
public boolean ffmpegAvailable() {
return commandAvailable("ffmpeg") && commandAvailable("ffprobe");
}
/**
* 提取视频文本:音轨转写([mm:ss] 前缀)+ 关键帧画面文字([画面 mm:ss] 前缀)。
* ASR 未配置且视觉不可用时返回空串,由上层按「待处理」落库(同图片语义)。
*/
public String extractText(Path video, String fileName, FrameOcr frameOcr) {
if (!ffmpegAvailable()) {
throw new ServiceException("服务器未安装 ffmpeg/ffprobe,无法解析视频");
}
long durationSeconds = probeDurationSeconds(video);
if (durationSeconds > MAX_DURATION_SECONDS) {
throw new ServiceException("视频超过 60 分钟,请切分后再上传");
}
Path workDir = null;
try {
workDir = Files.createTempDirectory("aihr-video-");
StringBuilder text = new StringBuilder();
String transcript = transcribeAudio(video, workDir, durationSeconds);
if (!transcript.isBlank()) {
text.append("【语音转写】\n").append(transcript).append('\n');
}
String frameText = ocrKeyFrames(video, workDir, durationSeconds, frameOcr);
if (!frameText.isBlank()) {
text.append("【画面文字】\n").append(frameText).append('\n');
}
return text.toString().strip();
} catch (ServiceException e) {
throw e;
} catch (Exception e) {
log.warn("video extract failed for {}: {}", fileName, e.getMessage());
throw new ServiceException("视频解析失败:" + e.getMessage());
} finally {
cleanDirQuietly(workDir);
}
}
private String transcribeAudio(Path video, Path workDir, long durationSeconds) throws IOException, InterruptedException {
if (!speechService.asrConfigured()) {
log.info("asr not configured, skip video transcription");
return "";
}
int segments = (int) Math.min(MAX_AUDIO_SEGMENTS, (durationSeconds + AUDIO_SEGMENT_SECONDS - 1) / AUDIO_SEGMENT_SECONDS);
StringBuilder transcript = new StringBuilder();
for (int i = 0; i < segments; i++) {
long offset = (long) i * AUDIO_SEGMENT_SECONDS;
Path segment = workDir.resolve("audio-" + i + ".mp3");
// 16kHz 单声道 mp3:ASR 友好且体积小
int exit = runFfmpeg(List.of(
"ffmpeg", "-y", "-nostdin", "-v", "error",
"-ss", String.valueOf(offset), "-t", String.valueOf(AUDIO_SEGMENT_SECONDS),
"-i", video.toString(),
"-vn", "-ac", "1", "-ar", "16000", "-b:a", "48k",
segment.toString()
));
if (exit != 0 || !Files.isRegularFile(segment) || Files.size(segment) < 1024) {
// 无音轨或该段为空(比如纯画面视频),跳过不算失败
continue;
}
Optional<String> piece = speechService.transcribe(Files.readAllBytes(segment), "segment-" + i + ".mp3", "audio/mpeg");
piece.filter(s -> !s.isBlank())
.ifPresent(s -> transcript.append('[').append(formatTime(offset)).append("] ").append(s.strip()).append('\n'));
Files.deleteIfExists(segment);
}
return transcript.toString();
}
private String ocrKeyFrames(Path video, Path workDir, long durationSeconds, FrameOcr frameOcr) throws IOException, InterruptedException {
if (frameOcr == null) {
return "";
}
long step = Math.max(FRAME_INTERVAL_SECONDS, durationSeconds / MAX_FRAMES + 1);
StringBuilder frameText = new StringBuilder();
String lastText = "";
for (long offset = 0; offset < Math.max(durationSeconds, 1); offset += step) {
Path frame = workDir.resolve("frame-" + offset + ".jpg");
int exit = runFfmpeg(List.of(
"ffmpeg", "-y", "-nostdin", "-v", "error",
"-ss", String.valueOf(offset), "-i", video.toString(),
"-frames:v", "1", "-q:v", "4", "-vf", "scale='min(1280,iw)':-2",
frame.toString()
));
if (exit != 0 || !Files.isRegularFile(frame)) {
continue;
}
String text = frameOcr.ocr(Files.readAllBytes(frame), "image/jpeg").orElse("").strip();
Files.deleteIfExists(frame);
// 相邻帧文字几乎相同(静止幻灯片)时去重,避免片段灌水
if (!text.isBlank() && !similarText(text, lastText)) {
frameText.append("[画面 ").append(formatTime(offset)).append("] ").append(text).append('\n');
lastText = text;
}
}
return frameText.toString();
}
private long probeDurationSeconds(Path video) {
try {
Process process = new ProcessBuilder(
"ffprobe", "-v", "error", "-show_entries", "format=duration",
"-of", "default=noprint_wrappers=1:nokey=1", video.toString()
).redirectErrorStream(true).start();
String output = new String(process.getInputStream().readAllBytes(), StandardCharsets.UTF_8).strip();
if (!process.waitFor(30, TimeUnit.SECONDS)) {
process.destroyForcibly();
throw new ServiceException("视频时长探测超时");
}
return (long) Double.parseDouble(output.lines().findFirst().orElse("0"));
} catch (ServiceException e) {
throw e;
} catch (Exception e) {
throw new ServiceException("视频文件无法识别:" + e.getMessage());
}
}
private int runFfmpeg(List<String> command) throws IOException, InterruptedException {
Process process = new ProcessBuilder(command).redirectErrorStream(true).start();
// 消费输出防止缓冲区塞满导致进程挂起
byte[] discarded = process.getInputStream().readAllBytes();
if (!process.waitFor(FFMPEG_TIMEOUT_SECONDS, TimeUnit.SECONDS)) {
process.destroyForcibly();
throw new ServiceException("ffmpeg 处理超时");
}
int exit = process.exitValue();
if (exit != 0) {
log.debug("ffmpeg exit {}: {}", exit, new String(discarded, StandardCharsets.UTF_8));
}
return exit;
}
private static boolean commandAvailable(String command) {
try {
Process process = new ProcessBuilder(command, "-version").redirectErrorStream(true).start();
process.getInputStream().readAllBytes();
return process.waitFor(10, TimeUnit.SECONDS) && process.exitValue() == 0;
} catch (Exception e) {
return false;
}
}
/** 简易相似判定:前 80 字符相同视为同一画面(幻灯片停留场景)。 */
private static boolean similarText(String a, String b) {
if (a.isBlank() || b.isBlank()) {
return false;
}
int len = Math.min(80, Math.min(a.length(), b.length()));
return a.regionMatches(0, b, 0, len);
}
private static String formatTime(long seconds) {
return "%02d:%02d".formatted(seconds / 60, seconds % 60);
}
private static void cleanDirQuietly(Path dir) {
if (dir == null) {
return;
}
try (Stream<Path> stream = Files.walk(dir)) {
stream.sorted(Comparator.reverseOrder()).forEach(path -> {
try {
Files.deleteIfExists(path);
} catch (IOException ignored) {
// 临时目录清理失败不影响业务
}
});
} catch (IOException ignored) {
// 同上
}
}
}