o2o-negosium-original/lps/worker/runner.py
민헌 abcb58ba05 feat(lps): 워커 루프 + LISTEN/NOTIFY — 큐 소비 파이프라인 가동
큐를 실제로 돌린다: 적재 → 워커 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>
2026-07-08 16:58:58 +09:00

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