Массовый источник за прогон берёт порцию, сохраняет позицию и останавливается по бюджету времени. Без внешнего толчка заливка шла рывками — ровно столько, сколько раз кто-то нажмёт «Запустить», и всё это время корпус стоял на месте. keep_ingesting.py раз в 5 минут проверяет, не простаивают ли массовые источники, и запускает следующую порцию. Идущие прогоны не трогает. Останавливается сам, когда источник исчерпан (прогон закрылся с нулём полученных статей), чтобы не крутить пустые запуски вечно. Поставлен в cron на app-хосте рядом с монитором. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
86 lines
4.1 KiB
Python
Executable File
86 lines
4.1 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, по умолчанию включено).
|
||
"""
|
||
|
||
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()
|