Files
vrsub/scripts/fix_zombie_batch_jobs.py
cat-shark dcdc5e8604 feat: 批量分块流水线、本地模型显存让渡与任务列表分工
批量引擎改为「分块流水线」:视频按 WOV_BATCH_STAGE_GROUP_SIZE(默认 8)分组,
组内按 DAG 拓扑序跑完全部视频(全部 extract → 全部 ASR → 全部翻译 → 全部 ASS)
再进入下一组,本地模型每组只加载一次、卸载一次,而不是每个视频来回加载卸载;
产物仍按组增量落到视频旁。调度器新增 execute_run(run_id, stop_after=节点):
该节点完成后任务保持 RUNNING 不收尾,下一次调用从产物表跳过已完成节点继续,
用于实现阶段边界。

- nodes/llm.py:翻译节点结束释放本机 Ollama 显存(node 参数 unload_after >
  LLM_UNLOAD_AFTER > 本机 loopback 端点默认卸载,云端端点不卸载;卸载失败只告警),
  新增 keep_model.flag 语义(阶段内保持常驻)与 release_local_model();
  新增节点内暂停(按批 20 行检查 paused.flag,抛 PauseRequested,调度器保持 PAUSED)。
- src/wov_app/batch.py:分组阶段执行与阶段末统一释放显存;失败视频只在它失败
  节点的那个阶段重试(避免 LLM 已常驻时重跑 ASR 抢显存);任务没有明细时保持
  QUEUED 等登记完成、仍有未完成视频时置回 QUEUED 自愈(原先留 RUNNING 会卡死:
  引擎只拾取 QUEUED,任务停在“运行中但没人推进”);无失败视频时删除任务级空目录;
  每个阶段开始前清理 paused.flag / keep_model.flag,避免强杀残留影响后续阶段。
- src/wov_app/config.py:新增 WOV_BATCH_STAGE_GROUP_SIZE(设为 1 即旧的每视频全链路)。
- 任务列表与批量页分工:GET /api/runs 默认排除 source=batch(一个批量任务会产生
  N 条单视频 run,会把 20 条窗口占满;且任务管理页的暂停/重试/删除对批量 run
  语义不成立),需要排查时用 include_batch=1;作为补偿批量页详情新增阶段列
  (阶段 i/N · 中文标签,由该视频 run 的 current_node_id 在 DAG 拓扑序中的位置
  推导,节点类型映射中文标签)。阶段只有节点边界粒度,句级进度不落库、只在日志。
- 顺带纳入此前未提交的批量僵尸状态恢复:recover_interrupted_batch_jobs 除 RUNNING
  外也把「COMPLETED 但仍含未结束视频」的任务置回 QUEUED;fix_zombie_batch_jobs.py
  改为按条件扫描并支持 --apply 预览;批量页明细只列本批真正处理过的视频。

测试新增/更新:分块流水线调用顺序(组内按节点跑完再下一组)、每组只释放一次模型、
阶段内保持常驻标志、翻译按批暂停、失败视频不跨阶段推进、任务无明细/中途登记视频时
置回 QUEUED、任务工作空间与残留信号清理、任务列表默认过滤批量 run、详情阶段字段、
前端阶段列渲染;全量 507 passed(唯一失败为既有素材缺失的 integration 用例)。
2026-09-18 10:31:52 +08:00

76 lines
3.3 KiB
Python
Raw Permalink 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.
"""修复僵尸批量任务脚本:把误标 COMPLETED 但仍有未结束视频的任务置回 QUEUED。
僵尸状态:`_run_job` 在视频收尾前就把任务置 COMPLETED,留下"N 个视频待处理
却已完成"的假完成;引擎只拾取 QUEUED,剩余视频永久无人处理。代码已修复
(置 COMPLETED 前校验无未结束明细 + 重启恢复把这类任务放回队列),本脚本用于
修复**历史遗留**数据(如 batch_351833b7d446COMPLETED/done=0 但仍有 1 个 PENDING)。
修复方式:把 COMPLETED 且仍有非终态(PENDING/RUNNING/PAUSED)明细的任务置回
QUEUED 并清空误导的 progress,保留 SKIPPED/COMPLETED/FAILED 明细与 run 记录;
批量引擎重启后会重新拾起,从剩余视频续跑,全部处理完才置 COMPLETED。
安全约束:只改"COMPLETED 且存在未结束明细"的任务,正常完成的任务不动。
默认只打印预览,确认后加 --apply 实际写入。
"""
from __future__ import annotations
import sys
from datetime import datetime, timezone
from pathlib import Path
from wov_app.db import Database
# 未结束的明细状态:出现任一个就不能算任务完成。
UNFINISHED_STATUSES = ("PENDING", "RUNNING", "PAUSED")
def _now_iso() -> str:
"""返回当前 UTC 时间的 ISO 格式字符串。"""
return datetime.now(timezone.utc).isoformat()
def find_zombie_jobs(db: Database) -> list[tuple[dict, list[dict]]]:
"""返回全部僵尸任务及其未结束明细(COMPLETED 但仍有未结束视频)。"""
zombies: list[tuple[dict, list[dict]]] = []
for job in db.list_batch_jobs(limit=1000):
if job["status"] != "COMPLETED":
continue
unfinished = [
v for v in db.list_batch_videos(str(job["id"]))
if v["status"] in UNFINISHED_STATUSES
]
if unfinished:
zombies.append((job, unfinished))
return zombies
def main(db_path: str, apply: bool = False) -> None:
"""诊断(默认)或修复(--apply)僵尸批量任务。"""
db = Database(Path(db_path))
print(f"数据库: {db_path}\n")
zombies = find_zombie_jobs(db)
if not zombies:
print(" [无] 没有 COMPLETED 但仍含未结束视频的批量任务")
return
for job, unfinished in zombies:
job_id = str(job["id"])
videos = db.list_batch_videos(job_id)
completed = sum(1 for v in videos if v["status"] == "COMPLETED")
print(f" [待修] {job_id}: status={job['status']} → QUEUED, "
f"COMPLETED={completed}, 未结束={len(unfinished)}"
f"{', '.join(Path(v['video_path']).name for v in unfinished[:3])}…)")
if apply:
db.update_batch_job(
job_id, status="QUEUED", progress=0,
total=int(job["total"] or 0) or len(unfinished),
done=completed, failed=int(job["failed"] or 0),
current_video=None, error=None, updated_at=_now_iso(),
)
print(f" ✓ 已置回 QUEUED(引擎将续跑剩余 {len(unfinished)} 个视频)")
if not apply:
print("\n以上为预览。确认无误后加 --apply 实际修复。")
if __name__ == "__main__":
main(sys.argv[1] if len(sys.argv) > 1 else "data/wov.db", "--apply" in sys.argv)