fix(worker): auto-resume after redeploy + bump concurrency 5->8
Auto-resume bug: kn and nspd_geo resume-on-worker_ready required heartbeat >5min / >10min stale. After redeploy worker boots in ~1-2 min, so jobs killed seconds before deploy had fresh heartbeat and were NEVER auto-resumed — required manual cancel + resume. Fix: drop time threshold entirely. On worker_ready ANY 'running' or 'paused' job is by definition a zombie (no worker exists yet), safe to resume all of them. Concurrency bump: 5 -> 8 prefork slots. Headroom for 5 geo jobs + 1 kn sweep + 2 objective tasks running simultaneously. Each slot ~150MB RSS -> ~1.2GB total, well within 6GB VPS RAM budget.
This commit is contained in:
parent
43672a0068
commit
5c92c2a0e9
2 changed files with 13 additions and 11 deletions
|
|
@ -212,9 +212,14 @@ celery_app.conf.beat_schedule = _build_beat_schedule()
|
||||||
@worker_ready.connect
|
@worker_ready.connect
|
||||||
def _resume_zombie_runs(sender=None, **_kwargs) -> None:
|
def _resume_zombie_runs(sender=None, **_kwargs) -> None:
|
||||||
"""When a worker finishes booting (after redeploy/restart), find any sweep
|
"""When a worker finishes booting (after redeploy/restart), find any sweep
|
||||||
that was 'running' with a stale heartbeat (>5 min) and re-enqueue a resume
|
that was 'running' or 'paused' and re-enqueue a resume task. Each resume
|
||||||
task. Each resume creates a fresh run_id linked via resumed_from_run_id;
|
creates a fresh run_id linked via resumed_from_run_id; the original row is
|
||||||
the original row is marked 'zombie' so the audit trail is preserved.
|
marked 'zombie' so the audit trail is preserved.
|
||||||
|
|
||||||
|
No time threshold: by definition, on worker_ready ANY 'running' row is a
|
||||||
|
zombie because there is no active worker. Previously we required heartbeat
|
||||||
|
>5 min, which skipped jobs interrupted seconds before redeploy — they
|
||||||
|
stayed in 'running' status forever and required manual cancel/resume.
|
||||||
"""
|
"""
|
||||||
from sqlalchemy import text
|
from sqlalchemy import text
|
||||||
|
|
||||||
|
|
@ -230,8 +235,6 @@ def _resume_zombie_runs(sender=None, **_kwargs) -> None:
|
||||||
FROM kn_scrape_runs
|
FROM kn_scrape_runs
|
||||||
WHERE status = 'running'
|
WHERE status = 'running'
|
||||||
AND objects_snapshot IS NOT NULL
|
AND objects_snapshot IS NOT NULL
|
||||||
AND COALESCE(heartbeat_at, started_at)
|
|
||||||
< NOW() - INTERVAL '5 minutes'
|
|
||||||
ORDER BY started_at ASC
|
ORDER BY started_at ASC
|
||||||
LIMIT 20
|
LIMIT 20
|
||||||
"""
|
"""
|
||||||
|
|
@ -286,9 +289,10 @@ def _resume_zombie_runs(sender=None, **_kwargs) -> None:
|
||||||
# старых runs сохранена в nspd_scrape_runs — никаких side-effects.
|
# старых runs сохранена в nspd_scrape_runs — никаких side-effects.
|
||||||
|
|
||||||
# NSPD geo-jobs: bulk-fetcher с собственной resume-логикой через
|
# NSPD geo-jobs: bulk-fetcher с собственной resume-логикой через
|
||||||
# nspd_geo_jobs / nspd_geo_targets. Жадный resume: status='running' со
|
# nspd_geo_jobs / nspd_geo_targets. Resume любых 'running' / 'paused' jobs
|
||||||
# stale heartbeat (>10мин) → re-enqueue с тем же job_id (task сам прочитает
|
# — на worker_ready по определению нет активных воркеров, всё running ==
|
||||||
# pending targets и продолжит).
|
# zombie. Раньше требовали heartbeat >10мин, что пропускало jobs убитых
|
||||||
|
# за минуту до редеплоя и оставляло их вечно висеть.
|
||||||
db = SessionLocal()
|
db = SessionLocal()
|
||||||
geo_resume_jobs: list[int] = []
|
geo_resume_jobs: list[int] = []
|
||||||
try:
|
try:
|
||||||
|
|
@ -300,8 +304,6 @@ def _resume_zombie_runs(sender=None, **_kwargs) -> None:
|
||||||
SET status = 'queued',
|
SET status = 'queued',
|
||||||
error = COALESCE(error, 'auto-resume at worker_ready')
|
error = COALESCE(error, 'auto-resume at worker_ready')
|
||||||
WHERE status IN ('running', 'paused')
|
WHERE status IN ('running', 'paused')
|
||||||
AND COALESCE(heartbeat_at, started_at, created_at)
|
|
||||||
< NOW() - INTERVAL '10 minutes'
|
|
||||||
RETURNING job_id
|
RETURNING job_id
|
||||||
"""
|
"""
|
||||||
)
|
)
|
||||||
|
|
|
||||||
|
|
@ -96,7 +96,7 @@ services:
|
||||||
# контейнера задан settings.objective_anton_sqlite_path (default
|
# контейнера задан settings.objective_anton_sqlite_path (default
|
||||||
# /data/anton-sqlite/analysis.db).
|
# /data/anton-sqlite/analysis.db).
|
||||||
- /opt/gendesign/site-finder:/data/anton-sqlite:ro
|
- /opt/gendesign/site-finder:/data/anton-sqlite:ro
|
||||||
command: ["celery", "-A", "app.workers.celery_app", "worker", "--loglevel=info", "--concurrency=5", "--queues=celery,scrape_kn,geo"]
|
command: ["celery", "-A", "app.workers.celery_app", "worker", "--loglevel=info", "--concurrency=8", "--queues=celery,scrape_kn,geo"]
|
||||||
|
|
||||||
beat:
|
beat:
|
||||||
# Lean backend-образ (без Chromium) — beat только триггерит таски в Redis.
|
# Lean backend-образ (без Chromium) — beat только триггерит таски в Redis.
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue