# LPS 워커 프로세스 진입점 (API 와 분리 실행 — 코드베이스 공유, 독립 스케일). # python worker_main.py # WORKER_CONCURRENCY=3 python worker_main.py # 상품 3개 동시 검색(권장 2~3, 로컬) # # 브라우저 어댑터는 컨텍스트당 직렬(lock)이라, 진짜 병렬을 위해 **워커마다 자기 브라우저 세트**를 준다: # 프로필 분리(user_data_dir_w{i}) + 워커별 다른 프록시 포트(=다른 IP). 동시성 N → 최대 4×N Chrome. import asyncio import os import signal import time from common.alerts import AlertManager from common.database.db_session_manager import DB_SESSION_MNG from common.logger import LOG from config.server_configs import web_server_config, openai_config, decodo_config from crud.job_crud import JobQueue from crud.negative_cache import NegativeCache from crud.bot_detection import BotDetectionLog from crud.ip_session import IpSessionLog from crud.price_history import PriceHistory from services.search.proxy import DecodoProxy from services.search.coupang.adapter import CoupangAdapter from services.search.naver.adapter import NaverAdapter from services.search.esm.adapter import EsmAdapter from services.search.st11.adapter import ElevenStAdapter from services.ai.similarity import SimilarityJudge from services.ai.keyword import KeywordGenerator 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") # 오픈마켓 폴백(G마켓·옥션·11번가)은 **기본 비활성** — 2026-07-10 협의 결정. # 실측상 크롤 몰이 최종 최저가를 바꾼 적이 없고(0회), 검색당 최대 15s + 프록시 대역폭의 # 대부분을 차지해 로직에서 제외했다(코드·테스트는 유지, 핸들러는 빈 폴백을 정상 처리). # 재가동: LPS_FALLBACKS=gmarket,auction,st11 (일부만도 가능) — 켜기 전 라이브 스모크로 # 셀렉터 드리프트 점검. 배경은 docs/decision-openmarket-crawler.md. _FALLBACK_SOURCES = ("gmarket", "auction", "st11") def _enabled_fallbacks() -> list[str]: names = [s.strip() for s in os.environ.get("LPS_FALLBACKS", "").split(",") if s.strip()] unknown = [n for n in names if n not in _FALLBACK_SOURCES] if unknown: LOG.w(f"LPS_FALLBACKS 무시된 값: {unknown} (가능: {list(_FALLBACK_SOURCES)})") return [n for n in names if n in _FALLBACK_SOURCES] def _build_worker(i: int, concurrency: int, has_openai: bool, neg_cache, history): """워커 1개의 자립 세트(브라우저 어댑터·AI·핸들러)를 만든다. 프로필 분리(user_data_dir_w{i}) + 워커별 다른 프록시 포트(=다른 IP)로 진짜 병렬을 보장한다.""" # 워커별 프록시(다른 포트=다른 IP). 100포트를 워커 수로 균등 분할해 시작점을 벌린다. proxy = DecodoProxy() if proxy.enabled and concurrency > 1: n = proxy.port_end - proxy.port_start + 1 proxy.seed_offset(i * max(1, n // concurrency)) bot_log = BotDetectionLog() ip_log = IpSessionLog() suffix = f"_w{i}" if concurrency > 1 else "" def _pf(source): # 워커별 Chrome 프로필 경로(중복 실행 시 ProcessSingleton 충돌 방지) # LPS_PROFILE_DIR 를 영속 볼륨으로 마운트하면 재시작해도 cf_clearance 등 쿠키 유지(재웜업 회피). base = os.environ.get("LPS_PROFILE_DIR", "/tmp") return f"{base}/lps_{source}{suffix}" adapters = { "coupang": CoupangAdapter(headless=False, user_data_dir=_pf("coupang"), proxy=proxy, on_detect=bot_log.record, on_session_end=ip_log.record), "naver": NaverAdapter(), # httpx 직접(프록시 미경유) — 워커별 인스턴스(last_bytes 경합 회피) } # 폴백은 기본 비활성(LPS_FALLBACKS 로 켬 — 상단 주석 참고). 켤 땐 데드라인이 상한이라 # 봇감지 재시도(챌린지 대기 2배)를 끈다(max_block_retries=0) — 빠르게 포기·스킵. fallback_adapters = {} for name in _enabled_fallbacks(): if name == "st11": fallback_adapters[name] = ElevenStAdapter(headless=False, user_data_dir=_pf("st11"), proxy=proxy, on_detect=bot_log.record, on_session_end=ip_log.record, max_block_retries=0) else: fallback_adapters[name] = EsmAdapter(name, headless=False, user_data_dir=_pf(name), proxy=proxy, on_detect=bot_log.record, on_session_end=ip_log.record, max_block_retries=0) # AI 도 워커별 인스턴스 — 공유 상태(last_usage) 경합 원천 제거 judge = SimilarityJudge() if has_openai else None keyword_gen = KeywordGenerator() if has_openai else None handler = build_search_handler( adapters, judge=judge, keyword_gen=keyword_gen, neg_cache=neg_cache, history=history, fallback_adapters=fallback_adapters, ai_model=openai_config.model, proxy_cost_per_gb=decodo_config.cost_per_gb, ) return handler, list(adapters.values()) + list(fallback_adapters.values()) async def _warmup_worker(worker_adapters, tries: int = 3, attempt_timeout: float = 60.0): """워커의 챌린지 소스(Turnstile/Akamai)를 미리 풀어 쿠키(cf_clearance 등)를 확보한다. 콜드 비용을 시작 시 몰아, 이후 실 작업은 웜(빠름). 백그라운드로 돌려 잡 처리를 막지 않는다. 나쁜 IP 는 인터랙티브 Turnstile 로 에스컬레이션되므로, 실패 시 **다른 IP 로 회전 재시도**한다. 시도당 타임아웃 필수 — 웜업은 search 중 어댑터 락을 쥐므로, 여기서 행하면 그 워커의 모든 실 검색이 락 대기로 함께 멈춘다(2026-07-10 부하테스트에서 15분 행 실측).""" for ad in worker_adapters: if ad.source not in ("gmarket", "auction", "coupang"): continue for attempt in range(tries): try: await asyncio.wait_for(ad.search("생수", limit=1), timeout=attempt_timeout) LOG.i(f"[warmup:{ad.source}] 챌린지 통과·쿠키 확보 (시도 {attempt + 1})") break except Exception as ex: if attempt < tries - 1: ad._rotate_ip(f"웜업 재시도({type(ex).__name__}) — 새 IP") else: LOG.w(f"[warmup:{ad.source}] {tries}회 실패(첫 잡에서 재시도): {type(ex).__name__}") def _proxy_ports_snapshot(adapters) -> tuple[int, int] | None: """워커 프록시들의 가용 포트 현황 — (최소 가용 수, 전체 포트 수). 프록시 미사용이면 None. 쿨다운 맵은 프록시 인스턴스(워커)별이라 가장 소진된 워커 기준(min)으로 본다.""" proxies = {} for ad in (adapters or []): p = getattr(ad, "_proxy", None) if p is not None and p.enabled: proxies[id(p)] = p if not proxies: return None any_p = next(iter(proxies.values())) total = any_p.port_end - any_p.port_start + 1 return min(p.available_ports() for p in proxies.values()), total async def run_ops_monitor(queue, bot_log, stop, interval: float = 30.0, adapters=None, alerts=None, ip_log=None): """워커 헬스 하트비트 + 임계 알림. 주기적으로 (1) 하트비트 파일 갱신(Docker HEALTHCHECK 가 행/좀비 워커 감지) (2) 큐/차단/DB풀/소스별 실패 지표 점검 → AlertManager 로 발화 (룰별 쿨다운으로 스팸 방지, 조건 해소 시 회복 알림).""" hb_path = os.environ.get("LPS_HEARTBEAT_FILE", "/tmp/lps_worker_heartbeat") th_dead = int(os.environ.get("LPS_ALERT_DEAD_1H", "20")) th_blocks = int(os.environ.get("LPS_ALERT_BLOCKS_1H", "80")) th_lag = int(os.environ.get("LPS_ALERT_QUEUE_LAG_SEC", "300")) th_pool = int(os.environ.get("LPS_ALERT_POOL_PCT", "90")) th_srcfail = int(os.environ.get("LPS_ALERT_SOURCE_FAIL_30M", "5")) th_deadline = int(os.environ.get("LPS_ALERT_DEADLINE_1H", "5")) th_cost = float(os.environ.get("LPS_ALERT_COST_1H_USD", "1.0")) th_ports = int(os.environ.get("LPS_ALERT_PORTS_LOW_PCT", "30")) th_leak = int(os.environ.get("LPS_ALERT_BLOCK_SESSIONS_6H", "1")) alerts = alerts or AlertManager(origin="worker") ip_log = ip_log or IpSessionLog() while not stop.is_set(): try: with open(hb_path, "w") as f: f.write(str(int(time.time()))) # 하트비트(mtime) — HEALTHCHECK 가 신선도 확인 except Exception: pass try: snap = await queue.ops() snap["blocks_1h"] = await bot_log.recent_count(60) pool = DB_SESSION_MNG.pool_status() snap["pool_pct"] = pool["pct"] await alerts.check("dead", snap["dead_1h"] >= th_dead, f"DEAD 1h={snap['dead_1h']}", snap) await alerts.check("blocks", snap["blocks_1h"] >= th_blocks, f"차단 1h={snap['blocks_1h']}", snap) await alerts.check("queue_lag", snap["oldest_pending_sec"] >= th_lag, f"큐지연={snap['oldest_pending_sec']}s", snap) await alerts.check("stuck", snap["stuck_running"] > 0, f"stuck={snap['stuck_running']}", snap) await alerts.check("db_pool", pool["pct"] >= th_pool, f"DB 풀 포화 {pool['pct']}% (checked_out {pool['checked_out']}/{pool['capacity']})", snap) await alerts.check("deadline", snap["deadline_1h"] >= th_deadline, f"잡 데드라인 강제종료 1h={snap['deadline_1h']} — 크롤 행 반복 신호", snap) await alerts.check("cost", snap["cost_1h_usd"] >= th_cost, f"검색원가 1h=${snap['cost_1h_usd']} — 비용 폭주(리소스차단 풀림·재시도 루프) 점검", snap) # 가용 프록시 포트 고갈 — 쿨다운 격리 누적. blocks_1h 보다 먼저 우는 대규모 차단 조기 신호. ports = _proxy_ports_snapshot(adapters) if ports: avail, total = ports snap["proxy_ports_avail"], snap["proxy_ports_total"] = avail, total await alerts.check("proxy_ports_low", avail * 100 <= total * th_ports, f"가용 프록시 포트 {avail}/{total} — 대규모 차단 진행 신호", snap) # 예산 누수 — 요청 예산을 지켰는데도 차단된 IP 세션 발생 = 현재 예산이 안전하지 않다는 신호. block_sessions = (await ip_log.recent_stats(360)).get("block", 0) snap["block_sessions_6h"] = block_sessions await alerts.check("budget_leak", block_sessions >= th_leak, f"예산 회전에도 차단된 IP 세션 6h={block_sessions} — LPS_IP_REQUEST_BUDGET 하향 검토", snap) # 소스별 장기 실패 — 최근 30분간 시도는 있는데 성공이 0건(쿼터 소진·셀렉터 드리프트·전면 차단 신호) per_source: dict[str, list[int]] = {} for ad in (adapters or []): tries, ok = ad.recent_stats(1800) agg = per_source.setdefault(ad.source, [0, 0]) agg[0] += tries agg[1] += ok for src, (tries, ok) in per_source.items(): await alerts.check(f"source_fail:{src}", tries >= th_srcfail and ok == 0, f"{src} 최근 30분 {tries}회 시도·성공 0건", snap) except Exception as ex: LOG.e_no_callstack(f"[ops-monitor] {type(ex).__name__}: {ex}") try: await asyncio.wait_for(stop.wait(), timeout=interval) except asyncio.TimeoutError: pass async def run_browser_reaper(adapters, stop, idle_sec: float = 120.0, interval: float = 30.0, close_timeout: float = 60.0): """유휴 브라우저 정리 루프 — 일정 시간 검색 없는 어댑터의 Chrome 을 닫아 메모리를 회수한다. 쿠키는 user_data_dir 에 남아, 다음 검색 때 재기동해도 (같은 IP면) 웜 유지. 순차 순회라 close 1건에도 타임아웃을 건다 — 한 어댑터의 close 행이 루프 전체를 멈춰 다른 워커의 브라우저까지 못 닫게 되는 것을 실측(2026-07-10 부하테스트)했다.""" while not stop.is_set(): try: await asyncio.wait_for(stop.wait(), timeout=interval) except asyncio.TimeoutError: pass for ad in adapters: close_if_idle = getattr(ad, "close_if_idle", None) if close_if_idle is None: # 네이버(httpx) 등 브라우저 없는 어댑터는 정리 대상 아님 continue try: await asyncio.wait_for(close_if_idle(idle_sec), timeout=close_timeout) except asyncio.TimeoutError: LOG.w(f"[browser-reaper] {getattr(ad, 'source', '?')} 정리 {close_timeout:.0f}s 초과 — 취소·스킵(다음 틱 재시도)") except Exception as ex: LOG.e_no_callstack(f"[browser-reaper] 정리 실패(무시): {ex}") async def main(concurrency: int = 1): queue = JobQueue() neg_cache, history = NegativeCache(), PriceHistory() # DB 기반 — 워커 공유 안전 has_openai = bool(openai_config.api_key) # 시작 프리플라이트: DECODO 게이트가 살아있는지(인증) 대표 프록시로 1회 확인. 포트는 워커별로 각자 잡음. probe = DecodoProxy() LOG.i(f"DECODO 프록시: {'ON(sticky ' + str(probe.session_minutes) + '분 회전)' if probe.enabled else 'OFF(미설정)'}") if probe.enabled: egress_ip, egress_port = await probe.healthcheck() LOG.i(f"DECODO 프리플라이트 OK — egress IP {egress_ip} (port {egress_port})") if egress_ip \ else LOG.w("DECODO 프리플라이트 실패 — 살아있는 포트를 못 찾음(런타임 회전으로 재시도)") fb = _enabled_fallbacks() LOG.i(f"AI(판정+검색어생성): {'ON' if has_openai else 'OFF(키 없음)'} · " f"오픈마켓 폴백: {', '.join(fb) if fb else 'OFF(기본 — LPS_FALLBACKS 로 활성화)'}") stop = asyncio.Event() listeners: list[JobListener] = [] tasks: list[asyncio.Task] = [] bg_tasks: list[asyncio.Task] = [] # 웜업 등 백그라운드(짧게 끝남, gather 대상 아님) all_adapters = [] # ── graceful shutdown: SIGINT(Ctrl+C)/SIGTERM(docker stop) → stop 이벤트 ── # asyncio.run 기본 동작(SIGINT=메인 태스크 즉시 cancel)은 하던 잡을 도중에 끊어 # RUNNING 인 채 lease 만료(120s)까지 묶어둔다. 대신 stop 을 set 해 "새 잡은 안 받고, # 하던 잡은 마무리"로 종료한다. 같은 신호를 한 번 더 받으면 강제 종료(태스크 취소). def _request_stop(sig_name: str): if not stop.is_set(): LOG.i(f"{sig_name} 수신 — graceful 종료: 새 잡 중단, 하던 잡 마무리 (한 번 더 = 강제 종료)") stop.set() for t in bg_tasks: # 웜업은 선택 작업 — 즉시 취소해 어댑터 락을 비운다 t.cancel() else: LOG.w(f"{sig_name} 재수신 — 강제 종료(실행 중 잡은 lease 만료 후 reaper 가 재큐)") for t in tasks: t.cancel() loop = asyncio.get_running_loop() for sig in (signal.SIGINT, signal.SIGTERM): loop.add_signal_handler(sig, _request_stop, sig.name) # 잡 1건 데드라인 — 정상 검색은 폴백 포함 수분 내 끝난다(실측 15~22s). 크롤 행 실측(15분) 대비 상한. job_deadline = float(os.environ.get("LPS_JOB_DEADLINE_SEC", "300")) for i in range(concurrency): handler, worker_adapters = _build_worker(i, concurrency, has_openai, neg_cache, history) all_adapters += worker_adapters bg_tasks.append(asyncio.create_task(_warmup_worker(worker_adapters))) # 챌린지 쿠키 선점(백그라운드) listener = JobListener() await listener.start() listeners.append(listener) worker = Worker(f"worker-{i}", queue, handler, job_deadline_sec=job_deadline) tasks.append(asyncio.create_task(worker.run(listener, stop))) tasks.append(asyncio.create_task(run_reaper(queue, stop))) tasks.append(asyncio.create_task(run_browser_reaper(all_adapters, stop))) # 유휴 브라우저 정리 tasks.append(asyncio.create_task(run_ops_monitor(queue, BotDetectionLog(), stop, adapters=all_adapters))) # 하트비트 + 임계 알림 LOG.i(f"LPS 워커 {concurrency}개 + reaper + 브라우저정리 + ops모니터(하트비트/알림) 기동 (워커별 세트 · 상품 {concurrency}개 동시)") # 종료 유예: stop 후 하던 잡이 이 시간 안에 끝나면 자연 종료, 초과하면 강제 취소. # docker stop 을 쓰면 compose 의 stop_grace_period 를 이보다 길게 잡아야 SIGKILL 전에 마무리된다. grace = float(os.environ.get("LPS_SHUTDOWN_GRACE_SEC", "60")) gathered = asyncio.gather(*tasks) stop_waiter = asyncio.create_task(stop.wait()) try: await asyncio.wait({gathered, stop_waiter}, return_when=asyncio.FIRST_COMPLETED) if gathered.done(): gathered.result() # 워커/리퍼가 예외로 죽은 경우 → 전파(finally 가 정리 후 종료) else: # 종료 신호 경로 — 워커 루프들이 stop 을 보고 하던 잡을 마친 뒤 스스로 끝나길 기다린다 try: await asyncio.wait_for(gathered, timeout=grace) LOG.i("graceful 종료 — 모든 워커가 하던 잡을 마무리함") except asyncio.TimeoutError: LOG.w(f"종료 유예 {grace:.0f}s 초과 — 남은 태스크 강제 취소(잡은 lease 만료 후 재큐)") except asyncio.CancelledError: # 신호 재수신(강제 종료)로 태스크가 취소된 경우 LOG.w("강제 종료 — 남은 리소스 정리 후 종료") finally: stop.set() stop_waiter.cancel() for t in (*tasks, *bg_tasks): t.cancel() # 취소 완주를 기다린 뒤 정리 — 실행 중 태스크가 브라우저/커넥션을 쓰는 채로 닫지 않게 await asyncio.gather(gathered, stop_waiter, *bg_tasks, return_exceptions=True) for listener in listeners: try: await listener.close() except Exception as ex: LOG.e_no_callstack(f"[shutdown] 리스너 정리 실패(무시): {ex}") for adapter in all_adapters: # 항목별 격리 — 하나가 실패해도 나머지 Chrome 은 닫는다 try: await adapter.close() except Exception as ex: LOG.e_no_callstack(f"[shutdown] {getattr(adapter, 'source', '?')} 정리 실패(무시): {ex}") LOG.i("LPS 워커 종료 완료") if __name__ == "__main__": asyncio.run(main(int(os.environ.get("WORKER_CONCURRENCY", "1"))))