Заливка корпуса была чёрным ящиком: у источника только last_status (idle/running/done/error), без «сколько из скольки», без причины падения и без способа остановить начатое. Теперь каждый запуск создаёт строку parse_runs, куда воркер раз в ~2с пишет стадию, счётчики и журнал событий. Админка: - шкала загрузки у каждого источника (0→50% выборка, 50→100% индексация), раскрытая строка — журнал прогона по шагам с таймингами; - «Запустить всё» / «Остановить всё» и остановка по одному источнику (кооперативная отмена: воркер останавливается сам, не рвя запись в базу); - пакетное добавление источников (тип + список тем), тип pmc в форме; - загрузка PDF/DOCX/TXT прямо в базу сравнения (index.ingest_upload); - страница «Отладка»: воркеры Celery и их текущие таски, очереди RabbitMQ, покрытие корпуса эмбеддингами, зависшие и упавшие прогоны, конфиг бэкендов. Защита от краш-лупа по consumer_timeout RabbitMQ (docs/DR-HA.md §6), без неё массовый запуск 170+ источников гарантированно ронял воркер: - PARSER_TIME_BUDGET_S (1500с) — прогон закругляется сам и помечается partial; - worker_prefetch_multiplier=1 — таймаут считается от ДОСТАВКИ сообщения, и с дефолтным префетчем очередь долгих run_parser убивала канал на задачах, которые ещё не начинались. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
178 lines
6.8 KiB
Python
178 lines
6.8 KiB
Python
"""Базовый класс для всех парсеров источников.
|
||
|
||
Определяет унифицированный интерфейс и формат документов.
|
||
"""
|
||
|
||
import json
|
||
import logging
|
||
from abc import ABC, abstractmethod
|
||
from collections.abc import Callable
|
||
from pathlib import Path
|
||
from typing import Any
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
# Колбэк прогресса выборки: вызывается после каждой страницы результатов с
|
||
# накопленным количеством документов. Возврат False — просьба остановиться
|
||
# (кооперативная отмена из админки или исчерпанный бюджет времени таска,
|
||
# см. worker-indexer/app/progress.py). Парсер обязан вернуть уже собранное,
|
||
# а не бросать исключение: частичная выборка — валидный результат.
|
||
ProgressCallback = Callable[[int], bool]
|
||
|
||
|
||
# Унифицированный формат документа
|
||
UNIFIED_SCHEMA = {
|
||
"source": str, # openalex, cyberleninka, arxiv, wikipedia_ru, wikipedia_en
|
||
"ext_id": str, # Внешний ID (уникален в рамках источника)
|
||
"doi": "str|None",
|
||
"title": str,
|
||
"authors": list, # [{"last_name": str, "first_name": str, "initials": str}]
|
||
"year": "int|None",
|
||
"lang": "str|None", # ru, en
|
||
"journal": "str|None",
|
||
"volume": "str|None",
|
||
"issue": "str|None",
|
||
"pages": "str|None",
|
||
"abstract": "str|None",
|
||
"url": "str|None",
|
||
"full_text": "str|None", # Для fingerprinting (не хранится в PostgreSQL)
|
||
}
|
||
|
||
|
||
class BaseParser(ABC):
|
||
"""Базовый класс парсера источника."""
|
||
|
||
source_name: str = "unknown"
|
||
|
||
def __init__(self) -> None:
|
||
self.logger = logging.getLogger(f"parser.{self.source_name}")
|
||
|
||
@abstractmethod
|
||
def fetch(self, **kwargs) -> list[dict[str, Any]]:
|
||
"""
|
||
Получить сырые документы из источника.
|
||
|
||
Args:
|
||
**kwargs: Специфичные для источника параметры (query, limit и т.д.).
|
||
Все парсеры принимают ещё и `progress_cb` (ProgressCallback) —
|
||
отчёт о прогрессе после каждой страницы и точка кооперативной
|
||
остановки.
|
||
|
||
Returns:
|
||
Список сырых словарей из API источника
|
||
"""
|
||
...
|
||
|
||
def transform(self, raw: dict[str, Any]) -> dict[str, Any]:
|
||
"""
|
||
Привести сырой документ к унифицированному формату.
|
||
|
||
Args:
|
||
raw: Сырой словарь из API источника
|
||
|
||
Returns:
|
||
Документ в унифицированном формате
|
||
"""
|
||
raise NotImplementedError(
|
||
f"Парсер {self.source_name!r} должен реализовать метод transform()"
|
||
)
|
||
|
||
def save_jsonl(self, docs: list[dict[str, Any]], output_path: Path) -> None:
|
||
"""
|
||
Сохранить документы в JSONL формате (один JSON объект на строку).
|
||
|
||
Args:
|
||
docs: Список документов
|
||
output_path: Путь к выходному файлу (дополняется, не перезаписывается)
|
||
"""
|
||
output_path.parent.mkdir(parents=True, exist_ok=True)
|
||
with output_path.open("a", encoding="utf-8") as f:
|
||
for doc in docs:
|
||
f.write(json.dumps(doc, ensure_ascii=False) + "\n")
|
||
|
||
def run(self, output_dir: Path, **kwargs) -> list[dict[str, Any]]:
|
||
"""
|
||
Запустить парсер: fetch → transform → save.
|
||
|
||
Args:
|
||
output_dir: Директория для сохранения результатов
|
||
**kwargs: Параметры для метода fetch()
|
||
|
||
Returns:
|
||
Список трансформированных документов
|
||
"""
|
||
self.logger.info(f"[{self.source_name}] Начало парсинга...")
|
||
|
||
raw_docs = self.fetch(**kwargs)
|
||
self.logger.info(f"[{self.source_name}] Получено {len(raw_docs)} документов")
|
||
|
||
transformed = []
|
||
errors = 0
|
||
for raw in raw_docs:
|
||
try:
|
||
doc = self.transform(raw)
|
||
if doc and doc.get("title") and doc.get("ext_id"):
|
||
transformed.append(doc)
|
||
except Exception as e:
|
||
errors += 1
|
||
self.logger.warning(f"Ошибка трансформации: {e}")
|
||
|
||
if errors:
|
||
self.logger.warning(f"[{self.source_name}] Ошибок трансформации: {errors}")
|
||
|
||
out_file = output_dir / f"{self.source_name}.jsonl"
|
||
self.save_jsonl(transformed, out_file)
|
||
|
||
self.logger.info(
|
||
f"[{self.source_name}] Сохранено {len(transformed)} документов → {out_file}"
|
||
)
|
||
|
||
return transformed
|
||
|
||
@staticmethod
|
||
def normalize_authors(raw_authors: list) -> list[dict[str, str]]:
|
||
"""
|
||
Нормализовать список авторов к единому формату.
|
||
|
||
Args:
|
||
raw_authors: Авторы из API (разный формат у каждого источника)
|
||
|
||
Returns:
|
||
Список {"last_name": str, "first_name": str, "initials": str}
|
||
"""
|
||
result = []
|
||
for author in raw_authors:
|
||
if isinstance(author, str):
|
||
parts = author.strip().split()
|
||
last_name = parts[0] if parts else ""
|
||
first_name = " ".join(parts[1:]) if len(parts) > 1 else ""
|
||
elif isinstance(author, dict):
|
||
last_name = (
|
||
author.get("last_name") or
|
||
author.get("family") or
|
||
author.get("surname") or ""
|
||
).strip()
|
||
first_name = (
|
||
author.get("first_name") or
|
||
author.get("given") or ""
|
||
).strip()
|
||
else:
|
||
continue
|
||
|
||
if not last_name:
|
||
continue
|
||
|
||
# Инициалы
|
||
initials = ""
|
||
if first_name:
|
||
parts = first_name.split()
|
||
initials = "".join(p[0].upper() + "." for p in parts if p)
|
||
|
||
result.append({
|
||
"last_name": last_name,
|
||
"first_name": first_name,
|
||
"initials": initials,
|
||
})
|
||
|
||
return result
|