From 0c81f5b04ba986069ad1a3535edfd703cb18b015 Mon Sep 17 00:00:00 2001 From: claude_dev Date: Wed, 26 Aug 2026 20:19:07 +0800 Subject: [PATCH] =?UTF-8?q?fix(orchestrator):=20worker=E8=A2=ABOOM?= =?UTF-8?q?=E5=87=BB=E6=9D=80=E5=90=8EProcessPoolExecutor=E6=B0=B8?= =?UTF-8?q?=E4=B9=85=E7=A0=B4=E6=8D=9F=E2=80=94=E2=80=94BrokenProcessPool?= =?UTF-8?q?=E8=AE=A9=E4=BB=BB=E5=8A=A1=E4=B8=AD=E5=BF=83=E7=98=AB=E7=97=AA?= =?UTF-8?q?=E5=88=B0=E5=AE=B9=E5=99=A8=E9=87=8D=E5=90=AF(factor=E6=89=B9?= =?UTF-8?q?=E6=AC=A1=E8=BF=9E=E6=92=9E=E4=B8=A4=E6=AC=A1500=E5=AE=9E?= =?UTF-8?q?=E9=94=A4);submit=5Fwork=E6=8D=95=E8=8E=B7=E5=90=8E=E9=87=8D?= =?UTF-8?q?=E5=BB=BA=E6=B1=A0=E9=87=8D=E8=AF=95=E4=B8=80=E6=AC=A1=20[nas]?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- sanguo_orchestrator/pool.py | 17 +++++++++++++++-- 1 file changed, 15 insertions(+), 2 deletions(-) diff --git a/sanguo_orchestrator/pool.py b/sanguo_orchestrator/pool.py index e23d283..ac0bc27 100644 --- a/sanguo_orchestrator/pool.py +++ b/sanguo_orchestrator/pool.py @@ -4,6 +4,7 @@ Provides task storage and status tracking (not actual multiprocessing) """ import logging from concurrent.futures import ProcessPoolExecutor, Future +from concurrent.futures.process import BrokenProcessPool from multiprocessing import get_context from .task import Task, TaskState @@ -28,9 +29,21 @@ class TaskPool: return task def submit_work(self, task_id: str, func, *args) -> Future: - """Submit work to the process pool executor""" + """Submit work to the process pool executor. + + worker 被 OOM 击杀后 ProcessPoolExecutor 永久破损(BrokenProcessPool), + 重建池重试一次,避免一次异常让任务中心瘫痪到容器重启。 + """ logger.debug("submit_work task_id=%s", task_id) - return self.executor.submit(func, *args) + try: + return self.executor.submit(func, *args) + except BrokenProcessPool: + logger.warning("pool broken (worker died?), rebuilding for %s", task_id) + self.executor.shutdown(wait=False) + self.executor = ProcessPoolExecutor( + max_workers=self.max_workers, mp_context=get_context("spawn") + ) + return self.executor.submit(func, *args) def update_stage(self, task_id: str, stage: str): """Update the stage of a task"""