웹빌더가 세 자리에서 연달아 죽었다 — 가게 등록 · 수집 시작 · 수집 완료. 전부 같은 뿌리다: DB 구조 재편이 테이블 이름을 옮기면서 **참조 세 종류 중 일부만** 따라갔다. - **생성자 12군데** (`place_links(...)` → `place_channels(...)`) import 와 `DBType()` 은 고쳤는데 생성자를 빠뜨렸다. 클래스가 없어도 import 는 통과하므로 기동은 정상이고, 그 줄이 실제로 실행되는 순간에만 터진다. - **raw SQL 12군데** (`job.jobs` → `jobs`) 잡 큐만 raw SQL 이라 ORM 이름 변경에 안 딸려 왔다. 큐가 안 도니 수집·비전·소개문·빌드· 지역데이터가 하나도 못 들어간다. 화면에는 "버튼만 안 먹는" 것으로 보였다. - **같은 이름의 속성 5군데** (`source.place_facts` → `source.facts` 등) 이름만 보고 일괄 치환해 테이블과 무관한 자리까지 바뀌었다. `RawSource` 는 수집기 결과 객체지 테이블이 아니다. - **뗀 표를 계속 부르던 5군데** (`ai_check_results`) 한 번도 쓰지 않아 마이그레이션이 뗀 표다. 부르면 SEO 진단이 통째로 죽는다. ★ 하나씩 터질 때마다 고치다가 멈추고 정적 검사로 남은 것을 한 번에 셌다 — pyflakes 가 19건을 짚었다. 이 종류는 import 도 타입검사도 안 잡는다. 테이블 이름을 옮긴 뒤에는 `python -m pyflakes services/ crud/ router/ worker/ common/ | grep "undefined name"` 을 돌린다. 검증: 직접 수집 실행(스테이,머뭄) — 잡 DONE · 재시도 0 · fact 2 · 사진 10 · 채널 2 저장. 정의 안 된 이름 0건. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
85 lines
3.5 KiB
Python
85 lines
3.5 KiB
Python
# o2o-web4ai 워커 프로세스 진입점 (API 와 분리 실행 — 코드베이스 공유, 독립 스케일).
|
|
# python worker_main.py
|
|
# WORKER_CONCURRENCY=3 python worker_main.py
|
|
#
|
|
# 수집·비전분석·빌드는 몇 분씩 걸려 동기 요청으로 처리할 수 없다. API 는 잡만 적재하고 즉시 응답하며,
|
|
# 실제 처리는 이 프로세스가 한다. 큐는 PostgreSQL(jobs) — 원자적 claim + lease 소유권이라
|
|
# 워커를 몇 개 띄우든(docker compose --scale) 같은 잡이 두 번 돌지 않는다.
|
|
|
|
import asyncio
|
|
import os
|
|
import signal
|
|
|
|
from common.database.db_session_manager import DB_SESSION_MNG
|
|
from common.logger import LOG
|
|
from config.server_configs import web_server_config
|
|
from crud.job_crud import JobQueue
|
|
from worker.handlers import build_handler, log_registry
|
|
from worker.notify import JobListener
|
|
from worker.runner import Worker, run_reaper
|
|
|
|
LOG.SetPrefix(f"{web_server_config.server_name}-worker")
|
|
|
|
# 잡 1건 처리 시간 상한. Perplexity(10~30s) + 크롤링 + Vision(사진 20~50장 배치)을 감안한 값.
|
|
JOB_DEADLINE_SEC = float(os.environ.get("JOB_DEADLINE_SEC", "900"))
|
|
# lease 임대 시간. heartbeat 가 lease_sec/3 마다 갱신하므로 짧아도 되지만,
|
|
# 워커가 죽었을 때 이만큼 지나야 reaper 가 회수한다.
|
|
LEASE_SEC = int(os.environ.get("JOB_LEASE_SEC", "120"))
|
|
|
|
|
|
async def main(concurrency: int = 1):
|
|
queue = JobQueue()
|
|
handler = build_handler()
|
|
log_registry()
|
|
|
|
stop = asyncio.Event()
|
|
listeners: list[JobListener] = []
|
|
tasks: list[asyncio.Task] = []
|
|
|
|
# ── graceful shutdown: SIGINT(Ctrl+C)/SIGTERM(docker stop) → stop 이벤트 ──
|
|
# 하던 잡은 마무리하고 새 잡은 받지 않는다. 강제 종료돼도 lease 만료 후 reaper 가 재큐한다.
|
|
def _request_stop(sig_name: str):
|
|
if not stop.is_set():
|
|
LOG.i(f"{sig_name} 수신 — graceful 종료: 새 잡 중단, 하던 잡 마무리 (한 번 더 = 강제 종료)")
|
|
stop.set()
|
|
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):
|
|
listener = JobListener()
|
|
await listener.start()
|
|
listeners.append(listener)
|
|
worker = Worker(f"worker-{i}", queue, handler, lease_sec=LEASE_SEC, job_deadline_sec=JOB_DEADLINE_SEC)
|
|
tasks.append(asyncio.create_task(worker.run(listener, stop)))
|
|
|
|
tasks.append(asyncio.create_task(run_reaper(queue, stop)))
|
|
LOG.i(f"워커 {concurrency}개 + reaper 기동 (lease {LEASE_SEC}s · 잡 데드라인 {JOB_DEADLINE_SEC:.0f}s)")
|
|
|
|
gathered = asyncio.gather(*tasks)
|
|
try:
|
|
await gathered
|
|
except asyncio.CancelledError:
|
|
LOG.w("강제 종료 — 남은 리소스 정리 후 종료")
|
|
finally:
|
|
stop.set()
|
|
for t in tasks:
|
|
t.cancel()
|
|
await asyncio.gather(gathered, return_exceptions=True)
|
|
for listener in listeners:
|
|
try:
|
|
await listener.close()
|
|
except Exception as ex:
|
|
LOG.e_no_callstack(f"[shutdown] 리스너 정리 실패(무시): {ex}")
|
|
await DB_SESSION_MNG.dispose_all()
|
|
LOG.i("워커 종료 완료")
|
|
|
|
|
|
if __name__ == "__main__":
|
|
asyncio.run(main(int(os.environ.get("WORKER_CONCURRENCY", "1"))))
|