feat: 文件夹批量处理引擎(后端)

- BatchWorker 单线程轮询 batch_jobs 表,处理 source=batch 的运行,
  与主调度器互不抢占(next_queued_run 排除 batch 来源)
- 直接读取用户所选文件夹下的视频逐个执行流水线,不上传到工作目录;
  中间态与产物落在视频旁同名文件夹,batch.done.json 完成标记去重
- 支持暂停/继续、失败容错(单视频失败不阻塞后续)、删除任务只清库
- 孤儿清理跳过 source=batch 运行,防止误删用户视频文件夹
- workflow_runs 新增 source 列(upload/batch),旧库自动迁移
This commit is contained in:
2026-08-23 16:25:09 +08:00
parent b16e0e9f3c
commit 711867e79f
10 changed files with 2056 additions and 9 deletions
+480
View File
@@ -0,0 +1,480 @@
"""文件夹批量处理引擎。
本地版的核心能力:**不把视频上传到工作目录**,而是直接读取用户所选文件夹
下的所有视频,逐个调用现有的工作流流水线(复用 WorkflowScheduler 的 DAG
执行与断点续跑逻辑)。
数据落盘约定:
- 每个视频的中间态数据(runs/、chunks/、帧图等)与最终产物都存放在**视频
所在目录的同名文件夹**里(movie.mp4 → movie/),源视频目录保持干净。
- 最终产物(SRT/ASS)在任务完成后从节点产物目录复制到同名文件夹根目录,
同时写入 `batch.done.json` 完成标记;再次批量处理同一文件夹时,已有完成
标记且产物文件齐全的视频直接跳过(已经处理过的不再处理)。
暂停/恢复语义(对应前端"暂停/继续"按钮):
- 暂停批量任务:把批量任务置为 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 datetime import datetime, timezone
from pathlib import Path
from wov_app import registry
from wov_app.config import BATCH_INTERVAL_SECONDS
from wov_app.db import Database
from wov_app.logging import get_logger
from wov_app.scheduler import WorkflowScheduler
from wov_sdk.models import WorkflowDefinition
# 批量引擎运行日志:任务进度、视频逐个处理与暂停/续跑等状态变化。
logger = get_logger("batch")
# 识别为视频文件的扩展名(大小写不敏感,扫描时统一转小写比较)。
VIDEO_EXTENSIONS = {
".mp4", ".mkv", ".avi", ".mov", ".webm", ".flv", ".ts",
".m4v", ".wmv", ".mpg", ".mpeg", ".3gp",
}
# 暂停信号文件名:与节点约定一致,位于 run 根目录(<work_dir>/runs/<run_id>/)。
PAUSE_FLAG = "paused.flag"
# 视频完成标记文件名:位于同名文件夹根目录,记录该视频已完成的工作流与最终
# 产物文件名,跨批量任务去重(已经处理过的不再处理)。
MARKER_NAME = "batch.done.json"
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 work_dir_for(video: Path) -> Path:
"""返回视频的同名文件夹:去掉扩展名,位于视频所在目录。
例如 movie.mp4 → 旁边的 movie/ 文件夹,中间态与最终产物都放这里。
"""
return video.parent / video.stem
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 ensure_videos(db: Database, job_id: str, videos: list[Path]) -> None:
"""为扫描到的视频补齐 batch_videos 明细;已存在的记录保持不变。
同一批量任务反复处理(暂停/续跑)时保留每个视频的状态,已完成的不重置。
"""
existing = {item["video_path"] for item in db.list_batch_videos(job_id)}
now = _now_iso()
for video in videos:
path = str(video)
if path in existing:
continue
db.create_batch_video({
"id": f"bv_{uuid.uuid4().hex[:12]}",
"job_id": job_id,
"video_path": path,
"work_dir": str(work_dir_for(video)),
"run_id": None,
"status": "PENDING",
"error": None,
"created_at": now,
"updated_at": now,
})
def create_job(
db: Database,
folder_path: str,
workflow_id: str,
recursive: bool = True,
) -> str:
"""创建批量任务:校验文件夹与工作流、扫描视频、落库明细,返回任务 ID。
校验失败抛出 ValueError(由路由层转为 422 响应);扫描到的视频全部
登记为 batch_videos 明细,引擎轮询到该任务后逐个处理。
"""
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),
"status": "QUEUED",
"progress": 0,
"total": 0,
"done": 0,
"failed": 0,
"current_video": None,
"error": None,
"created_at": now,
"updated_at": now,
})
ensure_videos(db, job_id, videos)
logger.info("创建批量任务 %s: 文件夹 %s, 工作流 %s, 视频 %d", job_id, folder, workflow_id, len(videos))
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()
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 统一处理)。"""
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()
# 扫描当前文件夹的视频,为新增视频补齐明细(已有明细保留原状态,
# 保证暂停/续跑时已完成与进行中的视频不被重置)。
self._ensure_videos(job_id, scan_videos(folder, bool(job["recursive"])))
items = self.db.list_batch_videos(job_id)
total = len(items)
self.db.update_batch_job(
job_id, status="RUNNING", total=total, progress=0,
current_video=None, error=None, updated_at=_now_iso(),
)
for index, item in enumerate(items):
# 暂停检查:批量任务被暂停后停止处理后续视频,等待用户继续。
current = self.db.get_batch_job(job_id)
if current is None or current["status"] == "PAUSED":
logger.info("批量任务 %s 已暂停,停止在视频 %s", job_id, item["video_path"])
return
# 已完成/已跳过的视频不再处理。
if item["status"] in ("COMPLETED", "SKIPPED"):
continue
video = Path(item["video_path"])
if not video.is_file():
self.db.update_batch_video(item["id"], status="FAILED", error="video file not found", updated_at=_now_iso())
continue
work_dir = Path(item["work_dir"])
# 同名文件夹里已有同工作流的完成标记且产物齐全 → 直接跳过。
if self._is_done(work_dir, job["workflow_id"]):
self.db.update_batch_video(item["id"], status="SKIPPED", updated_at=_now_iso())
continue
self.db.update_batch_job(
job_id, current_video=str(video),
progress=index / total if total else 0,
updated_at=_now_iso(),
)
try:
self._process_video(job, item, version, definition, work_dir)
except Exception as exc: # noqa: BLE001
# 单视频兜底:不中断整个批量任务,记录错误后继续下一个视频。
logger.exception("批量任务 %s 视频 %s 处理异常", job_id, video)
self.db.update_batch_video(item["id"], status="FAILED", error=str(exc), updated_at=_now_iso())
# 重新读取视频明细:_process_video 可能刚创建 run(快照里 run_id
# 还是 None),必须取最新记录才能拿到 run_id 判断暂停状态。
item = self.db.get_batch_video(item["id"])
# 视频处理中被暂停:批量任务整体保持 PAUSED,等待用户继续。
run = self.db.get_run(item["run_id"]) if item and item.get("run_id") else None
if run is not None and run["status"] == "PAUSED":
self.db.update_batch_video(item["id"], status="PAUSED", updated_at=_now_iso())
self.db.update_batch_job(job_id, status="PAUSED", updated_at=_now_iso())
return
# 全部视频处理完成:汇总已处理与失败数量,任务置为 COMPLETED。
items = self.db.list_batch_videos(job_id)
done = sum(1 for item in items if item["status"] in ("COMPLETED", "SKIPPED"))
failed = sum(1 for item in items if item["status"] == "FAILED")
self.db.update_batch_job(
job_id, status="COMPLETED", progress=1.0, done=done, failed=failed,
current_video=None, error=None, updated_at=_now_iso(),
)
logger.info(
"批量任务 %s 完成: 共 %d 个视频, 完成/跳过 %d, 失败 %d",
job_id, total, done, failed,
)
def _process_video(
self,
job: dict,
item: dict,
version: dict,
definition: WorkflowDefinition,
work_dir: Path,
) -> None:
"""处理单个视频:建 run(复用现有调度器)并执行,完成后复制最终产物。
per-video 的 WorkflowScheduler 以同名文件夹为 storage,中间态落在
<work_dir>/runs/<run_id>/steps/ 下;产物表记录全部节点输出,暂停后
续跑从产物表重建已完成节点(断点续跑)。
"""
video = Path(item["video_path"])
work_dir.mkdir(parents=True, exist_ok=True)
run_id = item.get("run_id")
if run_id is None:
# 首次处理:创建 source=batch 的运行,input_uri 直接指向本地视频
# (不再上传副本),调度器按工作流 DAG 动态组装节点执行。
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"] == "COMPLETED":
self._copy_finals(run_id, work_dir, definition, job["workflow_id"])
self.db.update_batch_video(item["id"], status="COMPLETED", error=None, updated_at=_now_iso())
return
# 暂停的 run 显式 resume 回 QUEUED,由 execute_run 从产物表断点续跑。
if run["status"] == "PAUSED":
self.db.resume_run(run_id, _now_iso())
elif run["status"] == "FAILED":
# 失败重跑:**保留**已完成节点的产物记录,只恢复 QUEUED——
# execute_run 从产物表重建已完成节点并跳过,只重跑失败节点。
# 不再 reset_run 清空产物:extract/ocr 等长耗时节点的成果(如
# ABP-885 的 22222 帧 OCR)会被白白丢弃重做(run_e2b74e89e232 实测)。
self.db.update_run(run_id, status="QUEUED", error=None, updated_at=_now_iso())
elif run["status"] == "RUNNING":
# 上次进程被杀残留:恢复 QUEUED(保留产物)由 execute_run 续跑。
self.db.update_run(run_id, status="QUEUED", updated_at=_now_iso())
# 清除可能残留的暂停信号(重启/异常中断后),避免本次执行误暂停。
(work_dir / "runs" / run_id / PAUSE_FLAG).unlink(missing_ok=True)
scheduler = WorkflowScheduler(self.db, work_dir)
scheduler.execute_run(run_id)
run = self.db.get_run(run_id)
if run["status"] == "COMPLETED":
# 复制最终产物到同名文件夹根目录并写完成标记(已处理过不再处理)。
self._copy_finals(run_id, work_dir, definition, job["workflow_id"])
self.db.update_batch_video(item["id"], status="COMPLETED", error=None, updated_at=_now_iso())
else:
# FAILED 或 PAUSED:由调用方根据 run 状态更新视频状态与任务状态。
self.db.update_batch_video(item["id"], status=run["status"], error=run.get("error"), updated_at=_now_iso())
# ------------------------------------------------------------------
# 明细与产物
# ------------------------------------------------------------------
def _ensure_videos(self, job_id: str, videos: list[Path]) -> None:
"""为扫描到的视频补齐 batch_videos 明细;已存在的记录保持不变。"""
ensure_videos(self.db, job_id, videos)
def _is_done(self, work_dir: Path, workflow_id: str) -> bool:
"""判断同名文件夹是否已完成当前工作流的处理。
完成标记记录 workflow_id 与最终产物文件名;只有工作流一致且产物文件
全部存在时才视为已处理(不同工作流的产物不互相误判为完成)。
"""
marker = load_marker(work_dir)
if marker is None or marker.get("workflow_id") != workflow_id:
return False
finals = marker.get("finals") or {}
return bool(finals) and all((work_dir / name).is_file() for name in finals.values())
def _copy_finals(
self,
run_id: str,
work_dir: Path,
definition: WorkflowDefinition,
workflow_id: str,
) -> None:
"""把最终产物从节点目录复制到同名文件夹根目录,并写完成标记。
调度器收尾时已把产物重命名为 上传文件名.标识.时间戳(如
movie.zh-CN.20260819120000.srt),这里原样复制,文件名保留辨识度。
"""
work_dir.mkdir(parents=True, exist_ok=True)
finals: dict[str, str] = {}
for alias in definition.final_outputs:
artifact = self.db.get_artifact(run_id, alias)
if artifact is None:
continue
source = Path(artifact["uri"])
if not source.is_file():
continue
target = work_dir / source.name
# 已存在的产物直接复用,避免重复复制。
if not target.is_file() or target.stat().st_size != source.stat().st_size:
shutil.copy2(source, target)
finals[alias] = source.name
marker = {
"workflow_id": workflow_id,
"workflow_version": definition.version,
"run_id": run_id,
"completed_at": _now_iso(),
"finals": finals,
}
(work_dir / MARKER_NAME).write_text(
json.dumps(marker, ensure_ascii=False, indent=2),
encoding="utf-8",
)
# ------------------------------------------------------------------
# 暂停/继续
# ------------------------------------------------------------------
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())