WORKER_CONCURRENCY 를 늘려도 공유 브라우저 lock 때문에 직렬화되던 문제를 고쳐 진짜 병렬 검색.
- worker_main: _build_worker(i) 로 워커마다 자립 세트(브라우저 어댑터·AI·핸들러) 생성.
프로필 분리(user_data_dir_w{i}, ProcessSingleton 충돌 회피) + 워커별 다른 프록시 포트
(proxy.seed_offset 로 100포트를 균등 분할=다른 IP). naver/judge/keyword 도 워커별(공유상태 경합 제거).
프리플라이트는 대표 프록시로 게이트 1회 확인.
- proxy.seed_offset(k): 워커 시작 포트 분산.
- loadtest.py: N개 상품 제출→폴링→처리량·지연(p50/p95)·AI/DECODO/총비용 집계.
실측(동시성2, 4상품): 순차합 323s→벽시계 181s(~1.8x), 상품당 $0.0071, 1000건 ~$7.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
110 lines
5.3 KiB
Python
110 lines
5.3 KiB
Python
# 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
|
||
|
||
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.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")
|
||
|
||
|
||
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()
|
||
suffix = f"_w{i}" if concurrency > 1 else ""
|
||
|
||
def _pf(source): # 워커별 Chrome 프로필 경로(중복 실행 시 ProcessSingleton 충돌 방지)
|
||
return f"/tmp/lps_{source}{suffix}"
|
||
|
||
adapters = {
|
||
"coupang": CoupangAdapter(headless=False, user_data_dir=_pf("coupang"), proxy=proxy, on_detect=bot_log.record),
|
||
"naver": NaverAdapter(), # httpx 직접(프록시 미경유) — 워커별 인스턴스(last_bytes 경합 회피)
|
||
}
|
||
fallback_adapters = {
|
||
"gmarket": EsmAdapter("gmarket", headless=False, user_data_dir=_pf("gmarket"), proxy=proxy, on_detect=bot_log.record),
|
||
"auction": EsmAdapter("auction", headless=False, user_data_dir=_pf("auction"), proxy=proxy, on_detect=bot_log.record),
|
||
"st11": ElevenStAdapter(headless=False, user_data_dir=_pf("st11"), proxy=proxy, on_detect=bot_log.record),
|
||
}
|
||
# 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 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 프리플라이트 실패 — 살아있는 포트를 못 찾음(런타임 회전으로 재시도)")
|
||
LOG.i(f"AI(판정+검색어생성): {'ON' if has_openai else 'OFF(키 없음)'} · 오픈마켓 폴백: gmarket, auction, st11")
|
||
|
||
stop = asyncio.Event()
|
||
listeners: list[JobListener] = []
|
||
tasks: list[asyncio.Task] = []
|
||
all_adapters = []
|
||
|
||
for i in range(concurrency):
|
||
handler, worker_adapters = _build_worker(i, concurrency, has_openai, neg_cache, history)
|
||
all_adapters += worker_adapters
|
||
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 기동 (워커별 브라우저 세트 — 상품 {concurrency}개 동시 검색)")
|
||
|
||
try:
|
||
await asyncio.gather(*tasks)
|
||
finally:
|
||
stop.set()
|
||
for listener in listeners:
|
||
await listener.close()
|
||
for adapter in all_adapters:
|
||
await adapter.close()
|
||
|
||
|
||
if __name__ == "__main__":
|
||
asyncio.run(main(int(os.environ.get("WORKER_CONCURRENCY", "1"))))
|