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>
This commit is contained in:
parent
bebb6a71e6
commit
abcb58ba05
@ -16,6 +16,10 @@ from common.database.db_session_manager import DB_SESSION_MNG
|
|||||||
from common.enums import DBType, DBWRType, JobStatus
|
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:
|
def compute_backoff(attempts: int, base: float = 5.0, cap: float = 600.0) -> float:
|
||||||
"""지수 백오프(초). attempts 회 시도 후 다음 재시도까지 대기 = base * 2^(attempts-1), cap 상한."""
|
"""지수 백오프(초). attempts 회 시도 후 다음 재시도까지 대기 = base * 2^(attempts-1), cap 상한."""
|
||||||
return min(cap, base * (2 ** max(0, attempts - 1)))
|
return min(cap, base * (2 ** max(0, attempts - 1)))
|
||||||
@ -53,6 +57,9 @@ class JobQueue:
|
|||||||
"t": job_type, "p": priority, "payload": json.dumps(payload),
|
"t": job_type, "p": priority, "payload": json.dumps(payload),
|
||||||
"dk": dedupe_key, "ma": max_attempts,
|
"dk": dedupe_key, "ma": max_attempts,
|
||||||
})).first()
|
})).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 str(row[0]) if row else None
|
||||||
|
|
||||||
return await self._tx(run)
|
return await self._tx(run)
|
||||||
|
|||||||
69
lps/tests/test_worker.py
Normal file
69
lps/tests/test_worker.py
Normal file
@ -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()
|
||||||
35
lps/worker/handlers.py
Normal file
35
lps/worker/handlers.py
Normal file
@ -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
|
||||||
47
lps/worker/notify.py
Normal file
47
lps/worker/notify.py
Normal file
@ -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
|
||||||
84
lps/worker/runner.py
Normal file
84
lps/worker/runner.py
Normal file
@ -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
|
||||||
52
lps/worker_main.py
Normal file
52
lps/worker_main.py
Normal file
@ -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"))))
|
||||||
Loading…
Reference in New Issue
Block a user