o2o-negosium-original/lps/worker/handlers.py
민헌 9557c1b1af feat(lps): 네이버 어댑터 + 다중 소스 병합 검색
네이버 쇼핑 오픈API 어댑터(크롤링 불필요) + 핸들러를 다중 소스로 확장.
쿠팡(브라우저)+네이버(API)를 동시 검색·병합해 교차 최저가를 뽑는다.

- search/naver/adapter: httpx + 오픈API + 키 로테이션(429/403 순환), .env 키 로드
- search/naver/transform: 순수 변환(태그/엔티티 정리, lprice). 가격비교(catalog) lprice 는
  '여러 판매자 중 최저가'라 최저가 솔루션엔 핵심 → 유지
- handler: asyncio.gather 동시 검색 + 소스별 실패 격리(일부 죽어도 결과) + 전체 실패 시 잡 실패
- worker_main: adapters={coupang, naver}
- tests: 네이버 변환 + 병합/실패격리/전체실패 4건 → 전체 31/31
- 라이브: 쿠팡30+네이버30 병합 top-N 최저가(소스 라벨 포함)

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-09 08:38:15 +09:00

56 lines
2.4 KiB
Python

"""잡 핸들러 — job_type 별 처리. 현재는 SEARCH(검색)만.
검색 핸들러: 여러 소스 어댑터를 동시 검색 → 병합 → 코어 파이프라인(필터·이상치·top-N 최저가).
소스별 실패는 격리한다(한 소스가 죽어도 나머지로 결과 산출). 모든 소스 실패 시에만 잡 실패(재시도).
AI 유사도 판정은 파이프라인 슬롯에 키 준비 시 결합한다.
"""
import asyncio
from common.enums import JobType
from common.logger import LOG
from services.search.contract import SearchAdapter
from services.search.util import parse_price
from services.pipeline.core import run_price_pipeline
def build_search_handler(adapters: dict[str, SearchAdapter], sources: list[str] | None = None, limit: int = 40, top_n: int = 5):
"""검색 핸들러 생성. adapters = {source: SearchAdapter}. sources 미지정 시 전체 사용."""
use = list(sources) if sources else list(adapters.keys())
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")
base_price = parse_price(payload.get("price"))
# 소스 동시 검색 (실패는 예외로 수거해 격리)
results = await asyncio.gather(
*[adapters[s].search(query, limit=limit) for s in use],
return_exceptions=True,
)
products = []
per_source: dict[str, dict] = {}
for src, res in zip(use, results):
if isinstance(res, Exception):
LOG.w(f"[{src}] 검색 실패: {type(res).__name__}: {res}")
per_source[src] = {"error": f"{type(res).__name__}: {res}"}
else:
products.extend(res)
per_source[src] = {"count": len(res)}
if not products and all("error" in v for v in per_source.values()):
raise RuntimeError(f"모든 소스 검색 실패: {per_source}") # 잡 실패 → 재시도
result = run_price_pipeline(products, base_price=base_price, top_n=top_n)
result["query"] = query
result["sources"] = per_source # 소스별 건수/에러 (관측)
return result
return handler