fix(lps): 워커 graceful shutdown — 신호 핸들러·잡 마무리 유예·정리 격리
- SIGINT/SIGTERM 핸들러 등록: cancel 대신 stop 이벤트 set → 새 잡 클레임 중단, 하던 잡은 마무리 후 자연 종료(트레이스백 없이 exit 0). 신호 재수신 시 강제 종료 - 종료 유예 LPS_SHUTDOWN_GRACE_SEC(기본 60s) 초과 시 강제 취소(잡은 lease 만료 후 재큐) - 워커/리퍼가 예외로 죽으면 기존처럼 전파하되, finally에서 남은 태스크 취소·완주 대기 후 정리 - 리스너·어댑터 정리를 항목별 try/except로 격리 — 하나 실패해도 나머지 Chrome 정리 - 웜업(bg) 태스크는 종료 신호 즉시 취소해 어댑터 락 해제 - compose lps-worker에 stop_grace_period: 75s (기본 10s면 드레인 전 SIGKILL) - 운영 가이드에 워커 종료 절차 문서화 검증: SIGTERM/SIGINT 실기동 테스트 — graceful 로그 후 exit 0, 기존 테스트 26개 통과 Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
parent
ddb6f6bd2f
commit
47c64eaf8b
@ -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
|
||||
|
||||
@ -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/총비용 집계 (워커 필요)
|
||||
|
||||
@ -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:
|
||||
try:
|
||||
await listener.close()
|
||||
for adapter in all_adapters:
|
||||
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__":
|
||||
|
||||
Loading…
Reference in New Issue
Block a user