From e0dd987dba49f0976b2fcff156775f924ed57bac Mon Sep 17 00:00:00 2001 From: audio2text dev Date: Mon, 6 Jul 2026 22:08:59 +0800 Subject: [PATCH] =?UTF-8?q?fix:=20ffmpeg=20=E4=B8=B2=E8=A1=8C=E5=8C=96?= =?UTF-8?q?=EF=BC=88=E6=9C=80=E5=A4=9A1=E4=B8=AA=E5=90=8C=E6=97=B6?= =?UTF-8?q?=E8=BF=90=E8=A1=8C=EF=BC=89=EF=BC=8C=E4=B8=8A=E4=BC=A0=E4=B8=8D?= =?UTF-8?q?=E9=99=90=E5=88=B6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 问题: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) --- app/services/scheduler.py | 78 ++++++++++++++++++++++++--------------- 1 file changed, 48 insertions(+), 30 deletions(-) diff --git a/app/services/scheduler.py b/app/services/scheduler.py index 3c000c3..24d7c71 100644 --- a/app/services/scheduler.py +++ b/app/services/scheduler.py @@ -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,23 +125,32 @@ def _ensure_scheduler_running() -> None: start_scheduler() -# ---------------- ffmpeg 异步提取 ---------------- +# ---------------- ffmpeg 串行提取 ---------------- -def _extract_audio_async(task_id: int) -> None: - """独立线程跑 ffmpeg 提取(CPU),完成后唤醒 GPU 调度线程。""" - db = get_session_local()() - try: - task = db.get(Task, task_id) - if task is None: - logger.error("任务 %d 不存在,ffmpeg 线程退出。", task_id) - return - pipeline.extract_phase(db, task) - _wake_event.set() # 通知 GPU 线程有新任务 - except Exception as exc: - logger.exception("任务 %d ffmpeg 提取失败:%s", task_id, exc) - pipeline.mark_failed(db, task_id, str(exc)) - finally: - db.close() +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) + continue + pipeline.extract_phase(db, task) + _wake_event.set() # 通知 GPU 线程有新任务 + except Exception as exc: + logger.exception("任务 %d ffmpeg 提取失败:%s", task_id, exc) + pipeline.mark_failed(db, task_id, str(exc)) + finally: + db.close() + _ffmpeg_queue.task_done() # ---------------- GPU 调度线程 ----------------