Сторож считал источник исчерпанным, если прошлый прогон не добавил и не выбрал ни одной статьи. Но ровно так же выглядит обрыв связи с API и прогон на устаревшем коде парсера: сегодня CORE после единственной неудачи был помечен исчерпанным и больше не запускался — молча, без ошибки в интерфейсе. Теперь нужно два пустых прогона подряд. Разовый сбой переживём, а реально кончившийся источник остановится всего на один прогон позже. Заодно таймаут запроса к CORE снижен с 60 до 30 секунд: три попытки по минуте отъедали 180 секунд из 1500 бюджета на одной залипшей странице — та же грабля, что чинили в pmc_bulk. Обычный ответ приходит за 7 секунд. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
161 lines
9.2 KiB
Python
Executable File
161 lines
9.2 KiB
Python
Executable File
#!/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()
|