#!/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()