Files
anti-plagiarism/scripts/ops/keep_ingesting.py
jze9 9123844a87
All checks were successful
Deploy / test (push) Successful in 11m54s
Deploy / deploy (push) Successful in 8s
fix(ops): не хоронить источник после одного пустого прогона
Сторож считал источник исчерпанным, если прошлый прогон не добавил и не
выбрал ни одной статьи. Но ровно так же выглядит обрыв связи с API и прогон
на устаревшем коде парсера: сегодня CORE после единственной неудачи был
помечен исчерпанным и больше не запускался — молча, без ошибки в интерфейсе.

Теперь нужно два пустых прогона подряд. Разовый сбой переживём, а реально
кончившийся источник остановится всего на один прогон позже.

Заодно таймаут запроса к CORE снижен с 60 до 30 секунд: три попытки по
минуте отъедали 180 секунд из 1500 бюджета на одной залипшей странице —
та же грабля, что чинили в pmc_bulk. Обычный ответ приходит за 7 секунд.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-09-16 18:34:14 +05:00

161 lines
9.2 KiB
Python
Executable File
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
#!/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", "core")
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:
# Источник исчерпан, если прогоны перестали добавлять статьи: дальше
# в дампе/бакете/выдаче для нас ничего нет.
#
# Но ОДИН пустой прогон — ещё не приговор: ровно так же выглядят
# обрыв связи с API и прогон на устаревшем коде парсера. На CORE
# 16.09.2026 из-за этого источник был помечен исчерпанным после
# единственной неудачи и больше не запускался — молча, без ошибки.
# Поэтому ждём два пустых прогона подряд: разовый сбой переживём,
# а реально кончившийся источник остановится всего на один прогон
# позже.
if last_run_id and not args.keep_going_empty:
last_two = s.execute(text("""
SELECT status, added, fetched FROM parse_runs
WHERE source_id = :sid ORDER BY id DESC LIMIT 2
"""), {"sid": source_id}).all()
if len(last_two) == 2 and all(
r.status == "done" and r.added == 0 and r.fetched == 0 for r in last_two
):
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()