"""queued 잡을 하나씩 집어 다음 게이트까지 돌리는 배경 실행기 큐가 DB에 있어 동시 실행 수를 올리거나 컨테이너를 늘려도 파이프라인 코드는 그대로 승인이 status를 queued로 되돌리면 다음 폴링에서 이어 감 """ import asyncio import traceback from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession from pipelines.runner import QUEUED, RUNNING, run_until_gate from tables.task import PlayreelTask, PosterAliveTask from utils.database import session_factory POLL_INTERVAL = 2.0 TABLES = (PlayreelTask, PosterAliveTask) _worker: asyncio.Task | None = None async def claim_next(session: AsyncSession): """queued 잡 하나를 running으로 바꿔 가져온다 워커가 하나인 동안은 경합이 없다. 늘릴 때 SELECT ... FOR UPDATE SKIP LOCKED로 바꾼다 """ for table in TABLES: found = await session.scalars( select(table).where(table.status == QUEUED) .order_by(table.created_at).limit(1)) task = found.first() if task is not None: task.status = RUNNING await session.commit() return task return None async def run_once() -> bool: async with session_factory() as session: task = await claim_next(session) if task is None: return False await run_until_gate(session, task) return True async def requeue_running() -> None: """기동 시 running으로 남은 잡을 되살린다 — 어디까지 갔는지는 state가 기억한다""" async with session_factory() as session: for table in TABLES: for task in await session.scalars(select(table).where(table.status == RUNNING)): task.status = QUEUED await session.commit() async def loop() -> None: await requeue_running() while True: try: if not await run_once(): await asyncio.sleep(POLL_INTERVAL) except asyncio.CancelledError: raise except Exception: # 워커는 죽지 않는다. DB가 끊긴 경우라 잠시 쉬고 다시 본다 traceback.print_exc() await asyncio.sleep(POLL_INTERVAL) def start() -> None: global _worker _worker = asyncio.create_task(loop(), name="pipeline-worker") async def stop() -> None: if _worker is not None: _worker.cancel() await asyncio.gather(_worker, return_exceptions=True)