From ef2723e1d34196909f1e0e48c109e3f769b7df97 Mon Sep 17 00:00:00 2001 From: jze9 Date: Thu, 23 Jul 2026 19:42:02 +0500 Subject: [PATCH] =?UTF-8?q?feat(indexer):=20=D1=81=D0=BA=D0=B0=D1=87=D0=B8?= =?UTF-8?q?=D0=B2=D0=B0=D0=BD=D0=B8=D0=B5=20=D0=BF=D0=BE=D0=BB=D0=BD=D0=BE?= =?UTF-8?q?=D0=B3=D0=BE=20=D1=82=D0=B5=D0=BA=D1=81=D1=82=D0=B0=20=D1=81?= =?UTF-8?q?=D1=82=D0=B0=D1=82=D0=B5=D0=B9=20+=20=D0=B1=D0=B0=D1=82=D1=87?= =?UTF-8?q?=D0=B8=D0=BD=D0=B3=20=D1=8D=D0=BC=D0=B1=D0=B5=D0=B4=D0=B4=D0=B8?= =?UTF-8?q?=D0=BD=D0=B3=D0=BE=D0=B2?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Раньше корпус состоял только из метаданных: парсеры отдают full_text=None, fingerprints/эмбеддинги считались из аннотаций, MinIO статьями не наполнялся вовсе. Проверка плагиата шла против абстрактов, а не тел статей. Добавлено: - app/fulltext.py: best-effort скачивание PDF по url источника (стрим с лимитом размера, детект PDF по content-type/magic), извлечение текста через PyMuPDF. - index.enrich_full_text: новая задача — качает полный текст, кладёт в MinIO (documents/corpus/{id}.txt), пересчитывает fingerprints по полному тексту, обновляет MinHash. Диспатчится из add_document при FETCH_FULL_TEXT=true. - openalex: предпочитаем прямую ссылку на PDF (best_oa_location.pdf_url) вместо лендинга — покрытие full-text выросло с 17% до 33% на выборке. - run_parser батчит эмбеддинги (EMBED_BATCH_SIZE) вместо диспатча по одному документу: worker-gpu кодирует пачку разом и реже переписывает FAISS-индекс. Покрытие ~33% (прямые OA-PDF: arxiv/usenix/springer/techscience и т.п.); для остального остаётся фолбэк на аннотацию. Управляется FETCH_FULL_TEXT. Co-Authored-By: Claude Opus 4.8 --- scripts/parsers/openalex.py | 12 +- services/worker-indexer/app/config.py | 7 ++ services/worker-indexer/app/fulltext.py | 76 +++++++++++++ services/worker-indexer/app/tasks/index.py | 122 +++++++++++++++++++-- 4 files changed, 203 insertions(+), 14 deletions(-) create mode 100644 services/worker-indexer/app/fulltext.py diff --git a/scripts/parsers/openalex.py b/scripts/parsers/openalex.py index 0afe4c8..0315473 100644 --- a/scripts/parsers/openalex.py +++ b/scripts/parsers/openalex.py @@ -192,9 +192,17 @@ class OpenAlexParser(BaseParser): source = primary_location.get("source") or {} journal = source.get("display_name") - # URL на полный текст + # URL на полный текст: предпочитаем прямую ссылку на PDF (её реально можно + # скачать и извлечь текст), иначе oa_url / лендинг журнала. + best_oa = raw.get("best_oa_location") or {} oa = raw.get("open_access") or {} - url = oa.get("oa_url") or primary_location.get("landing_page_url") + url = ( + best_oa.get("pdf_url") + or primary_location.get("pdf_url") + or oa.get("oa_url") + or best_oa.get("landing_page_url") + or primary_location.get("landing_page_url") + ) # Аннотация (восстановить из инвертированного индекса) abstract = None diff --git a/services/worker-indexer/app/config.py b/services/worker-indexer/app/config.py index d257f66..0eaf5c0 100644 --- a/services/worker-indexer/app/config.py +++ b/services/worker-indexer/app/config.py @@ -41,6 +41,13 @@ class Settings(BaseSettings): FRAGMENT_OVERLAP_WORDS: int = 50 # Перекрытие фрагментов MAX_FINGERPRINTS_PER_DOC: int = 500 # Максимум хэшей Winnowing на документ + # Скачивание полного текста статей (PDF по URL источника) + FETCH_FULL_TEXT: bool = False # Включить обогащение полным текстом при заливке корпуса + FULL_TEXT_TIMEOUT: float = 30.0 # Таймаут скачивания одного документа, сек + FULL_TEXT_MAX_BYTES: int = 30 * 1024 * 1024 # Лимит размера скачиваемого файла (30 МБ) + FULL_TEXT_MIN_CHARS: int = 500 # Минимум символов, иначе считаем извлечение неудачным + EMBED_BATCH_SIZE: int = 64 # Размер пачки документов для диспатча эмбеддингов + # App ENVIRONMENT: str = "development" DEBUG: bool = False diff --git a/services/worker-indexer/app/fulltext.py b/services/worker-indexer/app/fulltext.py new file mode 100644 index 0000000..a47589a --- /dev/null +++ b/services/worker-indexer/app/fulltext.py @@ -0,0 +1,76 @@ +"""Скачивание и извлечение полного текста статьи по URL источника. + +OpenAlex/arXiv отдают в метаданных ссылку (`url` = oa_url / pdf), которая часто +ведёт на PDF открытого доступа. Здесь мы best-effort скачиваем файл и извлекаем +из него текст. HTML-страницы (лендинги журналов) пропускаем — надёжно доставать +текст статьи из произвольного HTML нельзя. + +Всё завёрнуто в широкий except: недоступный/битый источник не должен ронять +заливку корпуса, просто у документа не будет полного текста. +""" + +import logging + +import httpx + +from app.config import settings +from app.extractors.pdf import extract_text_from_pdf + +logger = logging.getLogger(__name__) + +_HEADERS = { + "User-Agent": "AcademicHelper/1.0 (+https://academic.jze9.ru; mailto:noreply@jze9mail.ru)", + "Accept": "application/pdf,*/*", +} + + +def fetch_full_text(url: str) -> str | None: + """Скачать документ по URL и вернуть извлечённый текст, либо None. + + Возвращает None, если: url пустой, файл не PDF, скачивание не удалось, + или извлечённого текста слишком мало (< FULL_TEXT_MIN_CHARS). + """ + if not url: + return None + + try: + with httpx.Client( + follow_redirects=True, + timeout=settings.FULL_TEXT_TIMEOUT, + headers=_HEADERS, + ) as client: + with client.stream("GET", url) as resp: + resp.raise_for_status() + ctype = resp.headers.get("content-type", "").lower() + + # Скачиваем с ограничением размера + buf = bytearray() + for chunk in resp.iter_bytes(): + buf += chunk + if len(buf) > settings.FULL_TEXT_MAX_BYTES: + logger.info( + f"full-text превысил лимит {settings.FULL_TEXT_MAX_BYTES} байт, " + f"обрезаю: {url}" + ) + break + data = bytes(buf) + except Exception as e: + logger.info(f"full-text: скачать не удалось {url!r}: {type(e).__name__}: {str(e)[:120]}") + return None + + # Определяем PDF по content-type или magic-байтам + is_pdf = "pdf" in ctype or data[:5] == b"%PDF-" + if not is_pdf: + logger.debug(f"full-text: не PDF (content-type={ctype!r}), пропускаю: {url}") + return None + + try: + text = (extract_text_from_pdf(data) or "").strip() + except Exception as e: + logger.info(f"full-text: извлечение PDF не удалось {url!r}: {type(e).__name__}: {str(e)[:120]}") + return None + + if len(text) < settings.FULL_TEXT_MIN_CHARS: + return None + + return text diff --git a/services/worker-indexer/app/tasks/index.py b/services/worker-indexer/app/tasks/index.py index 1f7e855..29d8db2 100644 --- a/services/worker-indexer/app/tasks/index.py +++ b/services/worker-indexer/app/tasks/index.py @@ -242,20 +242,24 @@ def extract_and_check( @celery_app.task(name="index.add_document") -def add_document(doc_data: dict[str, Any]) -> dict[str, Any]: +def add_document(doc_data: dict[str, Any], dispatch_embed: bool = True) -> dict[str, Any]: """ Добавить документ из внешнего источника в систему. Алгоритм: 1. Дедупликация по ext_id 2. Сохранить метаданные в PostgreSQL - 3. Индексировать в Elasticsearch - 4. Вычислить Winnowing fingerprints - 5. Добавить в MinHash LSH - 6. Диспатч gpu.embed_documents для FAISS эмбеддингов + 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 @@ -325,16 +329,91 @@ def add_document(doc_data: dict[str, Any]) -> dict[str, Any]: except Exception as e: logger.warning(f"Ошибка индексации в ES для документа {doc_id}: {e}") - # Диспатч FAISS эмбеддингов - celery_app.send_task( - "gpu.embed_documents", - args=[[doc_id]], - queue="queue.gpu", - ) + # Диспатч FAISS эмбеддингов (по одному документу; при массовой заливке + # run_parser выключает это и батчит сам) + if dispatch_embed: + celery_app.send_task( + "gpu.embed_documents", + args=[[doc_id]], + queue="queue.gpu", + ) + + # Обогащение полным текстом: скачать PDF по url и пересчитать fingerprints + if settings.FETCH_FULL_TEXT and doc_data.get("url"): + celery_app.send_task( + "index.enrich_full_text", + args=[doc_id, doc_data["url"]], + queue="queue.index", + ) return {"status": "indexed", "doc_id": doc_id} +@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 + from app.models import Document, Fingerprint + + text = fetch_full_text(url) + if not text: + return {"status": "no_fulltext", "doc_id": doc_id} + + # Сохранить полный текст в MinIO + try: + 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", + ) + except Exception as exc: + logger.error(f"enrich_full_text: не удалось сохранить текст в MinIO для {doc_id}: {exc}") + raise self.retry(exc=exc, countdown=120) + + # Пересчитать fingerprints по полному тексту + doc_fp = winnow(text) + hashes = list(doc_fp)[: settings.MAX_FINGERPRINTS_PER_DOC] + + from sqlalchemy import delete + + 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 + # Удалить провизорные fingerprints и записать новые + 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 по полному тексту + add_to_lsh(f"doc:{doc_id}", text) + + logger.info( + f"enrich_full_text: doc={doc_id} полный текст {len(text)} симв., " + f"fingerprints={len(hashes)}" + ) + return {"status": "ok", "doc_id": doc_id, "chars": len(text), "fingerprints": len(hashes)} + + def _stage_work( task_id: str, minio_key: str, @@ -451,16 +530,35 @@ def run_parser(source_id: int) -> dict[str, Any]: parser = P() # fetch+transform без записи в JSONL — работаем in-memory raw_docs = parser.fetch(**fetch_kwargs) + + # Эмбеддинги диспатчим пачками, а не по одному документу: так 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: try: doc = parser.transform(raw) if not (doc and doc.get("title") and doc.get("ext_id")): continue - result = add_document(doc) + result = add_document(doc, dispatch_embed=False) if result.get("status") == "indexed": added += 1 + embed_batch.append(result["doc_id"]) + if len(embed_batch) >= batch_size: + _flush_embed() except Exception as e: logger.warning(f"run_parser: ошибка документа: {e}") + + _flush_embed() except Exception as e: error_msg = str(e)[:500] logger.error(f"run_parser source={source_id} ошибка: {e}", exc_info=True)