From abcb58ba050bbc1a7c594c0bc5b74f06416e36c3 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=EB=AF=BC=ED=97=8C?= Date: Wed, 8 Jul 2026 16:58:58 +0900 Subject: [PATCH] =?UTF-8?q?feat(lps):=20=EC=9B=8C=EC=BB=A4=20=EB=A3=A8?= =?UTF-8?q?=ED=94=84=20+=20LISTEN/NOTIFY=20=E2=80=94=20=ED=81=90=20?= =?UTF-8?q?=EC=86=8C=EB=B9=84=20=ED=8C=8C=EC=9D=B4=ED=94=84=EB=9D=BC?= =?UTF-8?q?=EC=9D=B8=20=EA=B0=80=EB=8F=99?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 큐를 실제로 돌린다: 적재 → 워커 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) --- lps/crud/job_crud.py | 7 ++++ lps/tests/test_worker.py | 69 +++++++++++++++++++++++++++++++++ lps/worker/handlers.py | 35 +++++++++++++++++ lps/worker/notify.py | 47 ++++++++++++++++++++++ lps/worker/runner.py | 84 ++++++++++++++++++++++++++++++++++++++++ lps/worker_main.py | 52 +++++++++++++++++++++++++ 6 files changed, 294 insertions(+) create mode 100644 lps/tests/test_worker.py create mode 100644 lps/worker/handlers.py create mode 100644 lps/worker/notify.py create mode 100644 lps/worker/runner.py create mode 100644 lps/worker_main.py diff --git a/lps/crud/job_crud.py b/lps/crud/job_crud.py index dcd4001..44633ad 100644 --- a/lps/crud/job_crud.py +++ b/lps/crud/job_crud.py @@ -16,6 +16,10 @@ from common.database.db_session_manager import DB_SESSION_MNG from common.enums import DBType, DBWRType, JobStatus +# 잡 적재 시 워커를 즉시 깨우는 LISTEN/NOTIFY 채널(폴링 제거). +JOB_NOTIFY_CHANNEL = "lps_job" + + def compute_backoff(attempts: int, base: float = 5.0, cap: float = 600.0) -> float: """지수 백오프(초). attempts 회 시도 후 다음 재시도까지 대기 = base * 2^(attempts-1), cap 상한.""" return min(cap, base * (2 ** max(0, attempts - 1))) @@ -53,6 +57,9 @@ class JobQueue: "t": job_type, "p": priority, "payload": json.dumps(payload), "dk": dedupe_key, "ma": max_attempts, })).first() + if row: + # 커밋 시 전달됨 → LISTEN 중인 유휴 워커를 즉시 깨운다(중복 스킵 시엔 알림 안 함). + await s.execute(text("SELECT pg_notify(:ch, '')"), {"ch": JOB_NOTIFY_CHANNEL}) return str(row[0]) if row else None return await self._tx(run) diff --git a/lps/tests/test_worker.py b/lps/tests/test_worker.py new file mode 100644 index 0000000..c8396ae --- /dev/null +++ b/lps/tests/test_worker.py @@ -0,0 +1,69 @@ +"""워커 루프 테스트 — drain→DONE, 실패→재시도→dead, reaper 회수 후 재처리, NOTIFY 깨움. +핸들러는 fake(브라우저 없이) — 워커 로직만 결정론적으로 검증한다.""" + +import asyncio + +import pytest_asyncio +from sqlalchemy import text + +from common.enums import JobStatus, JobType +from crud.job_crud import JobQueue +from worker.notify import JobListener +from worker.runner import Worker + + +@pytest_asyncio.fixture +async def q(db_engine): + async with db_engine.begin() as conn: + await conn.execute(text("TRUNCATE job")) + return JobQueue() + + +async def test_worker_drains_all_to_done(q): + for i in range(3): + await q.enqueue(JobType.SEARCH.value, {"product_name": f"item{i}"}) + + async def handler(job): + return {"ok": True, "q": job["payload"]["product_name"]} + + processed = await Worker("w1", q, handler).drain() + assert processed == 3 + assert (await q.counts())["DONE"] == 3 + + +async def test_worker_failure_retries_then_dead(q): + await q.enqueue(JobType.SEARCH.value, {"product_name": "x"}, max_attempts=2) + + async def boom(job): + raise RuntimeError("nope") + + w = Worker("w1", q, boom, backoff_fn=lambda a: 0) # 백오프 0 → 즉시 재시도 가능 + assert await w.process_one() is True # 1/2 실패 → PENDING + assert (await q.counts())["PENDING"] == 1 + assert await w.process_one() is True # 2/2 실패 → DEAD + counts = await q.counts() + assert counts["DEAD"] == 1 and counts["PENDING"] == 0 + assert await w.process_one() is False # DEAD 는 claim 대상 아님 + + +async def test_reaper_reclaims_then_worker_reprocesses(q): + jid = await q.enqueue(JobType.SEARCH.value, {"product_name": "x"}) + await q.claim("dead-worker", lease_sec=1) # 점유 후 사망 흉내 + await asyncio.sleep(1.3) + assert jid in await q.reap() # 회수 → PENDING + + async def handler(job): + return {"ok": True} + + assert await Worker("w2", q, handler).process_one() is True + assert (await q.counts())["DONE"] == 1 + + +async def test_enqueue_notifies_listener(q): + listener = JobListener() + await listener.start() + try: + await q.enqueue(JobType.SEARCH.value, {"product_name": "x"}) # pg_notify 발생 + assert await listener.wait(3.0) is True # 즉시 깨어남 + finally: + await listener.close() diff --git a/lps/worker/handlers.py b/lps/worker/handlers.py new file mode 100644 index 0000000..fda6e46 --- /dev/null +++ b/lps/worker/handlers.py @@ -0,0 +1,35 @@ +"""잡 핸들러 — job_type 별 처리. 현재는 SEARCH(검색)만. + +검색 핸들러는 소스 어댑터로 검색해 정규화 결과를 반환한다. +필터/이상치/AI 유사도(코어 파이프라인)는 다음 단계에서 이 핸들러 안에 결합한다. +""" + +from common.enums import JobType +from services.search.contract import SearchAdapter + + +def build_search_handler(adapters: dict[str, SearchAdapter], default_source: str = "coupang", limit: int = 40): + """검색 핸들러 생성. adapters = {source: SearchAdapter}.""" + + async def handler(job: dict) -> dict: + if job["job_type"] != JobType.SEARCH.value: + raise ValueError(f"unsupported job_type: {job['job_type']}") + + payload = job.get("payload") or {} + query = (payload.get("product_name") or "").strip() + if not query: + raise ValueError("empty product_name") + + adapter = adapters.get(default_source) + if adapter is None: + raise ValueError(f"no adapter for source: {default_source}") + + products = await adapter.search(query, limit=limit) + return { + "source": adapter.source, + "query": query, + "count": len(products), + "products": [p.model_dump() for p in products], + } + + return handler diff --git a/lps/worker/notify.py b/lps/worker/notify.py new file mode 100644 index 0000000..b90320d --- /dev/null +++ b/lps/worker/notify.py @@ -0,0 +1,47 @@ +"""LISTEN/NOTIFY 리스너 — 잡 적재 시 워커를 즉시 깨운다(폴링 낭비 제거). + +전용 asyncpg 연결로 LISTEN 한다(SQLAlchemy 풀과 분리). 알림이 오면 이벤트를 세팅하고, +워커는 claim 이 비었을 때 wait()로 알림 또는 짧은 타임아웃(안전망/reaper)까지 대기한다. +""" + +import asyncio + +import asyncpg + +from config.server_configs import main_db_config +from crud.job_crud import JOB_NOTIFY_CHANNEL + + +def _dsn() -> str: + c = main_db_config + pw = f":{c.write_pw}" if c.write_pw else "" + return f"postgresql://{c.write_id}{pw}@{c.write_host}:{c.write_port}/{c.name}" + + +class JobListener: + def __init__(self, channel: str = JOB_NOTIFY_CHANNEL): + self._channel = channel + self._conn: asyncpg.Connection | None = None + self._event = asyncio.Event() + + async def start(self): + self._conn = await asyncpg.connect(_dsn()) + await self._conn.add_listener(self._channel, self._on_notify) + + def _on_notify(self, *_args): + self._event.set() + + async def wait(self, timeout: float) -> bool: + """알림이 오거나 timeout 까지 대기. 알림으로 깨면 True, 타임아웃이면 False.""" + try: + await asyncio.wait_for(self._event.wait(), timeout) + return True + except asyncio.TimeoutError: + return False + finally: + self._event.clear() + + async def close(self): + if self._conn is not None: + await self._conn.close() + self._conn = None diff --git a/lps/worker/runner.py b/lps/worker/runner.py new file mode 100644 index 0000000..3ccd8db --- /dev/null +++ b/lps/worker/runner.py @@ -0,0 +1,84 @@ +"""워커 루프 + 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 diff --git a/lps/worker_main.py b/lps/worker_main.py new file mode 100644 index 0000000..e35900f --- /dev/null +++ b/lps/worker_main.py @@ -0,0 +1,52 @@ +# LPS 워커 프로세스 진입점 (API 와 분리 실행 — 코드베이스 공유, 독립 스케일). +# python worker_main.py +# WORKER_CONCURRENCY=2 python worker_main.py +# +# 쿠팡 검색은 브라우저(Chrome)라 무겁고 컨텍스트당 직렬이므로 기본 동시성은 1. +# 여러 브라우저로 늘리려면 CoupangAdapter 인스턴스를 워커마다 따로 준다. + +import asyncio +import os + +from common.logger import LOG +from config.server_configs import web_server_config +from crud.job_crud import JobQueue +from services.search.coupang.adapter import CoupangAdapter +from worker.handlers import build_search_handler +from worker.notify import JobListener +from worker.runner import Worker, run_reaper + +LOG.SetPrefix(f"{web_server_config.server_name}-worker") + + +async def main(concurrency: int = 1): + queue = JobQueue() + adapters = {"coupang": CoupangAdapter(headless=False)} # 워커들이 공유(내부 직렬화) + handler = build_search_handler(adapters) + + stop = asyncio.Event() + listeners: list[JobListener] = [] + tasks: list[asyncio.Task] = [] + + for i in range(concurrency): + listener = JobListener() + await listener.start() + listeners.append(listener) + worker = Worker(f"worker-{i}", queue, handler) + tasks.append(asyncio.create_task(worker.run(listener, stop))) + + tasks.append(asyncio.create_task(run_reaper(queue, stop))) + LOG.i(f"LPS 워커 {concurrency}개 + reaper 기동") + + try: + await asyncio.gather(*tasks) + finally: + stop.set() + for listener in listeners: + await listener.close() + for adapter in adapters.values(): + await adapter.close() + + +if __name__ == "__main__": + asyncio.run(main(int(os.environ.get("WORKER_CONCURRENCY", "1"))))