From d3106efeee46b272058837e1750da7a9acd48c06 Mon Sep 17 00:00:00 2001 From: jze9 Date: Sat, 5 Sep 2026 21:56:48 +0500 Subject: [PATCH] =?UTF-8?q?feat(ops):=20=D0=B7=D0=B0=D0=BB=D0=B8=D0=B2?= =?UTF-8?q?=D0=BA=D0=B0=20=D0=B8=D0=B4=D1=91=D1=82=20=D1=81=D0=B0=D0=BC?= =?UTF-8?q?=D0=B0=20=E2=80=94=20=D0=B0=D0=B2=D1=82=D0=BE=D0=BF=D0=B5=D1=80?= =?UTF-8?q?=D0=B5=D0=B7=D0=B0=D0=BF=D1=83=D1=81=D0=BA=20=D0=BF=D0=BE=D1=80?= =?UTF-8?q?=D1=86=D0=B8=D0=B9=20=D0=BF=D0=BE=20cron?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Массовый источник за прогон берёт порцию, сохраняет позицию и останавливается по бюджету времени. Без внешнего толчка заливка шла рывками — ровно столько, сколько раз кто-то нажмёт «Запустить», и всё это время корпус стоял на месте. keep_ingesting.py раз в 5 минут проверяет, не простаивают ли массовые источники, и запускает следующую порцию. Идущие прогоны не трогает. Останавливается сам, когда источник исчерпан (прогон закрылся с нулём полученных статей), чтобы не крутить пустые запуски вечно. Поставлен в cron на app-хосте рядом с монитором. Co-Authored-By: Claude Opus 5 --- scripts/ops/keep_ingesting.py | 85 +++++++++++++++++++++++++++++++++++ 1 file changed, 85 insertions(+) create mode 100755 scripts/ops/keep_ingesting.py diff --git a/scripts/ops/keep_ingesting.py b/scripts/ops/keep_ingesting.py new file mode 100755 index 0000000..8570274 --- /dev/null +++ b/scripts/ops/keep_ingesting.py @@ -0,0 +1,85 @@ +#!/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, по умолчанию включено). +""" + +import argparse + +BULK_TYPES = ("wikipedia_ru", "pmc_bulk") + + +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()) + + with db_session() as s: + 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()