Files
vrsub/src/wov_app/batch.py
T
cat-shark 171b088e7c feat: 批量流水线按 GPU 资源调度(非 GPU 阶段并行、GPU 阶段互斥)
此前组内阶段是串行的(先全部提音、再全部转写、再全部翻译),LLM 走线上端点时
翻译阶段不占显存、GPU 全程空转——实测占整轮挂钟约 40%(19.6W / 272MiB)。

- `_run_job` 改为按组启动在途流水线:每个视频独立推进自己的阶段,最多
  `WOV_BATCH_PIPELINE_WORKERS`(默认 4)个阶段在途。
- 派发只看资源:`stage_gpu_need_mb` 为 0 的阶段(提音、线上翻译、ASS)立刻派发,
  可与其它视频的转写并行;需要 GPU 的阶段由 `GpuGate` 互斥准入,并按"阶段索引
  最小者优先"派发,组内仍是先跑完全部转写再进翻译——本机 Ollama 模型每组只
  加载一次,不需要按"是否云端"写分支。
- 同一阶段只在途一份(派发即标记 running),单视频异常不带走整组;暂停沿用
  run 级 paused.flag,暂停后不再派发新阶段。
- 测试:远端翻译与其它视频转写重叠、本机端点下全部转写先于翻译且翻译互斥、
  提音与转写重叠,以及既有分组/暂停/失败隔离用例。
2026-09-18 22:40:53 +08:00

883 lines
41 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""文件夹批量处理引擎。
本地版核心能力:**不把视频上传到工作目录**,而是直接读取用户所选文件夹下的
全部视频,逐个调用现有的工作流流水线(复用 WorkflowScheduler 的 DAG 执行与
断点续跑逻辑)。
处理约定:
- **创建任务时一次性定位**:`create_job` 扫描文件夹并把每个视频登记为
batch_videos 明细;视频所在目录(视频旁)若已存在**文件名包含视频名**的
字幕文件(`.srt/.ass/.ssa/.vtt`),说明该视频已有字幕,直接记为 SKIPPED,
不为它触发任何流水线。运行时(BatchWorker)只消费已定位好的明细列表,
**不再重新扫描文件夹**(运行期间新增/删除的视频不会改变本次任务的范围)。
- **按资源调度的在途流水线**:视频按 `WOV_BATCH_STAGE_GROUP_SIZE` 分组,组内
每个视频独立推进自己的阶段(最多 `WOV_BATCH_PIPELINE_WORKERS` 个阶段在途)。
阶段是否需要 GPU 由 `wov_app.resources` 判定:提音、线上翻译与 ASS 不需要,
与其它视频的转写并行(GPU 不再空转);需要 GPU 的阶段由进程内门控
`GpuGate` 串行准入,并按"阶段索引最小者优先"派发,于是组内先跑完全部转写
再进翻译——本机 Ollama 模型每组只加载一次(线上端点则完全不受该顺序约束)。
阶段边界用 `execute_run(stop_after=节点)` 停在节点(任务保持 RUNNING),
LLM 阶段靠 `keep_model.flag` 让节点保持模型常驻,组末由引擎统一释放显存
(详见 docs/operations.md#文件夹批量处理)。
- **产物放在视频旁**:每个视频处理完成后,把工作流 `final_outputs` 对应的
最终产物文件(字幕流水线即日语转写 `.srt`、中文 `.srt` 与双目 `.ass`
**复制一份到视频的所在目录**,与 .mp4 放在一起;文件名按 `final_outputs`
别名**对齐媒体库既有约定**:日语转写存为 `<视频名>.JA.srt`、中文字幕存为
`<视频名>.CN.srt`、双目字幕存为 `<视频名>.CN_dual_eye.ass`(文件名稳定且
含视频主名,媒体库可自动匹配,下次批量扫描也会命中"已有字幕"规则跳过)。
- **过程文件清理**:视频收尾完成后删除该视频的整个工作空间与 run 记录,
中间产物(音频/分块/帧图/节点产物)不残留在媒体库,也不会被影视库软件
当作视频载入。工作空间位于应用私有目录 `storage/batch/<job_id>/<bv_id>/`
与用户的视频库目录天然隔离;暂停/失败的视频保留工作空间以便断点续跑。
暂停/恢复语义(对应前端"暂停/继续"按钮):
- 暂停批量任务:把批量任务置为 PAUSED,并暂停当前正在执行的 run(写
paused.flagwhisper 按分块、OCR 按帧检查后停止),处理中的视频保持
PAUSED,后续视频不再开始。
- 继续:批量任务恢复 QUEUED,引擎从断点继续——PAUSED 视频的 run 显式
resume 后由 execute_run 从产物表断点续跑,已完成节点不重复执行。
主调度器不会抢占批量 runnext_queued_run 排除 source=batch),批量引擎
使用 per-video 的 WorkflowScheduler 实例,storage 指向该视频的私有工作空间。
"""
from __future__ import annotations
import json
import shutil
import threading
import time
import uuid
from concurrent.futures import FIRST_COMPLETED, Future, ThreadPoolExecutor, wait
from datetime import datetime, timezone
from pathlib import Path
from wov_app import registry
from wov_app.config import (
BATCH_INTERVAL_SECONDS,
BATCH_PIPELINE_WORKERS,
BATCH_STAGE_GROUP_SIZE,
STORAGE_DIR,
)
from wov_app.db import Database
from wov_app.logging import get_logger
from wov_app.resources import GpuGate, stage_gpu_need_mb
from wov_app.scheduler import WorkflowScheduler, topological_sort
from wov_app.storage import atomic_copy
from wov_sdk.models import WorkflowDefinition, WorkflowNode
from nodes.llm import release_local_model
# 批量引擎运行日志:任务进度、视频逐个处理与暂停/续跑等状态变化。
logger = get_logger("batch")
# 识别为视频文件的扩展名(大小写不敏感,扫描时统一转小写比较)。
VIDEO_EXTENSIONS = {
".mp4", ".mkv", ".avi", ".mov", ".webm", ".flv", ".ts",
".m4v", ".wmv", ".mpg", ".mpeg", ".3gp",
}
# 字幕文件扩展名:批量扫描时按它识别"视频旁已有字幕";处理后放回视频旁的
# 最终产物(.srt/.ass)也在该集合内,保证下次扫描能命中同一规则直接跳过。
SUBTITLE_EXTENSIONS = {".srt", ".ass", ".ssa", ".vtt"}
# 组内流水线的轮询间隔(秒):等第一个阶段结束时顺带检查暂停与资源放行。
PIPELINE_POLL_SECONDS = 0.2
def _make_gate() -> GpuGate:
"""创建任务级 GPU 门控(测试通过替换本函数注入假探测结果)。"""
return GpuGate()
# 暂停信号文件名:与节点约定一致,位于 run 根目录(<work_dir>/runs/<run_id>/)。
PAUSE_FLAG = "paused.flag"
# 保持模型常驻信号文件名:LLM 阶段执行期间由引擎写入 run 根目录,节点据此不在
# 每次调用后卸载模型(同组视频共用一份已加载模型,减少加载/卸载次数)。
KEEP_MODEL_FLAG = "keep_model.flag"
# LLM 节点类型前缀:这类节点加载本地大模型,阶段结束后由引擎统一释放显存。
LLM_NODE_PREFIX = "llm"
# 兼容读取的历史完成标记文件名(旧任务用它记录产物路径)。当前逻辑不再
# 写入,产物直接放视频旁;保留读取能力以便旧任务的详情/下载仍可用。
MARKER_NAME = "batch.done.json"
# 批量处理私有工作空间根目录:位于应用存储目录下(data/storage/batch)。
# 每个视频的工作目录为 <根>/<job_id>/<bv_id>/,与用户视频库目录完全隔离,
# 媒体库软件只会看到最终放到视频旁的 .srt/.ass 字幕成品。
BATCH_WORK_ROOT = STORAGE_DIR / "batch"
def _now_iso() -> str:
"""返回当前 UTC 时间的 ISO 格式字符串。"""
return datetime.now(timezone.utc).isoformat()
def scan_videos(folder: Path, recursive: bool = True) -> list[Path]:
"""扫描文件夹下的全部视频文件,按路径排序保证处理顺序确定。
recursive=True 时递归扫描子文件夹;recursive=False 只扫描顶层。
"""
if recursive:
paths = [
p for p in folder.rglob("*")
if p.is_file() and p.suffix.lower() in VIDEO_EXTENSIONS
]
else:
paths = [
p for p in folder.glob("*")
if p.is_file() and p.suffix.lower() in VIDEO_EXTENSIONS
]
return sorted(paths)
def list_sidecar_subtitles(video: Path) -> list[Path]:
"""列出视频所在目录(视频旁)与视频"对应"的字幕文件。
判定规则:与视频同一目录、扩展名为字幕格式、且文件名包含视频主名
(大小写不敏感)的文件都视为该视频已带的字幕。典型命中如
`movie.srt`、`movie.CN.srt`、`movie.CN_dual_eye.ass`,以及本引擎处理
完成后放到视频旁的 `movie.CN.srt` / `movie.CN_dual_eye.ass`。
视频主名过短(单个字符)时只接受"主名."前缀,避免 a.mp4 误配 apple.srt。
目录不可读时保守返回空列表,不影响批量任务创建。
"""
stem = video.stem.lower()
try:
siblings = list(video.parent.iterdir())
except OSError:
return []
found: list[Path] = []
for item in siblings:
try:
if not item.is_file() or item.suffix.lower() not in SUBTITLE_EXTENSIONS:
continue
except OSError:
# 单个子项不可读(权限不足)时跳过,不拖垮整个目录。
continue
item_name = item.name.lower()
# 视频本身不是字幕扩展名,此处无需再排除同名文件;直接按主名匹配。
if len(stem) <= 1:
matched = item_name.startswith(stem + ".")
else:
matched = stem in item_name
if matched:
found.append(item)
return sorted(found)
# final_outputs 别名 → 视频旁文件名后缀:日语转写与中文译文同为 `.srt`,
# 只按扩展名映射会让两份产物撞名互相覆盖,必须按别名区分语言。
_PRODUCT_SUFFIXES = {
"ja_srt": ".JA.srt",
"cn_srt": ".CN.srt",
"ass": ".CN_dual_eye.ass",
}
def _sidecar_product_name(video: Path, source: Path, alias: str | None = None) -> str:
"""把最终产物映射为放在视频旁时的标准字幕文件名。
对齐媒体库既有约定(文件名稳定、无时间戳,媒体库可按视频主名自动匹配):
- 日语转写(`ja_srt`)→ `<视频名>.JA.srt`
- 中文 `.srt` 产物 → `<视频名>.CN.srt`
- 双目 `.ass` 产物 → `<视频名>.CN_dual_eye.ass`
- 其余扩展名的最终产物保留原文件名(含时间戳),避免误改语义。
别名来自 `final_outputs` 的键,未登记别名时按扩展名兜底(历史工作流兼容)。
"""
if alias in _PRODUCT_SUFFIXES:
return f"{video.stem}{_PRODUCT_SUFFIXES[alias]}"
suffix = source.suffix.lower()
if suffix == ".srt":
return f"{video.stem}.CN.srt"
if suffix == ".ass":
return f"{video.stem}.CN_dual_eye.ass"
return source.name
def load_marker(work_dir: Path) -> dict | None:
"""读取历史完成标记;不存在或损坏时返回 None。"""
path = work_dir / MARKER_NAME
if not path.is_file():
return None
try:
data = json.loads(path.read_text(encoding="utf-8"))
return data if isinstance(data, dict) else None
except (json.JSONDecodeError, OSError):
# 半行写入或权限异常时保守视为无标记,走旁挂字幕判定。
return None
def remove_job_workspace(job_id: str) -> None:
"""删除批量任务在应用私有存储下的工作空间目录(<storage>/batch/<job_id>)。
每个视频完成后工作空间已被逐视频清理;此处兜底清理任务级残留(删除任务
或任务异常中止时)。只作用于私有工作空间,绝不触碰用户视频目录。
"""
shutil.rmtree(BATCH_WORK_ROOT / job_id, ignore_errors=True)
def create_job(
db: Database,
folder_path: str,
workflow_id: str,
recursive: bool = True,
) -> str:
"""创建批量任务:校验文件夹与工作流、**一次性定位**视频并登记明细。
扫描到的每个视频都会登记为 batch_videos 明细:视频旁已有对应字幕文件
的直接记 SKIPPED(不触发流水线),否则记 PENDING(等待引擎处理)。
引擎运行时只消费这批已定位的明细,不再重新扫描文件夹。
校验失败抛出 ValueError(由路由层转为 422 响应)。
"""
folder = Path(folder_path).expanduser()
if not folder.is_dir():
raise ValueError("folder not found")
workflow = db.get_workflow(workflow_id)
if workflow is None or not workflow["published"]:
raise ValueError("published workflow not found")
if db.get_latest_workflow_version(workflow_id) is None:
raise ValueError("workflow has no version")
videos = scan_videos(folder, recursive)
if not videos:
raise ValueError("no videos found in folder")
job_id = f"batch_{uuid.uuid4().hex[:12]}"
now = _now_iso()
db.create_batch_job({
"id": job_id,
"folder_path": str(folder),
"workflow_id": workflow_id,
"recursive": int(recursive),
# CREATING:明细未写完前引擎看不见本任务。逐条登记 500+ 条明细要数秒,
# 若此刻已是 QUEUED,引擎会读到半个快照并在收尾时把任务误标 COMPLETED,
# 剩下的视频就再也不会被处理(自愈分支只碰运气)。写完明细立刻置 QUEUED。
"status": "CREATING",
"progress": 0,
"total": 0,
"done": 0,
"failed": 0,
"current_video": None,
"error": None,
"created_at": now,
"updated_at": now,
})
pending = 0
skipped = 0
try:
for video in videos:
# 视频所在目录已存在对应字幕文件 → 已处理过,直接跳过不触发流水线。
if list_sidecar_subtitles(video):
status = "SKIPPED"
skipped += 1
else:
status = "PENDING"
pending += 1
video_id = f"bv_{uuid.uuid4().hex[:12]}"
db.create_batch_video({
"id": video_id,
"job_id": job_id,
"video_path": str(video),
# 私有工作空间:storage/batch/<job_id>/<bv_id>/,与媒体库隔离。
"work_dir": str(BATCH_WORK_ROOT / job_id / video_id),
"run_id": None,
"status": status,
"error": None,
"created_at": now,
"updated_at": now,
})
except Exception as exc: # noqa: BLE001
# 明细写到一半失败:记 FAILED 留可见记录(CREATING 状态没人会拾起,
# 沉默的残留任务会让用户以为什么都没发生),然后把异常交给路由层。
db.update_batch_job(
job_id, status="FAILED", error=str(exc), updated_at=_now_iso(),
)
raise
# total = 本批真正需要处理(无字幕)的视频数;已有字幕被 SKIPPED 的
# 不计入总数也不计入完成数——进度条只反映"实际待处理"的这批。
if pending == 0:
# 整批都已有字幕、无任何可处理项:直接视为完成,不排队。
db.update_batch_job(
job_id, status="COMPLETED", total=0, progress=1.0,
updated_at=_now_iso(),
)
else:
# 明细全部就位后才排队,引擎从此拿到的快照一定是完整的。
db.update_batch_job(
job_id, status="QUEUED", total=pending, updated_at=_now_iso(),
)
logger.info(
"创建批量任务 %s: 文件夹 %s, 工作流 %s, 共 %d 个视频(%d 待处理, %d 已有字幕跳过)",
job_id, folder, workflow_id, len(videos), pending, skipped,
)
return job_id
class BatchWorker:
"""批量处理引擎:单线程轮询 QUEUED 批量任务,逐视频调用现有调度器执行。"""
def __init__(
self,
db: Database,
interval_seconds: float | None = None,
) -> None:
"""保存数据库依赖并初始化轮询线程控制字段。"""
self.db = db
self.interval_seconds = interval_seconds or BATCH_INTERVAL_SECONDS
self._thread: threading.Thread | None = None
self._stopping = False
def start(self) -> None:
"""启动批量处理线程;重复调用无副作用。"""
if self._thread is not None:
return
# 独立运行时确保节点已注册;重复注册幂等。
registry.register_all()
# 上一进程中断留下的 CREATING 任务(明细登记中途被杀/热重载)没人会推进,
# 启动时统一记为 FAILED,避免用户以为任务还在创建中。
stale = self.db.fail_creating_batch_jobs("创建明细中断(进程中断),请重新创建任务")
if stale:
logger.warning("启动清理 %d 个未登记完的批量任务(CREATING → FAILED", stale)
self._stopping = False
self._thread = threading.Thread(
target=self._loop,
name="wov-batch-worker",
daemon=True,
)
self._thread.start()
def stop(self) -> None:
"""请求停止并等待轮询线程退出。"""
self._stopping = True
if self._thread is not None:
self._thread.join(timeout=5)
self._thread = None
def _loop(self) -> None:
"""轮询循环:有排队中的批量任务就处理,否则休眠一个间隔。"""
while not self._stopping:
try:
job = self.db.next_queued_batch_job()
if job is not None:
self._process_job(job)
else:
time.sleep(self.interval_seconds)
except Exception: # noqa: BLE001
# 单次轮询异常不杀死线程,记录后跳过本轮(与主调度器一致)。
logger.exception("批量引擎轮询异常,跳过本轮")
time.sleep(self.interval_seconds)
# ------------------------------------------------------------------
# 任务执行
# ------------------------------------------------------------------
def _process_job(self, job: dict) -> None:
"""处理一个批量任务:校验、逐个消费已定位的视频并放置产物。
job 以 QUEUED 状态进入,处理期间置 RUNNING;全部视频处理完置
COMPLETED;被暂停时保持 PAUSED;校验失败置 FAILED。
"""
job_id = str(job["id"])
try:
self._run_job(job_id)
except Exception as exc: # noqa: BLE001
# 任务级兜底:任何未捕获异常都记录到任务而不是卡死在 RUNNING。
logger.exception("批量任务 %s 处理异常", job_id)
self.db.update_batch_job(job_id, status="FAILED", error=str(exc), updated_at=_now_iso())
def _run_job(self, job_id: str) -> None:
"""批量任务主流程(分块流水线,异常由 _process_job 统一处理)。
视频按 WOV_BATCH_STAGE_GROUP_SIZE 分组,组内按 DAG 拓扑顺序逐节点跑完
全部视频(先全部 extract、再全部 ASR、再全部 LLM 翻译、最后 ASS)再
进入下一组:本地模型每组只加载一次、卸载一次,产物按组增量落地。
只消费创建任务时已定位好的 batch_videos 明细,**不再扫描文件夹**。
"""
job = self.db.get_batch_job(job_id)
if job is None:
return
folder = Path(job["folder_path"])
if not folder.is_dir():
self.db.update_batch_job(job_id, status="FAILED", error="folder not found", updated_at=_now_iso())
return
workflow = self.db.get_workflow(job["workflow_id"])
if workflow is None or not workflow["published"]:
self.db.update_batch_job(job_id, status="FAILED", error="workflow not found or unpublished", updated_at=_now_iso())
return
version = self.db.get_latest_workflow_version(job["workflow_id"])
if version is None:
self.db.update_batch_job(job_id, status="FAILED", error="workflow has no version", updated_at=_now_iso())
return
definition = WorkflowDefinition.from_dict(version["definition"])
definition.validate()
order = topological_sort(definition)
node_by_id = {node.id: node for node in definition.nodes}
# 任务行先于视频明细写入(create_job 逐条插入),引擎可能在登记完成前就拾起
# 任务:此时没有任何明细,不能按“空任务”收尾,保持 QUEUED 等登记完成。
if not self.db.list_batch_videos(job_id):
logger.info("批量任务 %s 尚无视频明细(仍在登记),保持 QUEUED 稍后重试", job_id)
self.db.update_batch_job(job_id, status="QUEUED", updated_at=_now_iso())
return
# 待处理明细:SKIPPED 创建时已定(不参与 total/done),COMPLETED 无需重跑。
items = [
item for item in self.db.list_batch_videos(job_id)
if item["status"] not in ("COMPLETED", "SKIPPED")
]
# total 创建时已固定为"无字幕需处理的视频数",此处不覆盖;历史任务
# 的旧口径由 sync_batch_job_progress 在读取时修正为不含 SKIPPED。
self.db.update_batch_job(
job_id, status="RUNNING", progress=0,
current_video=None, error=None, updated_at=_now_iso(),
)
# total 为进度条分母;为 0 表示整批跳过(创建即 COMPLETED)。
total = int(job["total"] or 0)
group_size = max(1, int(BATCH_STAGE_GROUP_SIZE))
gate = _make_gate()
for start in range(0, len(items), group_size):
group = items[start:start + group_size]
outcome = self._run_group(
job=job,
group=group,
order=order,
node_by_id=node_by_id,
definition=definition,
version=version,
gate=gate,
)
if outcome == "PAUSED":
# 停下前先把已完成/失败项入账,让暂停中的前端看到真实进度。
self.db.sync_batch_job_progress(job_id)
logger.info("批量任务 %s 已暂停,等待用户继续", job_id)
return
# 先按明细实时对齐汇总(done 不计 SKIPPED),再判断能否收尾。
# 仍有未结束视频时不能标 COMPLETED,否则会出现“还有待处理视频却已完成”
# 的僵尸状态;此时保持 RUNNING,由引擎下一轮续跑。
self.db.sync_batch_job_progress(job_id)
job = self.db.get_batch_job(job_id)
if job is None:
return
leftovers = [
v for v in self.db.list_batch_videos(job_id)
if v["status"] not in ("SKIPPED", "COMPLETED", "FAILED")
]
if leftovers:
# 有未处理完的视频(常见于创建任务时明细还在逐条写入,本轮快照没包含
# 它们):置回 QUEUED 自愈,让引擎下一轮按最新明细重新分组续跑。
# 留在 RUNNING 不会被引擎再拾起(next_queued_batch_job 只取 QUEUED),
# 任务会停在“运行中但没人推进”的状态。
logger.warning(
"批量任务 %s 仍有 %d 个视频未处理完(%s…),置回 QUEUED 待下一轮续跑",
job_id, len(leftovers), Path(leftovers[0]["video_path"]).name,
)
self.db.update_batch_job(job_id, status="QUEUED", updated_at=_now_iso())
return
done = int(job["done"]) if job else 0
failed = int(job["failed"]) if job else 0
self.db.update_batch_job(
job_id, status="COMPLETED", progress=1.0,
current_video=None, error=None, updated_at=_now_iso(),
)
# 无失败视频时每个视频的工作空间已在收尾时删除,任务目录只剩空壳;
# 有失败视频则保留(它们的工作空间供断点重试)。
if failed == 0:
remove_job_workspace(job_id)
logger.info(
"批量任务 %s 完成: 待处理 %d 个视频, 完成 %d, 失败 %d",
job_id, total, done, failed,
)
def _is_job_paused(self, job_id: str) -> bool:
"""批量任务是否已被暂停(或记录已消失):暂停后不再派发新阶段。"""
job = self.db.get_batch_job(job_id)
return job is None or job["status"] == "PAUSED"
def _next_stage_node(
self,
pipeline: dict,
order: list[str],
node_by_id: dict[str, WorkflowNode],
) -> WorkflowNode | None:
"""返回该视频下一个待执行阶段的节点;已跑完或已失败时返回 None。"""
if pipeline["state"] != "ready" or pipeline["stage"] >= len(order):
return None
return node_by_id[order[pipeline["stage"]]]
def _submit_group_stage(
self,
job: dict,
pipeline: dict,
node_spec: WorkflowNode,
order: list[str],
definition: WorkflowDefinition,
version: dict,
running: dict,
pool: ThreadPoolExecutor,
lease_key: str | None = None,
) -> None:
"""把一个阶段交给线程池执行(异常在池内兜底,不让线程池任务抛出去)。"""
item = pipeline["item"]
is_last_stage = node_spec.id == order[-1]
# 标记在途:同一视频同一阶段只允许有一个执行体,否则会被重复派发。
pipeline["state"] = "running"
if lease_key:
pipeline["lease_key"] = lease_key
if node_spec.node_type.startswith(LLM_NODE_PREFIX):
pipeline["executed_llm"] = True
pipeline["llm_params"] = node_spec.params
self.db.update_batch_job(job["id"], current_video=str(item["video_path"]), updated_at=_now_iso())
logger.info(
"批量任务 %s 视频 %s 进入阶段 %s(在途 %d",
job["id"], Path(item["video_path"]).name, node_spec.id, len(running) + 1,
)
def _worker() -> str:
"""线程内的阶段执行:单视频异常不中断整组。"""
try:
outcome = self._run_stage(
job, item, version, definition, node_spec, order, is_last_stage,
)
except Exception as exc: # noqa: BLE001
logger.exception(
"批量任务 %s 视频 %s 阶段 %s 处理异常",
job["id"], item["video_path"], node_spec.id,
)
self.db.update_batch_video(
item["id"], status="FAILED", error=str(exc), updated_at=_now_iso(),
)
return "FAILED"
return outcome or "RUNNING"
running[pool.submit(_worker)] = pipeline
def _dispatch_group_stages(
self,
job: dict,
pipelines: list[dict],
order: list[str],
node_by_id: dict[str, WorkflowNode],
definition: WorkflowDefinition,
version: dict,
gate: GpuGate,
running: dict,
pool: ThreadPoolExecutor,
workers: int,
) -> None:
"""派发就绪阶段:非 GPU 阶段可并行,GPU 阶段互斥且按上游优先。"""
# 1) 不需要 GPU 的阶段(提音 / 线上翻译 / ASS):立刻派发,与其它视频的
# 转写并行,GPU 不再空转。
for pipeline in pipelines:
if len(running) >= workers:
break
node_spec = self._next_stage_node(pipeline, order, node_by_id)
if node_spec is None:
continue
if stage_gpu_need_mb(node_spec.node_type, node_spec.params) > 0:
continue
self._submit_group_stage(job, pipeline, node_spec, order, definition, version, running, pool)
# 2) 需要 GPU 的阶段:一次只跑一个,且选"阶段索引最小"的视频——组内因此
# 先把转写跑完再进翻译,本机 Ollama 模型仍每组只加载一次。
if len(running) >= workers or gate.holder is not None:
return
candidates = [
(pipeline["stage"], pipeline["index"], pipeline)
for pipeline in pipelines
if self._next_stage_node(pipeline, order, node_by_id) is not None
]
gpu_candidates = [
entry for entry in candidates
if stage_gpu_need_mb(
node_by_id[order[entry[2]["stage"]]].node_type,
node_by_id[order[entry[2]["stage"]]].params,
) > 0
]
if not gpu_candidates:
return
_, _, pipeline = min(gpu_candidates, key=lambda entry: (entry[0], entry[1]))
node_spec = node_by_id[order[pipeline["stage"]]]
need = stage_gpu_need_mb(node_spec.node_type, node_spec.params)
lease_key = f"{pipeline['item']['id']}:{node_spec.id}"
# 显存/在途不满足时本轮跳过,等其它阶段释放后再试(不阻塞派发线程)。
if not gate.try_acquire(lease_key, need):
return
self._submit_group_stage(
job, pipeline, node_spec, order, definition, version, running, pool, lease_key,
)
def _finish_group_stage(
self,
job: dict,
pipeline: dict,
future: Future,
order: list[str],
) -> None:
"""收集一个阶段的执行结果并推进该视频的流水线。"""
try:
outcome = future.result()
except Exception: # noqa: BLE001 - 池内已兜底,这里只保证组不被带崩
logger.exception("批量任务 %s 阶段执行线程异常", job["id"])
pipeline["state"] = "failed"
return
self.db.sync_batch_job_progress(job["id"])
if outcome == "PAUSED":
# 节点在边界(分块/批次/帧)停下:任务保持 PAUSED 等用户继续。
pipeline["state"] = "paused"
self.db.update_batch_job(job["id"], status="PAUSED", updated_at=_now_iso())
return
if outcome == "FAILED":
pipeline["state"] = "failed"
return
pipeline["stage"] += 1
pipeline["state"] = "done" if pipeline["stage"] >= len(order) else "ready"
def _run_group(
self,
job: dict,
group: list[dict],
order: list[str],
node_by_id: dict[str, WorkflowNode],
definition: WorkflowDefinition,
version: dict,
gate: GpuGate,
) -> str:
"""组内在途流水线:每个视频独立推进阶段,GPU 阶段按上游优先串行。
阶段是否需要 GPU 由 resources.stage_gpu_need_mb 判定(提音/线上翻译/ASS
不需要),于是它们与其它视频的转写并行;需要 GPU 的阶段由 GpuGate 准入,
并按"阶段索引最小者优先"派发,组内因此先把转写跑完再进翻译,本机 Ollama
模型仍每组只加载一次。
返回 "PAUSED" 表示组内被暂停(调用方停止任务),其余情况返回 "DONE"。
"""
job_id = job["id"]
workers = max(1, int(BATCH_PIPELINE_WORKERS))
pipelines = [
{"item": item, "stage": 0, "index": index, "state": "ready", "executed_llm": False}
for index, item in enumerate(group)
]
with ThreadPoolExecutor(max_workers=workers) as pool:
running: dict[Future, dict] = {}
while True:
paused = self._is_job_paused(job_id)
if not paused:
self._dispatch_group_stages(
job, pipelines, order, node_by_id, definition, version,
gate, running, pool, workers,
)
if not running:
break
done, _ = wait(
list(running), timeout=PIPELINE_POLL_SECONDS,
return_when=FIRST_COMPLETED,
)
for future in done:
pipeline = running.pop(future)
lease_key = pipeline.pop("lease_key", None)
if lease_key:
gate.release(lease_key)
self._finish_group_stage(job, pipeline, future, order)
if paused and not running:
return "PAUSED"
# 组末统一释放本地模型:否则下一组的 whisper 转写会 CUDA OOM。
llm_params = next((p["llm_params"] for p in pipelines if p.get("llm_params")), None)
if llm_params is not None:
self._release_llm_model(llm_params)
return "PAUSED" if self._is_job_paused(job_id) else "DONE"
def _run_stage(
self,
job: dict,
item: dict,
version: dict,
definition: WorkflowDefinition,
node_spec: WorkflowNode,
order: list[str],
is_last_stage: bool,
) -> str | None:
"""执行一个视频在一个阶段节点上的工作,返回执行后的视频状态。
返回 None 表示本阶段无需执行(视频已完成/已跳过)。per-video 的
WorkflowScheduler 以该视频的私有工作空间为 storage,产物表记录各节点
输出,所以同一视频的后续阶段直接从断点继续(不重跑已完成节点)。
LLM 阶段会写 keep_model.flag:阶段内保持模型常驻,阶段结束由引擎统一
释放(见 _release_llm_model),避免每个视频重新加载/卸载模型。
"""
fresh = self.db.get_batch_video(item["id"])
if fresh is None or fresh["status"] in ("COMPLETED", "SKIPPED"):
return None
video = Path(fresh["video_path"])
if not video.is_file():
self.db.update_batch_video(fresh["id"], status="FAILED", error="video file not found", updated_at=_now_iso())
return "FAILED"
work_dir = Path(fresh["work_dir"])
# 失败节点在当前阶段之前:本轮不再推进(否则会在本地 LLM 已常驻时重跑
# ASR 抢显存),留待下一次引擎循环从其失败节点重试。必须在复位 run 状态
# 之前判断,否则 _ensure_run 已把 FAILED 改成 QUEUED、判断会失效。
previous = self.db.get_run(fresh["run_id"]) if fresh.get("run_id") else None
if previous is not None and previous["status"] == "FAILED":
failed_node = previous["current_node_id"]
failed_index = order.index(failed_node) if failed_node in order else 0
if failed_index != order.index(node_spec.id):
return "FAILED"
run_id = self._ensure_run(job, fresh, version, work_dir)
run_dir = work_dir / "runs" / run_id
run_dir.mkdir(parents=True, exist_ok=True)
# 清除可能残留的信号(重启/强杀/异常中断后):暂停信号会让本次执行误暂停,
# 保持常驻信号会让后续单独重跑该节点时不再卸载模型。
(run_dir / PAUSE_FLAG).unlink(missing_ok=True)
(run_dir / KEEP_MODEL_FLAG).unlink(missing_ok=True)
if node_spec.node_type.startswith(LLM_NODE_PREFIX):
(run_dir / KEEP_MODEL_FLAG).write_text("", encoding="utf-8")
try:
scheduler = WorkflowScheduler(self.db, work_dir)
scheduler.execute_run(run_id, stop_after=None if is_last_stage else node_spec.id)
finally:
# 信号只在本阶段有效:残留会让后续单独重跑该节点时也不卸载模型。
(run_dir / KEEP_MODEL_FLAG).unlink(missing_ok=True)
run = self.db.get_run(run_id)
if run is None:
# execute_run 期间 run 记录被删除(极端外部操作),直接返回。
return None
if run["status"] == "COMPLETED":
# 放置最终产物到视频旁并清理过程文件。
self._finalize_video(fresh, run_id, video, work_dir, definition)
return "COMPLETED"
# FAILED 或 PAUSED:由调用方根据 run 状态更新视频状态与任务状态。
self.db.update_batch_video(fresh["id"], status=run["status"], error=run.get("error"), updated_at=_now_iso())
return run["status"]
def _ensure_run(self, job: dict, item: dict, version: dict, work_dir: Path) -> str:
"""确保视频有可执行的 run,返回 run_id。
run 缺失时新建(input_uri 指向本地视频,不上传副本);已存在的按断点
续跑语义复位:PAUSED/FAILED 显式置 QUEUED(保留产物,只重跑未完成
节点);RUNNING 是分阶段执行的上一个阶段或进程被杀的残留,同样置回。
"""
video = Path(item["video_path"])
work_dir.mkdir(parents=True, exist_ok=True)
run_id = item.get("run_id")
if run_id is not None and self.db.get_run(run_id) is None:
# run 记录已不存在(收尾异常删除了 run 但状态未同步):重新新建。
run_id = None
if run_id is None:
run_id = f"run_{uuid.uuid4().hex[:12]}"
now = _now_iso()
self.db.create_run({
"id": run_id,
"workflow_id": job["workflow_id"],
"workflow_version": int(version["version"]),
"status": "QUEUED",
"progress": 0,
"input_uri": str(video),
"param_overrides": None,
"source": "batch",
"created_at": now,
"updated_at": now,
})
self.db.update_batch_video(item["id"], run_id=run_id, updated_at=_now_iso())
run = self.db.get_run(run_id)
if run["status"] == "PAUSED":
# 暂停的 run 显式 resume 回 QUEUED,由 execute_run 从产物表断点续跑。
self.db.resume_run(run_id, _now_iso())
elif run["status"] in ("FAILED", "RUNNING"):
# 保留产物记录只置 QUEUEDexecute_run 跳过已完成节点、只重跑失败节点,
# 避免浪费抽帧/ASR 等长耗时成果。
self.db.update_run(run_id, status="QUEUED", error=None, updated_at=_now_iso())
return run_id
def _release_llm_model(self, params: dict) -> None:
"""LLM 阶段结束释放本机模型显存(分块流水线里每组一次,而非每视频一次)。"""
try:
release_local_model(params.get("model"))
except Exception: # noqa: BLE001 - 释放失败不影响批次推进
logger.warning("释放本地 LLM 模型失败", exc_info=True)
# ------------------------------------------------------------------
# 收尾:产物放置与过程文件清理
# ------------------------------------------------------------------
def _place_products(self, run_id: str, video: Path, definition: WorkflowDefinition) -> list[str]:
"""把最终产物文件复制到视频所在目录(视频旁),返回放置的文件名。
只为 `final_outputs` 声明的最终产物放置副本:字幕流水线的产物即日语
转写 `.srt`、中文 `.srt` 与双目 `.ass`,按库内约定命名(见 _sidecar_product_name),
文件名稳定且含视频主名——媒体库按主名匹配字幕,下次批量扫描也会命中
"已有字幕"规则跳过该视频。同名目标直接覆盖:可能是上一次运行/旧工作流
留下的旧内容,应以本次产物为准。
"""
# 先校验全部必需输出,缺文件时不覆盖任何视频旁成品,更不能继续清理。
products: list[tuple[Path, Path]] = []
for alias in definition.final_outputs:
artifact = self.db.get_artifact(run_id, alias)
if artifact is None:
raise ValueError(f"missing final artifact: {alias}")
source = Path(artifact["uri"])
if not source.is_file():
raise ValueError(f"missing final artifact file: {alias} ({source})")
target = video.parent / _sidecar_product_name(video, source, alias)
products.append((source, target))
placed: list[str] = []
for source, target in products:
# 稳定目标名 + 原子替换:复制失败保留旧成品,异常交调用方记录 FAILED。
# 多文件中途失败仍保留完整工作空间,下次可以幂等地重新放置。
atomic_copy(source, target)
placed.append(target.name)
return placed
def _finalize_video(
self,
item: dict,
run_id: str,
video: Path,
work_dir: Path,
definition: WorkflowDefinition,
) -> None:
"""视频成功处理后的收尾:产物放视频旁、清理 run 记录与过程文件。
顺序:先复制最终产物(失败则保持现状可重试),再删除 run 与产物
记录(产物已复制到视频旁不再依赖原文件),最后删除整个工作空间
(音频/分块/帧图等过程文件一并清除,防止影视库把切片数据当视频)。
"""
placed = self._place_products(run_id, video, definition)
self.db.delete_run(run_id)
self.db.update_batch_video(item["id"], status="COMPLETED", error=None, run_id=None, updated_at=_now_iso())
shutil.rmtree(work_dir, ignore_errors=True)
logger.info(
"视频 %s 处理完成,产物已放视频旁: %s,过程文件已清理",
video.name, ", ".join(placed) if placed else "(无)",
)
# ------------------------------------------------------------------
# 暂停/继续
# ------------------------------------------------------------------
def pause_job(self, job_id: str) -> None:
"""暂停批量任务:停止当前 run 与后续视频处理。
先把任务置为 PAUSED(引擎在视频间检查后停下),再暂停所有排队/运行
中的 run 并写 paused.flagwhisper 按分块、OCR 按帧检查后中止)。
"""
self.db.update_batch_job(job_id, status="PAUSED", updated_at=_now_iso())
for item in self.db.list_batch_videos(job_id):
if not item.get("run_id"):
continue
run = self.db.get_run(item["run_id"])
if run is None or run["status"] not in ("QUEUED", "RUNNING"):
continue
self.db.pause_run(item["run_id"], _now_iso())
run_dir = Path(item["work_dir"]) / "runs" / item["run_id"]
run_dir.mkdir(parents=True, exist_ok=True)
(run_dir / PAUSE_FLAG).write_text("", encoding="utf-8")
def resume_job(self, job_id: str) -> None:
"""继续批量任务:置回 QUEUED,引擎从断点续跑(PAUSED 视频逐个 resume)。"""
self.db.update_batch_job(job_id, status="QUEUED", updated_at=_now_iso())