78 lines
2.5 KiB
Python
78 lines
2.5 KiB
Python
"""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)
|