fix: ffmpeg 串行化(最多1个同时运行),上传不限制

问题:enqueue_task 每任务起一个 ffmpeg 线程,4 个文件同时跑 ffmpeg,
CPU/磁盘 IO 争抢导致整体变慢。

修复:
- scheduler.py: ffmpeg 改为单工作线程 + queue.Queue 串行队列
  enqueue_task 只入队不入线程,最多 1 个 ffmpeg 同时运行
- 上传接收不受限:complete 立即返回,不等待 ffmpeg
- GPU 调度线程不变(仍串行,与 ffmpeg 并行)

验证:4 文件并发上传,日志确认提取严格串行(40→41→42→43 无重叠),
排队任务显示 queued 状态,4/4 成功(21.0s)
This commit is contained in:
audio2text dev
2026-07-06 22:08:59 +08:00
parent 2f68c7e1f8
commit e0dd987dba

View File

@@ -1,11 +1,11 @@
"""任务调度器ffmpeg 异步提取 + GPU 阶段串行 + 模型复用。
"""任务调度器ffmpeg 串行提取 + GPU 阶段串行 + 模型复用。
设计动机多任务时不应串行等一个任务全跑完才下一个。ffmpeg 是纯 CPU可与 GPU 阶段
并行GPU 阶段ASR + 翻译)串行化(共享显存),但卸载模型前查队列,有同类待处理
任务就继续用当前模型,减少重复加载/卸载。
数据流:
enqueue_task ──► ffmpeg 线程每任务一个异步CPU
enqueue_task ──► ffmpeg 队列(串行,最多 1 个同时跑
│ 提取音频 → task.wav_path → status=transcribing
▼(唤醒 GPU 线程)
GPU 调度线程(单线程,常驻)
@@ -16,12 +16,18 @@
while 还有 translating 任务: translate_phase → status=done
(翻译队列空,回到 ① 等待)
并发约束:
- 文件上传无限制web 请求 + 线程池接收,不阻塞 ffmpeg/GPU
- ffmpeg 提取:最多 1 个同时运行(串行队列),避免 CPU/磁盘 IO 争抢
- GPU 阶段:串行(单 GPU 调度线程),模型不共驻
模型复用N 个任务的模型切换次数从 2N 降到最优 2 次(一批 ASR 全做完 → 切翻译 → 一批翻译全做完)。
"""
from __future__ import annotations
import logging
import queue
import threading
from datetime import datetime, timezone
@@ -37,38 +43,41 @@ logger = logging.getLogger("audio2text.scheduler")
# 唤醒 GPU 调度线程的事件(新任务入队或 ffmpeg 完成时 set
_wake_event = threading.Event()
# ffmpeg 任务队列(串行执行,最多 1 个同时跑)
_ffmpeg_queue: queue.Queue[int | None] = queue.Queue()
# GPU 调度线程单例
_scheduler_thread: threading.Thread | None = None
_ffmpeg_thread: threading.Thread | None = None
_scheduler_started = False
def enqueue_task(task_id: int) -> None:
"""任务入队: ffmpeg 线程提取音频 + 唤醒 GPU 调度线程。
"""任务入队:放入 ffmpeg 队列(串行提取) + 唤醒 GPU 调度线程。
替代旧 pipeline.enqueue_task每任务一个线程跑完整管线
ffmpeg 在独立线程跑CPU与 GPU 并行),完成后 GPU 调度线程接管 ASR+翻译
ffmpeg 最多 1 个同时运行CPU/磁盘 IO 限制),其余排队
上传接收不受限——complete 创建 Task 后立即返回,不等待 ffmpeg
"""
_ensure_scheduler_running()
t = threading.Thread(
target=_extract_audio_async, args=(task_id,),
name=f"ffmpeg-{task_id}", daemon=True,
)
t.start()
logger.info("任务 %d 已入队,开始音频提取。", task_id)
_ffmpeg_queue.put(task_id)
logger.info("任务 %d 已入队,等待音频提取。", task_id)
def start_scheduler() -> None:
"""启动 GPU 调度线程(应用启动时调一次,幂等)。"""
global _scheduler_thread, _scheduler_started
"""启动 ffmpeg + GPU 调度线程(应用启动时调一次,幂等)。"""
global _scheduler_thread, _ffmpeg_thread, _scheduler_started
if _scheduler_started:
return
_scheduler_started = True
_reset_stuck_tasks()
_ffmpeg_thread = threading.Thread(
target=_ffmpeg_worker, name="ffmpeg-worker", daemon=True,
)
_ffmpeg_thread.start()
_scheduler_thread = threading.Thread(
target=_gpu_scheduler, name="gpu-scheduler", daemon=True,
)
_scheduler_thread.start()
logger.info("GPU 调度线程已启动。")
logger.info("ffmpeg + GPU 调度线程已启动。")
def _reset_stuck_tasks() -> None:
@@ -116,16 +125,24 @@ def _ensure_scheduler_running() -> None:
start_scheduler()
# ---------------- ffmpeg 异步提取 ----------------
# ---------------- ffmpeg 串行提取 ----------------
def _extract_audio_async(task_id: int) -> None:
"""独立线程跑 ffmpeg 提取CPU完成后唤醒 GPU 调度线程。"""
def _ffmpeg_worker() -> None:
"""常驻 ffmpeg 工作线程:从队列取任务,串行提取音频(最多 1 个同时跑)。
队列收到 None 为停机信号(当前不使用,保留用于优雅关闭)。
"""
logger.info("ffmpeg 工作线程开始运行。")
while True:
task_id = _ffmpeg_queue.get()
if task_id is None:
break # 停机信号
db = get_session_local()()
try:
task = db.get(Task, task_id)
if task is None:
logger.error("任务 %d 不存在ffmpeg 线程退出", task_id)
return
logger.error("任务 %d 不存在ffmpeg 跳过", task_id)
continue
pipeline.extract_phase(db, task)
_wake_event.set() # 通知 GPU 线程有新任务
except Exception as exc:
@@ -133,6 +150,7 @@ def _extract_audio_async(task_id: int) -> None:
pipeline.mark_failed(db, task_id, str(exc))
finally:
db.close()
_ffmpeg_queue.task_done()
# ---------------- GPU 调度线程 ----------------