diff --git a/docker-compose.yml b/docker-compose.yml index b3829f9..d7d8585 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -159,6 +159,7 @@ services: extra_hosts: - "host.docker.internal:host-gateway" shm_size: "1gb" # Chrome 는 /dev/shm 을 많이 씀 — 부족하면 탭 크래시 + stop_grace_period: 75s # graceful 종료 유예(LPS_SHUTDOWN_GRACE_SEC=60 + 정리 여유) — 기본 10s 면 하던 잡 마무리 전에 SIGKILL restart: unless-stopped logging: driver: json-file diff --git a/lps/docs/operations.md b/lps/docs/operations.md index e2dc610..3e9808f 100644 --- a/lps/docs/operations.md +++ b/lps/docs/operations.md @@ -43,6 +43,11 @@ WORKER_CONCURRENCY=3 python worker_main.py # 동시성 2~3(로컬) > 워커 실행 시 쿠팡 크롤링용 **Chrome 창이 뜹니다**(정상). 기동 로그에 `DECODO 프리플라이트 OK — egress IP ...`, `AI: ON/OFF`가 표시됩니다. > 동시성 N이면 상품 N개가 진짜 병렬 처리됩니다(각 워커가 자기 프로필·프록시 IP 사용). +**워커 종료 (graceful)** +- `Ctrl+C`(SIGINT) 또는 `docker stop`(SIGTERM) 1회 → **새 잡은 안 받고, 하던 잡을 마무리한 뒤** 리스너·브라우저를 정리하고 종료합니다(`LPS 워커 종료 완료` 로그, 트레이스백 없음). +- 유예시간 `LPS_SHUTDOWN_GRACE_SEC`(기본 60s) 안에 안 끝나면 강제 취소되고, 그 잡은 lease 만료(120s) 후 reaper 가 재큐합니다. **한 번 더 신호를 보내면 즉시 강제 종료**입니다. +- Docker 는 compose 의 `stop_grace_period: 75s`(유예 60s + 정리 여유)가 SIGKILL 을 그만큼 미뤄줍니다 — 유예를 늘리면 이 값도 같이 늘리세요. + **부하 테스트** ```bash N=8 python loadtest.py # e2e: 상품 8개 제출→처리량·지연(p50/p95)·AI/DECODO/총비용 집계 (워커 필요) diff --git a/lps/worker_main.py b/lps/worker_main.py index 81f5dfd..6534b6f 100644 --- a/lps/worker_main.py +++ b/lps/worker_main.py @@ -7,6 +7,7 @@ import asyncio import os +import signal import time import httpx @@ -171,6 +172,25 @@ async def main(concurrency: int = 1): 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) + for i in range(concurrency): handler, worker_adapters = _build_worker(i, concurrency, has_openai, neg_cache, history) all_adapters += worker_adapters @@ -186,16 +206,42 @@ async def main(concurrency: int = 1): tasks.append(asyncio.create_task(run_ops_monitor(queue, BotDetectionLog(), stop))) # 하트비트 + 임계 알림 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.gather(*tasks) + 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() - for t in bg_tasks: + 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: - await listener.close() - for adapter in all_adapters: - await adapter.close() + 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__":