Files
anti-plagiarism/services/worker-indexer/app/tasks/index.py
jze9 d86bb606c7
All checks were successful
Deploy / test (push) Successful in 3m8s
Deploy / deploy (push) Successful in 30s
fix(wikipedia): позиция — байтовое смещение, а не номер статьи
Заливка Википедии останавливалась сама собой: позиция хранилась как номер
статьи, и каждый прогон перечитывал дамп с начала. На 20 тысячах это стоило
6 минут из 25 доступных, на 27 тысячах — уже около десяти, а на сотне тысяч
съело бы весь бюджет и заливка встала бы совсем. Ровно это и наблюдалось:
прогон висел с нулём полученных статей, воркер на 100% CPU.

Переход на multistream-вариант дампа: он состоит из независимых bz2-блоков по
~95 статей, к нему прилагается индекс со смещениями. Позиция продолжения стала
байтовым смещением, прогон стартует мгновенно с нужного места.

Проверено на реальных файлах, а не по предположению: формат индекса
(offset:page_id:title), 21 уникальное смещение на 2001 статью, прыжок seek на
смещение из середины файла даёт валидный XML со страницами.

wikipedia_resume_offset.py — разовый пересчёт позиции при переходе: находит по
индексу блок с максимальным залитым page_id (получилось 273 281 821).

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

988 lines
43 KiB
Python
Raw 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.
"""Celery задачи индексации документов и проверки плагиата (уровни 1-2)."""
import contextlib
import io
from datetime import UTC, datetime
from pathlib import Path
from typing import Any
from celery.utils.log import get_task_logger
from sqlalchemy import func, select, update
from sqlalchemy.exc import IntegrityError
from app.algorithms.minhash import add_to_lsh, find_similar
from app.algorithms.winnowing import sample_evenly, winnow, winnow_ordered
from app.celery_app import celery_app
from app.config import settings
from app.db import db_session, get_minio, refund_plagiarism_quota, update_task_status
from app.extractors.docx import extract_text_from_docx, extract_text_from_txt
from app.extractors.pdf import extract_text_from_pdf
from app.fragments import split_into_fragments
from app.staging import staged_work_to_doc_data
logger = get_task_logger(__name__)
@celery_app.task(
name="index.extract_and_check",
bind=True,
max_retries=3,
default_retry_delay=60,
)
def extract_and_check(
self,
task_id: str,
minio_key: str,
filename: str,
) -> dict[str, Any]:
"""
Извлечь текст из документа и проверить плагиат (уровни 1-2).
Алгоритм:
1. Скачать файл из MinIO
2. Извлечь текст (PDF/DOCX/TXT)
3. Разбить на фрагменты по 200 слов с перекрытием 50 слов
4. Уровень 1: Winnowing против базы fingerprints в PostgreSQL
5. Уровень 2: MinHash LSH нечёткий поиск
6. Диспатч gpu.check_plagiarism для уровней 3 и 4
Args:
task_id: ID задачи в PostgreSQL
minio_key: Ключ объекта в MinIO
filename: Оригинальное имя файла
"""
logger.info(f"Извлечение текста для задачи {task_id!r}, файл: {filename!r}")
update_task_status(task_id, "processing")
try:
# Скачать из MinIO
minio = get_minio()
response = minio.get_object(settings.MINIO_BUCKET_DOCS, minio_key)
file_data = response.read()
response.close()
response.release_conn()
logger.info(f"Файл скачан из MinIO: {minio_key} ({len(file_data)} байт)")
# Извлечь текст в зависимости от формата
ext = Path(filename).suffix.lower()
if ext == ".pdf":
text = extract_text_from_pdf(file_data)
elif ext == ".docx":
text = extract_text_from_docx(file_data)
else:
text = extract_text_from_txt(file_data)
if not text.strip():
raise ValueError("Не удалось извлечь текст из документа")
word_count = len(text.split())
logger.info(f"Текст извлечён: {len(text)} символов, {word_count} слов")
# ──── Автосбор в отстойник (буфер для пополнения базы) ────────────────
# Сохраняем извлечённый текст и метаданные для последующего ручного
# одобрения админом. Не дублируем в базу источников до решения.
_stage_work(task_id, minio_key, filename, text, word_count)
# Разбить на фрагменты
fragments = split_into_fragments(
text,
window=settings.FRAGMENT_WINDOW_WORDS,
overlap=settings.FRAGMENT_OVERLAP_WORDS,
)
logger.info(f"Фрагментов создано: {len(fragments)}")
# ──── Уровень 1: Winnowing fingerprints (пофрагментно, с локализацией) ───
# Winnow'им КАЖДЫЙ фрагмент отдельно и ищем, с каким источником и на
# какой позиции он совпадает. Так в отчёте видно не «документ похож на X»,
# а «фрагмент на символах A–B скопирован из источника X».
level1_matches: list[dict] = []
from app.models import Document, Fingerprint
with db_session() as session:
doc_cache: dict[int, Document] = {}
for frag in fragments:
frag_fp = winnow(frag["text"])
if not frag_fp:
continue
frag_hashes = list(frag_fp)
# Источник, разделяющий больше всего отпечатков с этим фрагментом.
# user_submission исключены: это чужие непубличные загрузки (и
# свои же прошлые прогоны того же файла) — сравнение с ними даёт
# ложные 100%-совпадения, а не реальный плагиат из источника.
row = session.execute(
select(
Fingerprint.doc_id,
func.count(Fingerprint.id).label("cnt"),
)
.join(Document, Document.id == Fingerprint.doc_id)
.where(
Fingerprint.hash_value.in_(frag_hashes),
Document.source != "user_submission",
)
.group_by(Fingerprint.doc_id)
.order_by(func.count(Fingerprint.id).desc())
.limit(1)
).first()
if row is None:
continue
doc_id, cnt = row
# Доля отпечатков фрагмента, найденных в источнике
similarity = cnt / len(frag_hashes) * 100
if similarity < settings.EXACT_FRAGMENT_THRESHOLD:
continue
doc = doc_cache.get(doc_id)
if doc is None:
doc = session.get(Document, doc_id)
if doc is None:
continue
doc_cache[doc_id] = doc
level1_matches.append({
"fragment": frag["text"][:300],
"position_start": frag["start"],
"position_end": frag["end"],
"similarity": round(min(similarity, 100.0), 1),
"method": "exact",
"source_title": doc.title,
"source_url": doc.url,
"source_db": doc.source,
})
logger.info(
f"Уровень 1 (Winnowing): {len(level1_matches)} совпадений-фрагментов"
)
# ──── Уровень 2: MinHash LSH ────────────────────────────────────────────
level2_matches: list[dict] = []
similar_keys = find_similar(text)
if similar_keys:
from app.models import Document
with db_session() as session:
for key in similar_keys[:10]:
# Ключ формата "doc:{id}"
try:
doc_id = int(key.split(":")[-1])
doc = session.get(Document, doc_id)
if not doc or doc.source == "user_submission":
continue
level2_matches.append({
"fragment": text[:200],
"position_start": 0,
"position_end": len(text),
"similarity": 60.0, # MinHash даёт только факт похожести
"method": "fuzzy",
"source_title": doc.title,
"source_url": doc.url,
"source_db": doc.source,
})
except (ValueError, IndexError):
continue
logger.info(f"Уровень 2 (MinHash): {len(level2_matches)} совпадений")
# ──── Диспатч GPU задачи (уровни 3-4) ─────────────────────────────────
# Передаём только текстовые данные (JSON-сериализуемые)
celery_app.send_task(
"gpu.check_plagiarism",
args=[task_id, text, fragments],
kwargs={
"level1_matches": level1_matches,
"level2_matches": level2_matches,
},
queue="queue.gpu",
)
logger.info(f"GPU задача отправлена для задачи {task_id!r}")
return {
"task_id": task_id,
"fragments": len(fragments),
"level1": len(level1_matches),
"level2": len(level2_matches),
}
except ValueError as exc:
# Детерминированная ошибка (битый файл, не тот формат, пустой текст) —
# повтор не поможет: результат будет тем же на 2-й и 3-й попытке. Не
# ретраим (быстрее фидбек юзеру) и возвращаем месячную квоту — реальная
# проверка так и не началась, списывать не за что.
logger.warning(f"Задача {task_id!r}: не удалось обработать файл: {exc}")
update_task_status(task_id, "failed", str(exc))
refund_plagiarism_quota(task_id)
return {"task_id": task_id, "status": "failed", "error": str(exc)}
except Exception as exc:
logger.error(f"Ошибка при обработке задачи {task_id!r}: {exc}", exc_info=True)
update_task_status(task_id, "failed", str(exc))
raise self.retry(exc=exc, countdown=60) from exc
@celery_app.task(name="index.add_document")
def add_document(doc_data: dict[str, Any], dispatch_embed: bool = True) -> dict[str, Any]:
"""
Добавить документ из внешнего источника в систему.
Алгоритм:
1. Дедупликация по ext_id
2. Сохранить метаданные в PostgreSQL
3. Вычислить провизорные Winnowing fingerprints (из аннотации)
4. Добавить в MinHash LSH
5. Индексировать в Elasticsearch
6. Диспатч gpu.embed_documents для FAISS эмбеддингов (если dispatch_embed)
7. Если включён FETCH_FULL_TEXT и есть url — диспатч index.enrich_full_text,
который скачает полный текст и пересчитает fingerprints по нему.
Args:
doc_data: Словарь с метаданными документа
dispatch_embed: Диспатчить ли эмбеддинг по одному документу. При массовой
заливке run_parser выключает это и батчит эмбеддинги сам.
Returns:
dict со статусом операции и doc_id
"""
from app.models import Document, Fingerprint
ext_id = doc_data.get("ext_id")
if not ext_id:
return {"status": "error", "reason": "ext_id обязателен"}
with db_session() as session:
# Проверить дублирование
existing = session.execute(
select(Document).where(Document.ext_id == ext_id)
).scalar_one_or_none()
if existing:
return {"status": "duplicate", "doc_id": existing.id}
# Создать документ
allowed_fields = {c.key for c in Document.__table__.columns}
doc_kwargs = {k: v for k, v in doc_data.items() if k in allowed_fields}
doc = Document(**doc_kwargs)
session.add(doc)
try:
session.flush()
except IntegrityError:
# Гонка: другой воркер параллельно вставил этот ext_id между нашей
# проверкой existing и flush. Уникальный индекс отработал — это дубль.
session.rollback()
existing = session.execute(
select(Document).where(Document.ext_id == ext_id)
).scalar_one_or_none()
return {"status": "duplicate", "doc_id": existing.id if existing else None}
doc_id = doc.id
# Вычислить fingerprints
text = doc_data.get("full_text") or doc_data.get("abstract", "") or ""
if text:
fingerprints_to_add = sample_evenly(
winnow_ordered(text), settings.MAX_FINGERPRINTS_PER_DOC
)
for i, hash_val in enumerate(fingerprints_to_add):
session.add(Fingerprint(doc_id=doc_id, hash_value=hash_val, position=i))
# Добавить в MinHash LSH (in-memory)
add_to_lsh(f"doc:{doc_id}", text)
session.commit()
logger.info(f"Документ {doc_id} добавлен в PostgreSQL: {doc_data.get('title', '')[:50]!r}")
# Индексация в Elasticsearch
try:
from elasticsearch import Elasticsearch
from app.config import settings as cfg
es = Elasticsearch(cfg.ELASTICSEARCH_URL)
es_doc = {
"doc_id": doc_id,
"source": doc_data.get("source"),
"title": doc_data.get("title", ""),
"abstract": doc_data.get("abstract", ""),
"authors": " ".join(
f"{a.get('last_name', '')} {a.get('first_name', '')}"
for a in doc_data.get("authors", [])
),
"year": doc_data.get("year"),
"lang": doc_data.get("lang"),
"journal": doc_data.get("journal"),
"doi": doc_data.get("doi"),
"url": doc_data.get("url"),
}
es.index(index="documents", id=str(doc_id), document=es_doc)
logger.info(f"Документ {doc_id} проиндексирован в Elasticsearch")
except Exception as e:
logger.warning(f"Ошибка индексации в ES для документа {doc_id}: {e}")
# Диспатч FAISS эмбеддингов (по одному документу; при массовой заливке
# run_parser выключает это и батчит сам)
if dispatch_embed:
celery_app.send_task(
"gpu.embed_documents",
args=[[doc_id]],
queue="queue.gpu",
)
# Обогащение полным текстом: скачать PDF по url и пересчитать fingerprints.
# Если full_text уже есть (например, OCR-фрагмент из ответа поиска
# КиберЛенинки) — не дублируем запрос: страница статьи там HTML, а не PDF,
# enrich_full_text гарантированно ничего не найдёт, только зря нагрузит
# источник (и получит 503 при массовой заливке).
if settings.FETCH_FULL_TEXT and doc_data.get("url") and not doc_data.get("full_text"):
celery_app.send_task(
"index.enrich_full_text",
args=[doc_id, doc_data["url"]],
queue="queue.index",
)
return {"status": "indexed", "doc_id": doc_id}
def store_full_text(doc_id: int, text: str) -> dict[str, Any]:
"""Сохранить полный текст документа и переиндексировать его по нему.
Общая часть для всех путей получения полного текста: скачанный PDF
(enrich_full_text) и текст, пришедший прямо из API источника (бэкфилл PMC —
scripts/ops/). Провизорные fingerprints, посчитанные по аннотации,
заменяются на посчитанные по полному тексту — ради этого всё и делается:
L1 начинает видеть тело статьи, а не только её краткое описание.
Args:
doc_id: документ в PostgreSQL
text: полный текст статьи
Returns:
dict со статусом, объёмом текста и числом отпечатков
"""
from sqlalchemy import delete
from app.models import Document, Fingerprint
minio = get_minio()
key = f"corpus/{doc_id}.txt"
data = text.encode("utf-8")
minio.put_object(
settings.MINIO_BUCKET_DOCS, key, io.BytesIO(data), length=len(data),
content_type="text/plain; charset=utf-8",
)
hashes = sample_evenly(winnow_ordered(text), settings.MAX_FINGERPRINTS_PER_DOC)
with db_session() as session:
doc = session.get(Document, doc_id)
if doc is None:
return {"status": "doc_gone", "doc_id": doc_id}
doc.minio_key = key
session.execute(delete(Fingerprint).where(Fingerprint.doc_id == doc_id))
for i, hash_val in enumerate(hashes):
session.add(Fingerprint(doc_id=doc_id, hash_value=hash_val, position=i))
session.commit()
# MinHash LSH (L2) тоже должен считаться по полному тексту
add_to_lsh(f"doc:{doc_id}", text)
logger.info(f"store_full_text: doc={doc_id} {len(text)} симв., fingerprints={len(hashes)}")
return {"status": "ok", "doc_id": doc_id, "chars": len(text), "fingerprints": len(hashes)}
@celery_app.task(
name="index.enrich_full_text",
bind=True,
max_retries=2,
default_retry_delay=120,
)
def enrich_full_text(self, doc_id: int, url: str) -> dict[str, Any]:
"""Скачать полный текст статьи и пересчитать по нему fingerprints.
Метаданные и провизорные fingerprints (из аннотации) уже сохранены в
add_document. Здесь мы:
1. Скачиваем PDF по url и извлекаем текст (best-effort).
2. Сохраняем полный текст в MinIO (bucket documents, префикс corpus/).
3. Заменяем fingerprints документа на посчитанные по полному тексту.
4. Обновляем MinHash LSH.
Недоступный/не-PDF источник — не ошибка: возвращаем no_fulltext.
"""
from app.fulltext import fetch_full_text
text = fetch_full_text(url)
if not text:
return {"status": "no_fulltext", "doc_id": doc_id}
try:
return store_full_text(doc_id, text)
except Exception as exc:
logger.error(f"enrich_full_text: не удалось сохранить текст для {doc_id}: {exc}")
raise self.retry(exc=exc, countdown=120) from exc
@celery_app.task(
name="index.auto_approve_submission",
bind=True,
max_retries=2,
default_retry_delay=60,
)
def auto_approve_submission(self, task_id: str) -> dict[str, Any]:
"""Автоматически добавить проверенную работу студента в базу для сравнения.
Вызывается из gpu.check_plagiarism ПОСЛЕ сохранения результата проверки —
так работа не сматчится сама с собой. Идемпотентно: пропускает, если
StagedWork не 'pending' (уже одобрена/отклонена вручную или повторный вызов)
или если AUTO_APPROVE_SUBMISSIONS выключен настройкой — тогда работа так и
остаётся ждать ручного решения в админке.
"""
from app.models import StagedWork
if not settings.AUTO_APPROVE_SUBMISSIONS:
return {"status": "skipped", "reason": "auto-approve выключен настройкой"}
with db_session() as session:
sw = session.execute(
select(StagedWork).where(StagedWork.task_id == task_id)
).scalar_one_or_none()
if sw is None:
return {"status": "skipped", "reason": "StagedWork не найден"}
if sw.status != "pending":
return {"status": "skipped", "reason": f"уже {sw.status}"}
full_text = ""
if sw.text_key:
try:
minio = get_minio()
obj = minio.get_object(settings.MINIO_BUCKET_STAGING, sw.text_key)
full_text = obj.read().decode("utf-8", errors="replace")
except Exception as e:
logger.warning(f"auto_approve_submission: не прочитан текст {task_id!r}: {e}")
return {"status": "error", "reason": str(e)}
doc_data = staged_work_to_doc_data(sw, full_text)
sw.status = "approved"
sw.reviewed_by = None # None = одобрено автоматически, не человеком
sw.reviewed_at = datetime.now(UTC)
session.commit()
result = add_document(doc_data, dispatch_embed=True)
logger.info(f"auto_approve_submission: задача {task_id!r} → {result}")
return result
def _stage_work(
task_id: str,
minio_key: str,
filename: str,
text: str,
word_count: int,
) -> None:
"""Сохранить проверенную работу в отстойник (StagedWork) + текст в MinIO.
Идемпотентно по task_id: повторный вызов (например при requeue) не дублирует.
Ошибки логируются, но не валят основную проверку плагиата.
"""
try:
from app.models import StagedWork
# Сохранить извлечённый текст в bucket staging
text_key = f"text/{task_id}.txt"
minio = get_minio()
bucket = settings.MINIO_BUCKET_STAGING
if not minio.bucket_exists(bucket):
minio.make_bucket(bucket)
data = text.encode("utf-8")
minio.put_object(
bucket, text_key, io.BytesIO(data), length=len(data),
content_type="text/plain; charset=utf-8",
)
with db_session() as session:
existing = session.execute(
select(StagedWork).where(StagedWork.task_id == task_id)
).scalar_one_or_none()
if existing:
return
# Достать user_id из задачи
from app.models import Task
task = session.get(Task, task_id)
user_id = task.user_id if task else None
session.add(StagedWork(
user_id=user_id,
task_id=task_id,
filename=filename,
minio_key=minio_key,
text_key=text_key,
title=Path(filename).stem if filename else None,
word_count=word_count,
status="pending",
))
session.commit()
logger.info(f"Работа задачи {task_id!r} добавлена в отстойник")
except Exception as e:
logger.warning(f"Не удалось добавить работу {task_id!r} в отстойник: {e}")
def _last_pmc_key() -> str | None:
"""Ключ бакета после последней залитой статьи PMC — старт листинга.
Бакет отдаётся в лексикографическом порядке ключей, и ext_id вида
`pmc:PMC10000000` соответствует ключу `PMC10000000.1`. Берём максимальный
залитый и продолжаем после него.
"""
from sqlalchemy import text as sql_text
with db_session() as session:
row = session.execute(sql_text(
"SELECT max(ext_id) FROM documents WHERE ext_id LIKE 'pmc:PMC%'"
)).first()
if not row or not row[0]:
return None
return row[0].split(":", 1)[1] + ".1"
def _parser_for(source_type: str, cfg: dict[str, Any]) -> tuple[Any, dict[str, Any]]:
"""Инстанс парсера и совместимые с его fetch() аргументы по типу источника."""
limit = cfg["limit"]
query = cfg.get("query") or ""
if source_type == "openalex":
from openalex import OpenAlexParser as P
kwargs = {
"query": query,
"limit": limit,
"lang": cfg.get("lang"),
"year_from": cfg.get("year_from"),
"year_to": cfg.get("year_to"),
"open_access_only": settings.INGEST_OPEN_ACCESS_ONLY,
}
elif source_type == "cyberleninka":
from cyberleninka import CyberLeninkaParser as P
kwargs = {"query": query, "limit": limit}
elif source_type == "arxiv":
from arxiv import ArxivParser as P
kwargs = {"query": query, "limit": limit, "year_from": cfg.get("year_from")}
elif source_type == "pmc":
from pmc import PMCParser as P
kwargs = {
"query": query,
"limit": limit,
"year_from": cfg.get("year_from"),
"year_to": cfg.get("year_to"),
}
elif source_type == "wikipedia_ru":
# Массовый источник: позиция — байтовое смещение в multistream-дампе.
# Раньше хранился номер статьи, и каждый прогон перечитывал дамп с
# начала: на 20 тысячах это стоило 6 минут из 25, дальше росло линейно
from wikipedia_ru import WikipediaRuParser as P
kwargs = {
"limit": limit,
"dump_path": query or None, # query = путь к дампу, если задан
"start_offset": int(cfg.get("resume_token") or 0),
}
elif source_type == "pmc_bulk":
from pmc_bulk import PMCBulkParser as P
kwargs = {
"limit": limit,
"resume_token": cfg.get("resume_token"),
# Токена ещё нет (первый прогон после ручных заливок) — начинаем
# после последней уже залитой статьи, иначе листинг часами
# перемалывает существующие как дубли
"start_after": None if cfg.get("resume_token") else _last_pmc_key(),
}
else:
raise ValueError(f"неизвестный тип источника: {source_type}")
return P(), kwargs
def _ingest_one_by_one(parser: Any, raw_docs: Any, prog: Any, run_id: int) -> tuple[bool, bool]:
"""Обычный путь: документ за документом через add_document.
Подходит источникам, отдающим сотни статей: работает ORM-логика с
дедупликацией, индексацией в Elasticsearch и обогащением полным текстом.
Returns:
(отменён, исчерпан бюджет времени)
"""
cancelled = exhausted = False
prog.stage = "index"
prog.count_fetched(len(raw_docs))
prog.log("info", f"получено {len(raw_docs)} документов, индексация")
_write_run(run_id, **prog.snapshot())
# Эмбеддинги диспатчим пачками, а не по одному документу: так worker-gpu
# кодирует батч разом и переписывает FAISS-индекс на диск раз в N добавлений,
# а не на каждый документ.
embed_batch: list[int] = []
batch_size = settings.EMBED_BATCH_SIZE
def _flush_embed() -> None:
if embed_batch:
celery_app.send_task(
"gpu.embed_documents", args=[list(embed_batch)], queue="queue.gpu"
)
embed_batch.clear()
for raw in raw_docs:
if not cancelled and prog.over_budget:
exhausted = True
prog.log("warning", f"бюджет времени {prog.budget_s:.0f}с исчерпан на индексации")
if cancelled or exhausted:
break
try:
doc = parser.transform(raw)
if not (doc and doc.get("title") and doc.get("ext_id")):
prog.count_result("skipped")
continue
result = add_document(doc, dispatch_embed=False)
prog.count_result(result.get("status", "error"))
if result.get("status") == "indexed":
embed_batch.append(result["doc_id"])
if len(embed_batch) >= batch_size:
_flush_embed()
except Exception as e:
prog.count_result("error")
prog.log("warning", f"документ пропущен: {e}")
logger.warning(f"run_parser: ошибка документа: {e}")
if prog.should_flush():
cancelled = _write_run(run_id, **prog.snapshot())
prog.mark_flushed()
if cancelled:
prog.log("warning", "отмена по запросу из админки")
_flush_embed()
return cancelled, exhausted
def _ingest_bulk(
parser: Any, raw_docs: Any, prog: Any, run_id: int, source_id: int
) -> tuple[bool, bool]:
"""Массовый путь: статьи приходят потоком и пишутся пачками через COPY.
Отличия от обычного пути и почему они нужны:
- парсер отдаёт генератор, поэтому счётчик «получено» растёт по ходу дела,
а не известен заранее;
- запись идёт пакетно (bulk_writer), иначе миллионы отпечатков занимают
недели вместо суток;
- позиция продолжения сохраняется в источнике, чтобы следующий прогон
начинал с места остановки, а не перечитывал дамп с начала.
Эмбеддинги здесь не диспатчатся: они считаются заметно медленнее заливки и
делали бы её узким местом. Векторы досчитываются отдельно —
`scripts/ops/reembed_missing.py`.
Returns:
(отменён, исчерпан бюджет времени)
"""
from app.bulk_writer import connect, write_batch
from app.models import ParseSource
cancelled = exhausted = False
prog.stage = "index"
prog.log("info", "массовый источник: запись пачками")
_write_run(run_id, **prog.snapshot())
conn = connect()
batch: list[dict[str, Any]] = []
resume_token: str | None = None
batch_size = settings.BULK_WRITE_BATCH
def save_resume() -> None:
"""Запомнить позицию в источнике — с неё продолжит следующий прогон."""
if resume_token is None:
return
with db_session() as session:
src = session.get(ParseSource, source_id)
if src:
src.resume_token = str(resume_token)
session.commit()
def flush() -> bool:
"""Записать накопленное; True — попросили остановиться."""
nonlocal conn, batch
if not batch:
return False
conn, counts = write_batch(conn, batch, settings.MAX_FINGERPRINTS_PER_DOC)
for _ in range(counts["added"]):
prog.count_result("indexed")
for _ in range(counts["duplicates"]):
prog.count_result("duplicate")
batch = []
save_resume()
stop = _write_run(run_id, **prog.snapshot())
prog.mark_flushed()
return stop
try:
for raw in raw_docs:
if prog.over_budget:
exhausted = True
prog.log("warning", f"бюджет времени {prog.budget_s:.0f}с исчерпан")
break
try:
doc = parser.transform(raw)
except Exception as e:
prog.count_result("error")
logger.warning(f"run_parser bulk: ошибка документа: {e}")
continue
if not (doc and doc.get("ext_id") and doc.get("text")):
prog.count_result("skipped")
continue
# Позиция продолжения: у Википедии — номер статьи в дампе,
# у PMC — токен страницы бакета
resume_token = doc.get("dump_offset") or doc.get("resume_token") or resume_token
batch.append(doc)
prog.count_fetched(prog.fetched + 1)
if len(batch) >= batch_size and flush():
cancelled = True
prog.log("warning", "отмена по запросу из админки")
break
if not cancelled and flush():
cancelled = True
finally:
with contextlib.suppress(Exception):
conn.close()
return cancelled, exhausted
def _write_run(run_id: int, **fields: Any) -> bool:
"""Записать прогресс прогона и вернуть True, если запрошена отмена.
Отмена кооперативная: админка ставит cancel_requested, воркер узнаёт о ней
на ближайшем тике прогресса и останавливается сам. Так прогон завершается
с осмысленным статусом и не оставляет источник висеть в running (что было
бы при жёстком revoke уже начатого таска).
"""
from app.models import ParseRun
with db_session() as session:
row = session.execute(
update(ParseRun)
.where(ParseRun.id == run_id)
.values(heartbeat_at=datetime.now(UTC), **fields)
.returning(ParseRun.cancel_requested)
).first()
return bool(row and row[0])
@celery_app.task(name="index.run_parser", bind=True)
def run_parser(self, source_id: int, run_id: int | None = None) -> dict[str, Any]:
"""Запустить парсинг источника и наполнить базу документов.
Переиспользует парсеры из scripts/parsers и задачу add_document для
каждого полученного документа. Ход заливки пишется в parse_runs (шкала
загрузки и журнал в админке); прогон можно отменить из админки и он сам
закругляется по бюджету времени, не доводя воркер до падения по
consumer_timeout RabbitMQ (см. app/progress.py).
Args:
source_id: ID источника в parse_sources
run_id: ID заранее созданной строки parse_runs (её создаёт админка,
чтобы прогон был виден в очереди ещё до старта). Без него строка
создаётся здесь — для запусков из скриптов.
"""
import sys
from app.models import ParseRun, ParseSource
from app.progress import RunProgress
# Парсеры лежат в /parsers (bind-mount scripts/parsers, см. compose)
if "/parsers" not in sys.path:
sys.path.insert(0, "/parsers")
with db_session() as session:
src = session.get(ParseSource, source_id)
if src is None:
return {"status": "error", "reason": "источник не найден"}
cfg = {
"source_type": src.source_type,
"query": src.query,
"lang": src.lang,
"year_from": src.year_from,
"year_to": src.year_to,
"limit": src.limit,
# Массовые источники продолжают с сохранённой позиции
"resume_token": src.resume_token,
}
src.last_status = "running"
src.last_error = None
src.last_run_at = datetime.now(UTC)
if run_id is None:
run = ParseRun(source_id=source_id, status="running", stage="fetch")
session.add(run)
session.flush()
run_id = run.id
src.last_run_id = run_id
session.commit()
prog = RunProgress(target=cfg["limit"], budget_s=settings.PARSER_TIME_BUDGET_S)
prog.stage = "fetch"
prog.log("info", f"старт: {cfg['source_type']} q={cfg.get('query') or '—'} limit={cfg['limit']}")
_write_run(
run_id,
status="running",
celery_task_id=self.request.id,
error=None,
# Отдельно от started_at (постановка в очередь): при массовом запуске
# между ними часы ожидания, и без этой отметки «длительность прогона»
# в отладке показывала очередь, а не работу
run_started_at=datetime.now(UTC),
**prog.snapshot(),
)
cancelled = False
exhausted = False # остановлены бюджетом времени, а не концом выдачи
error_msg = None
try:
parser, fetch_kwargs = _parser_for(cfg["source_type"], cfg)
def on_fetch_progress(fetched: int) -> bool:
"""Тик прогресса выборки; False — парсеру пора остановиться."""
nonlocal cancelled, exhausted
prog.count_fetched(fetched)
if prog.over_budget:
exhausted = True
prog.log("warning", f"бюджет времени {prog.budget_s:.0f}с исчерпан на выборке")
return False
if prog.should_flush():
if _write_run(run_id, **prog.snapshot()):
cancelled = True
prog.log("warning", "отмена по запросу из админки")
return False
prog.mark_flushed()
return True
# fetch+transform без записи в JSONL — работаем in-memory
raw_docs = parser.fetch(**fetch_kwargs, progress_cb=on_fetch_progress)
if getattr(parser, "bulk", False):
# Массовые источники (дамп Википедии, бакет PMC) отдают генератор:
# миллионы статей нельзя ни держать в памяти, ни писать по одной
# через ORM — для них отдельный путь с пакетной записью
cancelled, exhausted = _ingest_bulk(parser, raw_docs, prog, run_id, source_id)
else:
cancelled, exhausted = _ingest_one_by_one(parser, raw_docs, prog, run_id)
except Exception as e:
error_msg = str(e)[:500]
prog.log("error", f"прогон упал: {error_msg}")
logger.error(f"run_parser source={source_id} ошибка: {e}", exc_info=True)
if error_msg:
status = "error"
elif cancelled:
status = "cancelled"
elif exhausted:
status = "partial"
else:
status = "done"
prog.stage = "finished"
prog.log(
"info" if status in ("done", "partial") else "warning",
f"итог: {status}, добавлено {prog.added}, дублей {prog.duplicates}, "
f"ошибок {prog.failed}, за {prog.elapsed:.0f}с",
)
_write_run(
run_id,
status=status,
error=error_msg,
finished_at=datetime.now(UTC),
**prog.snapshot(),
)
with db_session() as session:
src = session.get(ParseSource, source_id)
if src:
src.last_status = status
src.last_error = error_msg
src.docs_added = (src.docs_added or 0) + prog.added
src.last_run_at = datetime.now(UTC)
session.commit()
return {
"status": status,
"run_id": run_id,
"added": prog.added,
"duplicates": prog.duplicates,
"failed": prog.failed,
"error": error_msg,
}
@celery_app.task(name="index.ingest_upload", bind=True, max_retries=2, default_retry_delay=60)
def ingest_upload(
self,
minio_key: str,
filename: str,
meta: dict[str, Any] | None = None,
) -> dict[str, Any]:
"""Добавить загруженный админом файл в базу документов (корпус для сравнения).
Это не проверка плагиата: файл сразу становится источником, с которым
сравниваются работы студентов. Текст извлекается тем же кодом, что и в
extract_and_check, дальше — обычный add_document (дедуп, fingerprints,
LSH, Elasticsearch, эмбеддинги).
"""
meta = meta or {}
try:
minio = get_minio()
response = minio.get_object(settings.MINIO_BUCKET_DOCS, minio_key)
file_data = response.read()
response.close()
response.release_conn()
ext = Path(filename).suffix.lower()
if ext == ".pdf":
text = extract_text_from_pdf(file_data)
elif ext == ".docx":
text = extract_text_from_docx(file_data)
else:
text = extract_text_from_txt(file_data)
if not text.strip():
raise ValueError("не удалось извлечь текст")
except ValueError as exc:
logger.warning(f"ingest_upload {minio_key}: {exc}")
return {"status": "failed", "minio_key": minio_key, "error": str(exc)}
except Exception as exc:
logger.error(f"ingest_upload {minio_key}: {exc}", exc_info=True)
raise self.retry(exc=exc, countdown=60) from exc
doc_data = {
"source": meta.get("source") or "manual_upload",
# Ключ MinIO уникален (UUID в имени) — годится как ext_id для дедупа
"ext_id": f"upload:{minio_key}",
"title": meta.get("title") or Path(filename).stem,
"authors": meta.get("authors") or [],
"year": meta.get("year"),
"lang": meta.get("lang"),
"abstract": text[:2000],
"full_text": text,
"minio_key": minio_key,
}
result = add_document(doc_data)
logger.info(
f"ingest_upload: {filename!r} → {result.get('status')} "
f"(doc_id={result.get('doc_id')}, {len(text)} симв.)"
)
return {"status": result.get("status"), "doc_id": result.get("doc_id"), "filename": filename}