#!/usr/bin/env python3 """Держать массовую заливку идущей: перезапускать источники, когда они закончили. Массовый источник за один прогон берёт столько статей, сколько влезает в бюджет времени таска, сохраняет позицию и останавливается. Без внешнего толчка заливка идёт рывками — ровно столько, сколько раз кто-то нажмёт «Запустить». Этот скрипт и есть толчок: раз в несколько минут проверяет, не простаивает ли источник, и запускает следующую порцию. Ставится в cron на app-хосте: */5 * * * * cd /home/user/anti-plagiarism && /usr/bin/docker compose \ -f docker-compose.prod.yml exec -T worker-indexer python - \ < scripts/ops/keep_ingesting.py >> ~/.antiplag-monitor/ingest.log 2>&1 Останавливается сам, когда источник исчерпан: парсер перестаёт отдавать статьи, прогон закрывается с нулём добавленных, и скрипт больше его не трогает (--stop-when-empty, по умолчанию включено). Перед подсчётом «занято ли» снимает зависшие прогоны: если worker-indexer упал или перезапустился, прогон навсегда остаётся в статусе running с замолчавшим heartbeat — без этой чистки сторож видит «занято» бесконечно и не запускает вообще ничего, для всех источников сразу (порог тот же, что и в панели отладки — services/api/app/api/admin.py, stale_cutoff). Молчащий heartbeat сам по себе ещё не значит «мёртв»: пачка COPY на большой базе идёт десятки минут без единого тика. Поэтому прогон снимается, только если его таска нет среди выполняющихся на воркерах queue.index; не ответил хоть один такой воркер — не снимается ничего. Иначе сторож снимал живой медленный прогон и тут же запускал его дубль по тем же статьям. """ import argparse BULK_TYPES = ("wikipedia_ru", "pmc_bulk") STALE_MINUTES = 10 def live_task_ids(celery_app) -> set[str] | None: """ID тасков, выполняющихся на воркерах queue.index; None — кто-то из них не ответил.""" queues = celery_app.control.inspect(timeout=10).active_queues() or {} workers = [w for w, qs in queues.items() if any(q["name"] == "queue.index" for q in qs)] if not workers: return None active = celery_app.control.inspect(destination=workers, timeout=10).active() if not active or set(active) != set(workers): return None return {t["id"] for tasks in active.values() for t in tasks} def main() -> None: ap = argparse.ArgumentParser(description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter) ap.add_argument("--types", default=",".join(BULK_TYPES), help="типы источников через запятую") ap.add_argument("--keep-going-empty", action="store_true", help="запускать даже если прошлый прогон ничего не добавил") args = ap.parse_args() from app.celery_app import celery_app from app.db import db_session from app.models import ParseRun, ParseSource from sqlalchemy import text types = tuple(t.strip() for t in args.types.split(",") if t.strip()) live = live_task_ids(celery_app) with db_session() as s: if live is None: print("воркеры queue.index не ответили — живость прогонов не проверить, не снимаю") reaped = [] else: reaped = s.execute(text(f""" UPDATE parse_runs SET status='error', stage='finished', error='зависший прогон снят автоматически: таск не выполняется, heartbeat молчал дольше {STALE_MINUTES} минут', finished_at=now() WHERE status='running' AND (heartbeat_at IS NULL OR heartbeat_at < now() - interval '{STALE_MINUTES} minutes') AND (celery_task_id IS NULL OR NOT (celery_task_id = ANY(CAST(:live AS text[])))) RETURNING id, source_id """), {"live": sorted(live)}).all() # queued без heartbeat — тот же зомби другого вида: send_task не дошёл # или упал между двумя commit при постановке в очередь (celery_task_id # так и остался пустым), воркер о такой задаче никогда не узнает reaped += s.execute(text(f""" UPDATE parse_runs SET status='error', stage='finished', error='зависший прогон снят автоматически: висел в очереди дольше {STALE_MINUTES} минут', finished_at=now() WHERE status='queued' AND celery_task_id IS NULL AND started_at < now() - interval '{STALE_MINUTES} minutes' RETURNING id, source_id """)).all() if reaped: for source_id in {row.source_id for row in reaped}: s.execute(text(""" UPDATE parse_sources SET last_status='error', last_error='зависший прогон снят автоматически' WHERE id=:sid AND last_status='running' """), {"sid": source_id}) s.commit() print(f"снято зависших прогонов: {len(reaped)} ({', '.join(str(r.id) for r in reaped)})") active = s.execute(text( "SELECT count(*) FROM parse_runs WHERE status IN ('queued','running')" )).scalar_one() if active: print(f"уже идёт прогонов: {active} — ждём") return sources = s.execute(text(""" SELECT id, source_type, last_run_id FROM parse_sources WHERE enabled AND source_type = ANY(:types) ORDER BY id """), {"types": list(types)}).all() started = [] for source_id, source_type, last_run_id in sources: # Источник исчерпан, если прошлый прогон не добавил ни одной статьи # и не был отменён: дальше в дампе/бакете для нас ничего нет if last_run_id and not args.keep_going_empty: prev = s.get(ParseRun, last_run_id) if prev and prev.status == "done" and prev.added == 0 and prev.fetched == 0: print(f"{source_type}: источник исчерпан, пропускаем") continue src = s.get(ParseSource, source_id) run = ParseRun(source_id=source_id, status="queued", stage="queued", target=src.limit, log=[]) s.add(run) s.flush() src.last_status = "running" src.last_error = None src.last_run_id = run.id s.commit() result = celery_app.send_task("index.run_parser", args=[source_id, run.id], queue="queue.index") run.celery_task_id = result.id s.commit() started.append(f"{source_type}(прогон {run.id})") print("запущено:", ", ".join(started) if started else "нечего запускать") if __name__ == "__main__": main()