From 5ceec1e2e3aede6de58b039f57c8d74f2757ab07 Mon Sep 17 00:00:00 2001 From: jze9 Date: Tue, 15 Sep 2026 13:39:36 +0500 Subject: [PATCH] =?UTF-8?q?fix(ops):=20=D1=81=D1=82=D0=BE=D1=80=D0=BE?= =?UTF-8?q?=D0=B6=20=D0=BD=D0=B5=20=D1=81=D0=BD=D0=B8=D0=BC=D0=B0=D0=B5?= =?UTF-8?q?=D1=82=20=D0=B6=D0=B8=D0=B2=D0=BE=D0=B9=20=D0=BF=D1=80=D0=BE?= =?UTF-8?q?=D0=B3=D0=BE=D0=BD=20=D0=BD=D0=B0=20=D0=B4=D0=BE=D0=BB=D0=B3?= =?UTF-8?q?=D0=BE=D0=BC=20COPY?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Молчащий heartbeat сторож трактовал как смерть прогона и снимал его по порогу 10 минут. Но запись большой пачки отпечатков в базу на HDD идёт COPY'ем 6-15 минут без единого тика heartbeat — и сторож убивал вполне живой прогон, тут же запуская дубль по тем же статьям. Корпус часами топтался на месте: из лога видно десятки перезапусков подряд с нулевым приростом. Теперь перед снятием сторож спрашивает воркеров queue.index через inspect().active(), выполняется ли ещё таск прогона. Снимаются только настоящие зомби — те, чьего celery_task_id нет ни на одном воркере. Если хоть один воркер не ответил, живость не проверить и не снимается ничего (fail-safe). Проверено на живом прогоне: его таск попал в набор активных, прогон помечен защищённым. Co-Authored-By: Claude Opus 4.8 --- scripts/ops/keep_ingesting.py | 62 +++++++++++++++++++++++++++++++++++ 1 file changed, 62 insertions(+) diff --git a/scripts/ops/keep_ingesting.py b/scripts/ops/keep_ingesting.py index 8570274..f47f44a 100755 --- a/scripts/ops/keep_ingesting.py +++ b/scripts/ops/keep_ingesting.py @@ -15,11 +15,36 @@ Останавливается сам, когда источник исчерпан: парсер перестаёт отдавать статьи, прогон закрывается с нулём добавленных, и скрипт больше его не трогает (--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: @@ -38,7 +63,44 @@ def main() -> None: 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()