큐를 실제로 돌린다: 적재 → 워커 claim → 핸들러 실행 → 결과 저장 → DONE. API(적재)와 워커(소비)를 분리 프로세스로(코드베이스 공유, 독립 스케일). - worker/notify: 전용 asyncpg LISTEN 리스너. enqueue 에서 pg_notify → 유휴 워커 즉시 기상(폴링 제거) - worker/runner: Worker(claim→처리, 처리중 heartbeat 로 lease 갱신, complete/fail) + run_reaper - worker/handlers: job_type 별 핸들러(주입식). SEARCH=소스 어댑터 검색→정규화 결과 - worker_main: API 분리 워커 진입점(브라우저 무거워 기본 동시성 1) - job_crud.enqueue: 삽입 시 pg_notify (중복 스킵 시엔 미발생) - tests: drain→DONE·실패→재시도→dead·reaper 회수 후 재처리·NOTIFY 기상 4건 (전체 20/20) - 라이브 E2E 확인: 적재→워커가 실제 쿠팡 검색(8건)→DONE Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
85 lines
3.3 KiB
Python
85 lines
3.3 KiB
Python
"""워커 루프 + reaper.
|
|
|
|
워커는 큐에서 잡을 원자적으로 claim → 핸들러 실행 → complete/fail 한다.
|
|
- 처리 중 heartbeat 로 lease 를 갱신(긴 잡이 reaper 에 회수되지 않게).
|
|
- claim 이 비면 LISTEN 알림 또는 짧은 타임아웃까지 대기(폴링 최소화).
|
|
- 핸들러는 주입식(async def(job)->dict) — 프로덕션은 검색 파이프라인, 테스트는 fake.
|
|
"""
|
|
|
|
import asyncio
|
|
|
|
from common.enums import JobStatus
|
|
from common.logger import LOG
|
|
from crud.job_crud import JobQueue, compute_backoff
|
|
|
|
|
|
class Worker:
|
|
def __init__(self, worker_id: str, queue: JobQueue, handler, lease_sec: int = 120, backoff_fn=compute_backoff):
|
|
self.worker_id = worker_id
|
|
self.queue = queue
|
|
self.handler = handler
|
|
self.lease_sec = lease_sec
|
|
self.backoff_fn = backoff_fn
|
|
|
|
async def process_one(self) -> bool:
|
|
"""대기 잡 1건을 claim·처리. 처리했으면 True, 없으면 False."""
|
|
job = await self.queue.claim(self.worker_id, self.lease_sec)
|
|
if not job:
|
|
return False
|
|
await self._process(job)
|
|
return True
|
|
|
|
async def drain(self) -> int:
|
|
"""큐가 빌 때까지 처리(테스트/일회성 배치용). 처리한 잡 수 반환."""
|
|
n = 0
|
|
while await self.process_one():
|
|
n += 1
|
|
return n
|
|
|
|
async def run(self, listener=None, stop: asyncio.Event | None = None, idle_timeout: float = 5.0):
|
|
"""상시 루프. stop 이 설정될 때까지 처리하고, 유휴 시 알림/타임아웃까지 대기."""
|
|
stop = stop or asyncio.Event()
|
|
while not stop.is_set():
|
|
worked = await self.process_one()
|
|
if not worked:
|
|
if listener is not None:
|
|
await listener.wait(idle_timeout)
|
|
else:
|
|
await asyncio.sleep(idle_timeout)
|
|
|
|
async def _process(self, job: dict):
|
|
jid = job["job_id"]
|
|
hb = asyncio.create_task(self._heartbeat(jid))
|
|
try:
|
|
result = await self.handler(job)
|
|
await self.queue.complete(jid, self.worker_id, result)
|
|
LOG.d(f"[{self.worker_id}] done {jid}")
|
|
except Exception as ex:
|
|
backoff = self.backoff_fn(job["attempts"])
|
|
st = await self.queue.fail(jid, self.worker_id, f"{type(ex).__name__}: {ex}", backoff)
|
|
LOG.w(f"[{self.worker_id}] fail {jid} → {JobStatus(st).name if st else '?'} ({type(ex).__name__}: {ex})")
|
|
finally:
|
|
hb.cancel()
|
|
try:
|
|
await hb
|
|
except asyncio.CancelledError:
|
|
pass
|
|
|
|
async def _heartbeat(self, jid: str):
|
|
interval = max(1, self.lease_sec // 3)
|
|
while True:
|
|
await asyncio.sleep(interval)
|
|
await self.queue.renew_lease(jid, self.worker_id, self.lease_sec)
|
|
|
|
|
|
async def run_reaper(queue: JobQueue, stop: asyncio.Event, interval: float = 30.0):
|
|
"""만료 lease(워커 사망) 잡을 주기적으로 회수. 재시도 남으면 재큐, 소진되면 DEAD."""
|
|
while not stop.is_set():
|
|
reclaimed = await queue.reap()
|
|
if reclaimed:
|
|
LOG.w(f"[reaper] reclaimed {len(reclaimed)} stale job(s)")
|
|
try:
|
|
await asyncio.wait_for(stop.wait(), interval)
|
|
except asyncio.TimeoutError:
|
|
pass
|