Compare commits
8 Commits
d86bb606c7
...
main
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
d137074ed5 | ||
|
|
9123844a87 | ||
|
|
263a521d63 | ||
|
|
7adccd0362 | ||
|
|
7348051c7b | ||
|
|
3d4fcecbb2 | ||
|
|
20c1d00443 | ||
|
|
5ceec1e2e3 |
@@ -152,6 +152,10 @@ services:
|
|||||||
depends_on:
|
depends_on:
|
||||||
elasticsearch:
|
elasticsearch:
|
||||||
condition: service_healthy
|
condition: service_healthy
|
||||||
|
# CORE не отвечает на ключ с российских адресов — парсер ходит через
|
||||||
|
# sing-box (см. scripts/parsers/core.py)
|
||||||
|
singbox-proxy:
|
||||||
|
condition: service_started
|
||||||
|
|
||||||
worker-notifier:
|
worker-notifier:
|
||||||
build:
|
build:
|
||||||
|
|||||||
@@ -43,6 +43,7 @@ fi
|
|||||||
|
|
||||||
REBUILD=""
|
REBUILD=""
|
||||||
FRONTEND_CHANGED=0
|
FRONTEND_CHANGED=0
|
||||||
|
PARSERS_CHANGED=0
|
||||||
if [ "$FULL" = 1 ]; then
|
if [ "$FULL" = 1 ]; then
|
||||||
REBUILD="$ALL_BACKEND"
|
REBUILD="$ALL_BACKEND"
|
||||||
FRONTEND_CHANGED=1
|
FRONTEND_CHANGED=1
|
||||||
@@ -51,6 +52,12 @@ else
|
|||||||
echo "$CHANGED" | grep -qE "^services/$svc/" && REBUILD="$REBUILD $svc"
|
echo "$CHANGED" | grep -qE "^services/$svc/" && REBUILD="$REBUILD $svc"
|
||||||
done
|
done
|
||||||
echo "$CHANGED" | grep -qE '^services/frontend/' && FRONTEND_CHANGED=1
|
echo "$CHANGED" | grep -qE '^services/frontend/' && FRONTEND_CHANGED=1
|
||||||
|
# Парсеры не вшиты в образ, а смонтированы в worker-indexer (см. compose),
|
||||||
|
# поэтому пересобирать нечего. Но процесс воркера держит модуль парсера в
|
||||||
|
# памяти с прошлого прогона: без перезапуска правка молча не применяется —
|
||||||
|
# заливка продолжает ходить по старому коду, и это не видно ниоткуда, кроме
|
||||||
|
# тела запросов в логах. Ловилось на CORE 16.09.2026.
|
||||||
|
echo "$CHANGED" | grep -qE '^scripts/parsers/' && PARSERS_CHANGED=1
|
||||||
fi
|
fi
|
||||||
REBUILD="$(echo "$REBUILD" | xargs || true)"
|
REBUILD="$(echo "$REBUILD" | xargs || true)"
|
||||||
|
|
||||||
@@ -111,6 +118,13 @@ echo "==> [4/6] Применение секретов к неизменивши
|
|||||||
# shellcheck disable=SC2086
|
# shellcheck disable=SC2086
|
||||||
$COMPOSE up -d $ALL_BACKEND
|
$COMPOSE up -d $ALL_BACKEND
|
||||||
|
|
||||||
|
# Правка парсера доезжает монтированием, но модуль уже загружен в память
|
||||||
|
# воркера — перезапускаем, если образ и так не пересобирался выше.
|
||||||
|
if [ "$PARSERS_CHANGED" = 1 ] && ! echo " $REBUILD " | grep -q " worker-indexer "; then
|
||||||
|
echo "==> [4.5/6] Парсеры менялись — перезапуск worker-indexer"
|
||||||
|
$COMPOSE restart worker-indexer
|
||||||
|
fi
|
||||||
|
|
||||||
echo "==> [5/6] Миграции БД (с ретраем)"
|
echo "==> [5/6] Миграции БД (с ретраем)"
|
||||||
for i in $(seq 1 10); do
|
for i in $(seq 1 10); do
|
||||||
$COMPOSE exec -T api alembic upgrade head && break
|
$COMPOSE exec -T api alembic upgrade head && break
|
||||||
|
|||||||
@@ -15,11 +15,36 @@
|
|||||||
Останавливается сам, когда источник исчерпан: парсер перестаёт отдавать статьи,
|
Останавливается сам, когда источник исчерпан: парсер перестаёт отдавать статьи,
|
||||||
прогон закрывается с нулём добавленных, и скрипт больше его не трогает
|
прогон закрывается с нулём добавленных, и скрипт больше его не трогает
|
||||||
(--stop-when-empty, по умолчанию включено).
|
(--stop-when-empty, по умолчанию включено).
|
||||||
|
|
||||||
|
Перед подсчётом «занято ли» снимает зависшие прогоны: если worker-indexer упал
|
||||||
|
или перезапустился, прогон навсегда остаётся в статусе running с замолчавшим
|
||||||
|
heartbeat — без этой чистки сторож видит «занято» бесконечно и не запускает
|
||||||
|
вообще ничего, для всех источников сразу (порог тот же, что и в панели
|
||||||
|
отладки — services/api/app/api/admin.py, stale_cutoff).
|
||||||
|
|
||||||
|
Молчащий heartbeat сам по себе ещё не значит «мёртв»: пачка COPY на большой
|
||||||
|
базе идёт десятки минут без единого тика. Поэтому прогон снимается, только
|
||||||
|
если его таска нет среди выполняющихся на воркерах queue.index; не ответил
|
||||||
|
хоть один такой воркер — не снимается ничего. Иначе сторож снимал живой
|
||||||
|
медленный прогон и тут же запускал его дубль по тем же статьям.
|
||||||
"""
|
"""
|
||||||
|
|
||||||
import argparse
|
import argparse
|
||||||
|
|
||||||
BULK_TYPES = ("wikipedia_ru", "pmc_bulk")
|
BULK_TYPES = ("wikipedia_ru", "pmc_bulk", "core")
|
||||||
|
STALE_MINUTES = 10
|
||||||
|
|
||||||
|
|
||||||
|
def live_task_ids(celery_app) -> set[str] | None:
|
||||||
|
"""ID тасков, выполняющихся на воркерах queue.index; None — кто-то из них не ответил."""
|
||||||
|
queues = celery_app.control.inspect(timeout=10).active_queues() or {}
|
||||||
|
workers = [w for w, qs in queues.items() if any(q["name"] == "queue.index" for q in qs)]
|
||||||
|
if not workers:
|
||||||
|
return None
|
||||||
|
active = celery_app.control.inspect(destination=workers, timeout=10).active()
|
||||||
|
if not active or set(active) != set(workers):
|
||||||
|
return None
|
||||||
|
return {t["id"] for tasks in active.values() for t in tasks}
|
||||||
|
|
||||||
|
|
||||||
def main() -> None:
|
def main() -> None:
|
||||||
@@ -38,7 +63,44 @@ def main() -> None:
|
|||||||
|
|
||||||
types = tuple(t.strip() for t in args.types.split(",") if t.strip())
|
types = tuple(t.strip() for t in args.types.split(",") if t.strip())
|
||||||
|
|
||||||
|
live = live_task_ids(celery_app)
|
||||||
|
|
||||||
with db_session() as s:
|
with db_session() as s:
|
||||||
|
if live is None:
|
||||||
|
print("воркеры queue.index не ответили — живость прогонов не проверить, не снимаю")
|
||||||
|
reaped = []
|
||||||
|
else:
|
||||||
|
reaped = s.execute(text(f"""
|
||||||
|
UPDATE parse_runs SET status='error', stage='finished',
|
||||||
|
error='зависший прогон снят автоматически: таск не выполняется, heartbeat молчал дольше {STALE_MINUTES} минут',
|
||||||
|
finished_at=now()
|
||||||
|
WHERE status='running'
|
||||||
|
AND (heartbeat_at IS NULL OR heartbeat_at < now() - interval '{STALE_MINUTES} minutes')
|
||||||
|
AND (celery_task_id IS NULL OR NOT (celery_task_id = ANY(CAST(:live AS text[]))))
|
||||||
|
RETURNING id, source_id
|
||||||
|
"""), {"live": sorted(live)}).all()
|
||||||
|
# queued без heartbeat — тот же зомби другого вида: send_task не дошёл
|
||||||
|
# или упал между двумя commit при постановке в очередь (celery_task_id
|
||||||
|
# так и остался пустым), воркер о такой задаче никогда не узнает
|
||||||
|
reaped += s.execute(text(f"""
|
||||||
|
UPDATE parse_runs SET status='error', stage='finished',
|
||||||
|
error='зависший прогон снят автоматически: висел в очереди дольше {STALE_MINUTES} минут',
|
||||||
|
finished_at=now()
|
||||||
|
WHERE status='queued'
|
||||||
|
AND celery_task_id IS NULL
|
||||||
|
AND started_at < now() - interval '{STALE_MINUTES} minutes'
|
||||||
|
RETURNING id, source_id
|
||||||
|
""")).all()
|
||||||
|
if reaped:
|
||||||
|
for source_id in {row.source_id for row in reaped}:
|
||||||
|
s.execute(text("""
|
||||||
|
UPDATE parse_sources SET last_status='error',
|
||||||
|
last_error='зависший прогон снят автоматически'
|
||||||
|
WHERE id=:sid AND last_status='running'
|
||||||
|
"""), {"sid": source_id})
|
||||||
|
s.commit()
|
||||||
|
print(f"снято зависших прогонов: {len(reaped)} ({', '.join(str(r.id) for r in reaped)})")
|
||||||
|
|
||||||
active = s.execute(text(
|
active = s.execute(text(
|
||||||
"SELECT count(*) FROM parse_runs WHERE status IN ('queued','running')"
|
"SELECT count(*) FROM parse_runs WHERE status IN ('queued','running')"
|
||||||
)).scalar_one()
|
)).scalar_one()
|
||||||
@@ -53,12 +115,25 @@ def main() -> None:
|
|||||||
|
|
||||||
started = []
|
started = []
|
||||||
for source_id, source_type, last_run_id in sources:
|
for source_id, source_type, last_run_id in sources:
|
||||||
# Источник исчерпан, если прошлый прогон не добавил ни одной статьи
|
# Источник исчерпан, если прогоны перестали добавлять статьи: дальше
|
||||||
# и не был отменён: дальше в дампе/бакете для нас ничего нет
|
# в дампе/бакете/выдаче для нас ничего нет.
|
||||||
|
#
|
||||||
|
# Но ОДИН пустой прогон — ещё не приговор: ровно так же выглядят
|
||||||
|
# обрыв связи с API и прогон на устаревшем коде парсера. На CORE
|
||||||
|
# 16.09.2026 из-за этого источник был помечен исчерпанным после
|
||||||
|
# единственной неудачи и больше не запускался — молча, без ошибки.
|
||||||
|
# Поэтому ждём два пустых прогона подряд: разовый сбой переживём,
|
||||||
|
# а реально кончившийся источник остановится всего на один прогон
|
||||||
|
# позже.
|
||||||
if last_run_id and not args.keep_going_empty:
|
if last_run_id and not args.keep_going_empty:
|
||||||
prev = s.get(ParseRun, last_run_id)
|
last_two = s.execute(text("""
|
||||||
if prev and prev.status == "done" and prev.added == 0 and prev.fetched == 0:
|
SELECT status, added, fetched FROM parse_runs
|
||||||
print(f"{source_type}: источник исчерпан, пропускаем")
|
WHERE source_id = :sid ORDER BY id DESC LIMIT 2
|
||||||
|
"""), {"sid": source_id}).all()
|
||||||
|
if len(last_two) == 2 and all(
|
||||||
|
r.status == "done" and r.added == 0 and r.fetched == 0 for r in last_two
|
||||||
|
):
|
||||||
|
print(f"{source_type}: источник исчерпан (два пустых прогона), пропускаем")
|
||||||
continue
|
continue
|
||||||
|
|
||||||
src = s.get(ParseSource, source_id)
|
src = s.get(ParseSource, source_id)
|
||||||
|
|||||||
282
scripts/parsers/core.py
Normal file
282
scripts/parsers/core.py
Normal file
@@ -0,0 +1,282 @@
|
|||||||
|
"""Парсер CORE (core.ac.uk) — полные тексты открытых репозиториев.
|
||||||
|
|
||||||
|
Зачем нужен: полный текст есть у малой доли корпуса, и по-русски дыра самая
|
||||||
|
большая — CyberLeninka отдаёт только HTML-страницу статьи, тела не даёт. CORE
|
||||||
|
агрегирует открытые репозитории (в том числе вузовские русско- и
|
||||||
|
белорусскоязычные) и кладёт готовый текст прямо в выдачу поиска — его не надо
|
||||||
|
ни скачивать отдельно, ни извлекать из PDF.
|
||||||
|
|
||||||
|
Замер на живом API 16.09.2026: страница в 100 записей приходит за ~7с и
|
||||||
|
содержит 47 статей с текстом длиннее 1500 символов, медиана — 12.7 тыс.
|
||||||
|
символов. То есть один запрос ≈ 47 документов корпуса.
|
||||||
|
|
||||||
|
Чем сузить выборку до русских текстов: в общем потоке CORE кириллицы нет
|
||||||
|
вовсе (проверено: 0 из 48 полных текстов за 2019). Язык у работ размечен
|
||||||
|
плохо — `language` сплошь None или "zz", фильтровать по нему нельзя. Зато
|
||||||
|
работает отбор по архиву-поставщику: запрос вида
|
||||||
|
`(repositories.id:1298 OR repositories.id:21908 OR ...)` по вузовским
|
||||||
|
репозиториям России и Беларуси даёт 79 полных текстов из 100 записей, 78 из
|
||||||
|
них — кириллица. Сам список архивов живёт в поле `query` источника, а не в
|
||||||
|
коде: его правят из админки, не пересобирая образ.
|
||||||
|
|
||||||
|
Массовый парсер (`bulk = True`): отдаёт статьи генератором и пишется пачками
|
||||||
|
через COPY, как Википедия и PMC (см. bulk_writer.py). Статьи без текста
|
||||||
|
пропускаются молча — они уже есть у нас как метаданные из OpenAlex, а
|
||||||
|
обогащать корпус нечем.
|
||||||
|
|
||||||
|
Грабли API, все проверены живьём:
|
||||||
|
- **с российских адресов ключ не работает**: прямой запрос с прода виснет без
|
||||||
|
ошибки (соединение есть, тело ответа не приходит), через `singbox-proxy` тот
|
||||||
|
же запрос отвечает за 1.7с. Анонимные короткие запросы проходят и напрямую,
|
||||||
|
из-за чего поломка выглядит как сетевая. Адрес прокси — `CORE_PROXY_URL`;
|
||||||
|
- у `/v3/search/works` обязателен слэш на конце: иначе 301, а при редиректе
|
||||||
|
теряется заголовок Authorization и запрос уходит анонимным;
|
||||||
|
- `offset` упирается в 100 000 (под капотом Azure Search, глубже — 400),
|
||||||
|
поэтому выборка режется на части — по одному архиву-поставщику на запрос;
|
||||||
|
- **условия через `AND` не связываются**: `(repositories.id:...) AND yearPublished:2026`
|
||||||
|
отдаёт 2.3 млн работ вперемешку по годам, то есть условия объединяются по
|
||||||
|
«или» и год работает лишь как подсказка ранжированию. По отдельности каждое
|
||||||
|
условие фильтрует честно (`yearPublished:2026` — ровно 2026, `repositories.id:1298`
|
||||||
|
— ровно этот архив). Поэтому в запрос идёт РОВНО ОДНО условие: номер архива.
|
||||||
|
Годы, если заданы, отсекаются уже на нашей стороне;
|
||||||
|
- `fullText` — не фильтруемое поле, `_exists_:fullText` отвечает 500.
|
||||||
|
Отбирать статьи с текстом приходится на своей стороне, оплачивая трафиком;
|
||||||
|
- CORE агрегирует репозитории и отдаёт одну статью несколько раз под разными
|
||||||
|
id (проверено: 90810208 и 354263199 — один текст). Дедупликация корпуса
|
||||||
|
идёт только по ext_id, поэтому такие пары легли бы отдельными документами
|
||||||
|
и потом ловились бы как заимствование друг у друга. Отсеиваем по хешу
|
||||||
|
текста в пределах прогона: копии приходят в соседних строках выдачи.
|
||||||
|
"""
|
||||||
|
|
||||||
|
import hashlib
|
||||||
|
import logging
|
||||||
|
import os
|
||||||
|
import re
|
||||||
|
import time
|
||||||
|
from collections.abc import Iterator
|
||||||
|
from typing import Any
|
||||||
|
|
||||||
|
import httpx
|
||||||
|
from base import BaseParser, ProgressCallback
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
API_URL = "https://api.core.ac.uk/v3/search/works/" # слэш обязателен, см. докстринг
|
||||||
|
PAGE = 100 # максимум записей за запрос
|
||||||
|
MAX_OFFSET = 100_000 # потолок глубины у Azure Search
|
||||||
|
MIN_CHARS = 1500 # короче — обрывок или аннотация, а не статья
|
||||||
|
RETRIES = 3
|
||||||
|
# CORE отдаёт данные по ключу только за пределами РФ: с прода прямой запрос
|
||||||
|
# молча виснет (TCP есть, тело ответа не приходит), а тот же запрос через
|
||||||
|
# sing-box возвращается за 1.7с. Анонимные мелкие запросы проходят и напрямую,
|
||||||
|
# поэтому со стороны это выглядит как поломка сети, а не блокировка. Тот же
|
||||||
|
# приём, что для OpenRouter в worker-gpu. Пусто = ходить напрямую.
|
||||||
|
DEFAULT_PROXY = "socks5://singbox-proxy:1080"
|
||||||
|
|
||||||
|
|
||||||
|
class COREParser(BaseParser):
|
||||||
|
"""Полные тексты из CORE. Отдаёт статьи потоком, режет выборку по годам."""
|
||||||
|
|
||||||
|
source_name = "core"
|
||||||
|
bulk = True
|
||||||
|
|
||||||
|
def __init__(self, api_key: str | None = None, proxy_url: str | None = None) -> None:
|
||||||
|
super().__init__()
|
||||||
|
self.api_key = api_key or os.environ.get("CORE_API_KEY", "")
|
||||||
|
# Вне compose-сети (локальный прогон, тесты) имени singbox-proxy нет —
|
||||||
|
# там прокси отключают, выставив CORE_PROXY_URL пустым
|
||||||
|
self.proxy_url = self._proxy(proxy_url)
|
||||||
|
if not self.api_key:
|
||||||
|
logger.warning("CORE: ключ не задан (CORE_API_KEY) — API ответит 401")
|
||||||
|
# Страница весит ~5 МБ, обычный ответ ~7с. 30с — запас на джиттер в
|
||||||
|
# десять раз, но без риска съесть бюджет прогона: три попытки по 60с
|
||||||
|
# отъедали 180с из 1500с на одной залипшей странице (поймано 16.09.2026,
|
||||||
|
# та же грабля, что чинили в pmc_bulk)
|
||||||
|
self.client = httpx.Client(
|
||||||
|
timeout=30,
|
||||||
|
proxy=self.proxy_url or None,
|
||||||
|
headers={
|
||||||
|
"Authorization": f"Bearer {self.api_key}",
|
||||||
|
"User-Agent": "AcademicHelper/1.0 (noreply@jze9.ru)",
|
||||||
|
},
|
||||||
|
)
|
||||||
|
self.last_token: str | None = None
|
||||||
|
|
||||||
|
@staticmethod
|
||||||
|
def _proxy(proxy_url: str | None) -> str:
|
||||||
|
"""Адрес прокси: явный аргумент, иначе CORE_PROXY_URL, иначе умолчание."""
|
||||||
|
if proxy_url is not None:
|
||||||
|
return proxy_url
|
||||||
|
return os.environ.get("CORE_PROXY_URL", DEFAULT_PROXY)
|
||||||
|
|
||||||
|
def fetch( # type: ignore[override]
|
||||||
|
self,
|
||||||
|
limit: int = 1000,
|
||||||
|
query: str = "",
|
||||||
|
year_from: int | None = None,
|
||||||
|
year_to: int | None = None,
|
||||||
|
resume_token: str | None = None,
|
||||||
|
progress_cb: ProgressCallback | None = None,
|
||||||
|
**_ignored: Any,
|
||||||
|
) -> Iterator[dict[str, Any]]:
|
||||||
|
"""Статьи с полным текстом, архив за архивом.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
limit: сколько статей с текстом отдать за прогон
|
||||||
|
query: список архивов — номера через запятую/пробел или выражение
|
||||||
|
`(repositories.id:N OR ...)`. Пустой — обычный поиск по словам
|
||||||
|
year_from: отсекать статьи старше этого года (необязательно)
|
||||||
|
year_to: отсекать статьи новее этого года (необязательно)
|
||||||
|
resume_token: позиция вида "1298:4200" — архив и смещение в нём
|
||||||
|
progress_cb: см. base.ProgressCallback
|
||||||
|
"""
|
||||||
|
parts = self._repos(query)
|
||||||
|
start_part, start_offset = self._parse_token(resume_token)
|
||||||
|
if start_part in parts:
|
||||||
|
parts = parts[parts.index(start_part):]
|
||||||
|
else:
|
||||||
|
start_offset = 0
|
||||||
|
|
||||||
|
given = 0
|
||||||
|
seen: set[str] = set()
|
||||||
|
for part in parts:
|
||||||
|
offset = start_offset
|
||||||
|
start_offset = 0 # смещение относится только к архиву из токена
|
||||||
|
|
||||||
|
while given < limit and offset < MAX_OFFSET:
|
||||||
|
results = self._page(part, query, offset)
|
||||||
|
if results is None: # API не ответил — прогон закончен
|
||||||
|
return
|
||||||
|
if not results:
|
||||||
|
logger.info("CORE: архив %s исчерпан на смещении %d", part, offset)
|
||||||
|
break
|
||||||
|
|
||||||
|
token = f"{part}:{offset}"
|
||||||
|
for work in results:
|
||||||
|
text = work.get("fullText") or ""
|
||||||
|
if len(text) < MIN_CHARS:
|
||||||
|
continue
|
||||||
|
year = work.get("yearPublished")
|
||||||
|
if year and ((year_from and year < year_from)
|
||||||
|
or (year_to and year > year_to)):
|
||||||
|
continue
|
||||||
|
# Дубли CORE: один текст под разными id (см. докстринг).
|
||||||
|
# Хеша начала текста хватает — совпадение первых 5 тыс.
|
||||||
|
# символов у разных статей практически исключено
|
||||||
|
digest = hashlib.sha1(text[:5000].encode()).hexdigest()
|
||||||
|
if digest in seen:
|
||||||
|
continue
|
||||||
|
seen.add(digest)
|
||||||
|
given += 1
|
||||||
|
# Токен указывает на страницу, из которой пришла статья, а
|
||||||
|
# не на следующую: пачка может прерваться на её середине, и
|
||||||
|
# тогда следующий прогон перечитает страницу целиком.
|
||||||
|
# Лишний запрос дешевле потерянных статей, дубли отсекает
|
||||||
|
# ON CONFLICT на ext_id
|
||||||
|
work["resume_token"] = token
|
||||||
|
yield work
|
||||||
|
if given >= limit:
|
||||||
|
break
|
||||||
|
|
||||||
|
offset += PAGE
|
||||||
|
self.last_token = f"{part}:{offset}"
|
||||||
|
if progress_cb and not progress_cb(given):
|
||||||
|
logger.info("CORE: выборка остановлена по запросу (%d)", given)
|
||||||
|
return
|
||||||
|
|
||||||
|
logger.info("CORE: пройдено архивов: %d, отдано %d", len(parts), given)
|
||||||
|
|
||||||
|
@staticmethod
|
||||||
|
def _repos(query: str) -> list[str]:
|
||||||
|
"""Номера архивов из настройки источника.
|
||||||
|
|
||||||
|
Принимает и выражение `(repositories.id:1298 OR ...)`, и простой список
|
||||||
|
«1298, 21908». Если номеров нет вовсе — единственная часть с именем
|
||||||
|
`q`: тогда парсер просто ищет по словам запроса.
|
||||||
|
"""
|
||||||
|
ids = re.findall(r"repositories\.id:(\d+)", query or "")
|
||||||
|
if not ids:
|
||||||
|
ids = re.findall(r"\b\d{2,7}\b", query or "")
|
||||||
|
return ids or ["q"]
|
||||||
|
|
||||||
|
@staticmethod
|
||||||
|
def _parse_token(token: str | None) -> tuple[str | None, int]:
|
||||||
|
"""Разобрать позицию "архив:смещение"; мусор — начать сначала."""
|
||||||
|
if not token or ":" not in token:
|
||||||
|
return None, 0
|
||||||
|
part, _, offset = token.rpartition(":")
|
||||||
|
try:
|
||||||
|
return part, int(offset)
|
||||||
|
except ValueError:
|
||||||
|
logger.warning("CORE: непонятная позиция %r, начинаем сначала", token)
|
||||||
|
return None, 0
|
||||||
|
|
||||||
|
def _page(self, part: str, query: str, offset: int) -> list[dict[str, Any]] | None:
|
||||||
|
"""Одна страница выдачи; None — API не отвечает, прогон пора кончать.
|
||||||
|
|
||||||
|
В запросе РОВНО одно условие: связка через AND в этом API не работает
|
||||||
|
(см. докстринг модуля), поэтому годы отсекаются уже после выборки.
|
||||||
|
"""
|
||||||
|
q = f"repositories.id:{part}" if part != "q" else (query or "*")
|
||||||
|
params = {"q": q, "limit": PAGE, "offset": offset}
|
||||||
|
|
||||||
|
for attempt in range(RETRIES):
|
||||||
|
try:
|
||||||
|
resp = self.client.get(API_URL, params=params)
|
||||||
|
if resp.status_code == 429:
|
||||||
|
# Лимит тарифа: подождать и повторить, а не ронять прогон
|
||||||
|
wait = int(resp.headers.get("retry-after", 30))
|
||||||
|
logger.warning("CORE: лимит запросов, ждём %dс", wait)
|
||||||
|
time.sleep(min(wait, 60))
|
||||||
|
continue
|
||||||
|
resp.raise_for_status()
|
||||||
|
return resp.json().get("results") or []
|
||||||
|
except Exception as e:
|
||||||
|
logger.warning(
|
||||||
|
"CORE: страница %s:%d не удалась (%d/%d): %s",
|
||||||
|
part, offset, attempt + 1, RETRIES, e,
|
||||||
|
)
|
||||||
|
time.sleep(2 * (attempt + 1))
|
||||||
|
|
||||||
|
logger.error("CORE: страница %s:%d не далась за %d попыток", part, offset, RETRIES)
|
||||||
|
return None
|
||||||
|
|
||||||
|
def transform(self, raw: dict[str, Any]) -> dict[str, Any]:
|
||||||
|
"""Привести статью CORE к унифицированному формату корпуса."""
|
||||||
|
ext = raw.get("id")
|
||||||
|
text = raw.get("fullText") or ""
|
||||||
|
if not ext or len(text) < MIN_CHARS:
|
||||||
|
return {}
|
||||||
|
|
||||||
|
lang = raw.get("language") or {}
|
||||||
|
code = lang.get("code") if isinstance(lang, dict) else lang
|
||||||
|
journals = raw.get("journals") or []
|
||||||
|
journal = journals[0].get("title") if journals and isinstance(journals[0], dict) else None
|
||||||
|
|
||||||
|
return {
|
||||||
|
"source": self.source_name,
|
||||||
|
"ext_id": f"core:{ext}",
|
||||||
|
"title": (raw.get("title") or f"CORE {ext}")[:1000],
|
||||||
|
"authors": self.normalize_authors(self._authors(raw)),
|
||||||
|
"doi": raw.get("doi"),
|
||||||
|
"year": raw.get("yearPublished"),
|
||||||
|
# zz у CORE значит «язык не определён» — лучше пусто, чем мусор
|
||||||
|
"lang": code if code and code != "zz" else None,
|
||||||
|
"journal": journal[:300] if journal else None,
|
||||||
|
"url": raw.get("downloadUrl") or f"https://core.ac.uk/works/{ext}",
|
||||||
|
"abstract": (raw.get("abstract") or text[:2000])[:2000],
|
||||||
|
"text": text,
|
||||||
|
"resume_token": raw.get("resume_token"),
|
||||||
|
}
|
||||||
|
|
||||||
|
@staticmethod
|
||||||
|
def _authors(raw: dict[str, Any]) -> list[dict[str, str]]:
|
||||||
|
"""Авторы CORE приходят одной строкой "Фамилия, Имя Отчество"."""
|
||||||
|
out = []
|
||||||
|
for author in raw.get("authors") or []:
|
||||||
|
name = author.get("name") if isinstance(author, dict) else author
|
||||||
|
if not name:
|
||||||
|
continue
|
||||||
|
last, _, first = str(name).partition(",")
|
||||||
|
out.append({"last_name": last.strip(), "first_name": first.strip()})
|
||||||
|
return out
|
||||||
@@ -14,7 +14,7 @@
|
|||||||
import logging
|
import logging
|
||||||
import re
|
import re
|
||||||
from collections.abc import Iterator
|
from collections.abc import Iterator
|
||||||
from concurrent.futures import ThreadPoolExecutor
|
from concurrent.futures import ThreadPoolExecutor, as_completed
|
||||||
from typing import Any
|
from typing import Any
|
||||||
from urllib.parse import quote
|
from urllib.parse import quote
|
||||||
|
|
||||||
@@ -39,8 +39,11 @@ class PMCBulkParser(BaseParser):
|
|||||||
def __init__(self, workers: int = 12) -> None:
|
def __init__(self, workers: int = 12) -> None:
|
||||||
super().__init__()
|
super().__init__()
|
||||||
self.workers = workers
|
self.workers = workers
|
||||||
|
# 60с на файл при бюджете задачи в 1500с — один залипший запрос съедал
|
||||||
|
# почти весь бюджет. Обычный GET сюда укладывается в доли секунды,
|
||||||
|
# 12с — с большим запасом на джиттер, но без риска съесть весь прогон
|
||||||
self.client = httpx.Client(
|
self.client = httpx.Client(
|
||||||
timeout=60,
|
timeout=12,
|
||||||
headers={"User-Agent": "AcademicHelper/1.0 (noreply@jze9.ru)"},
|
headers={"User-Agent": "AcademicHelper/1.0 (noreply@jze9.ru)"},
|
||||||
)
|
)
|
||||||
# Позиция листинга: бакет отдаётся страницами, и продолжать прогон
|
# Позиция листинга: бакет отдаётся страницами, и продолжать прогон
|
||||||
@@ -82,8 +85,13 @@ class PMCBulkParser(BaseParser):
|
|||||||
logger.info("PMC bulk: бакет закончился")
|
logger.info("PMC bulk: бакет закончился")
|
||||||
return
|
return
|
||||||
|
|
||||||
|
# as_completed вместо map(): map() отдаёт результаты строго по
|
||||||
|
# порядку отправки, поэтому один залипший запрос блокирует все
|
||||||
|
# уже готовые — даже если остальные 11 потоков давно отработали
|
||||||
with ThreadPoolExecutor(max_workers=self.workers) as pool:
|
with ThreadPoolExecutor(max_workers=self.workers) as pool:
|
||||||
for raw in pool.map(self._fetch_article, ids):
|
futures = {pool.submit(self._fetch_article, i): i for i in ids}
|
||||||
|
for future in as_completed(futures):
|
||||||
|
raw = future.result()
|
||||||
if raw is None:
|
if raw is None:
|
||||||
continue
|
continue
|
||||||
given += 1
|
given += 1
|
||||||
|
|||||||
156
scripts/parsers/tests/test_core.py
Normal file
156
scripts/parsers/tests/test_core.py
Normal file
@@ -0,0 +1,156 @@
|
|||||||
|
"""Юнит-тесты парсера CORE — чистая логика, без сети.
|
||||||
|
|
||||||
|
Стерегут то, ради чего парсер написан иначе остальных: отбор статей с полным
|
||||||
|
текстом, отсев дублей CORE (одна статья под разными id), нарезку выборки по
|
||||||
|
архивам-поставщикам и позицию продолжения «архив:смещение».
|
||||||
|
|
||||||
|
Нарезка по архивам — не украшение: связка условий через `AND` в API CORE не
|
||||||
|
работает (запрос по архивам вместе с годом отдаёт годы вперемешку), поэтому в
|
||||||
|
запрос идёт ровно одно условие, а годы отсекаются уже на нашей стороне.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from core import DEFAULT_PROXY, MIN_CHARS, COREParser
|
||||||
|
|
||||||
|
LONG = "слово " * 400 # заведомо длиннее MIN_CHARS
|
||||||
|
|
||||||
|
SAMPLE = {
|
||||||
|
"id": 123456,
|
||||||
|
"title": "Нейронные сети в медицине",
|
||||||
|
"authors": [{"name": "Кузьмин, Ярослав Вадимович"}, {"name": "Савенок А."}],
|
||||||
|
"doi": "10.1234/x",
|
||||||
|
"yearPublished": 2019,
|
||||||
|
"language": {"code": "ru"},
|
||||||
|
"journals": [{"title": "Вестник"}],
|
||||||
|
"downloadUrl": "https://core.ac.uk/download/1.pdf",
|
||||||
|
"abstract": "Аннотация",
|
||||||
|
"fullText": LONG,
|
||||||
|
"resume_token": "1298:100",
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
def test_repos_from_or_expression():
|
||||||
|
q = "(repositories.id:1298 OR repositories.id:21908 OR repositories.id:949)"
|
||||||
|
assert COREParser._repos(q) == ["1298", "21908", "949"]
|
||||||
|
|
||||||
|
|
||||||
|
def test_repos_from_plain_list():
|
||||||
|
assert COREParser._repos("1298, 21908 949") == ["1298", "21908", "949"]
|
||||||
|
|
||||||
|
|
||||||
|
def test_repos_without_numbers_falls_back_to_word_search():
|
||||||
|
# Номеров нет — единственная часть «q»: обычный поиск по словам
|
||||||
|
assert COREParser._repos("нейронные сети") == ["q"]
|
||||||
|
assert COREParser._repos("") == ["q"]
|
||||||
|
|
||||||
|
|
||||||
|
def test_parse_token():
|
||||||
|
assert COREParser._parse_token("1298:4200") == ("1298", 4200)
|
||||||
|
assert COREParser._parse_token("q:300") == ("q", 300)
|
||||||
|
assert COREParser._parse_token(None) == (None, 0)
|
||||||
|
assert COREParser._parse_token("мусор") == (None, 0)
|
||||||
|
assert COREParser._parse_token("архив:смещение") == (None, 0)
|
||||||
|
|
||||||
|
|
||||||
|
def test_transform_maps_fields():
|
||||||
|
t = COREParser(proxy_url="").transform(SAMPLE)
|
||||||
|
assert t["ext_id"] == "core:123456"
|
||||||
|
assert t["source"] == "core"
|
||||||
|
assert t["year"] == 2019 and t["lang"] == "ru"
|
||||||
|
assert t["journal"] == "Вестник"
|
||||||
|
assert t["text"] == LONG
|
||||||
|
assert t["resume_token"] == "1298:100"
|
||||||
|
# Автор приходит одной строкой «Фамилия, Имя Отчество»
|
||||||
|
assert t["authors"][0]["last_name"] == "Кузьмин"
|
||||||
|
assert t["authors"][0]["initials"] == "Я.В."
|
||||||
|
|
||||||
|
|
||||||
|
def test_transform_skips_short_text():
|
||||||
|
assert COREParser(proxy_url="").transform(dict(SAMPLE, fullText="коротко")) == {}
|
||||||
|
assert COREParser(proxy_url="").transform(dict(SAMPLE, fullText=None)) == {}
|
||||||
|
assert COREParser(proxy_url="").transform(dict(SAMPLE, id=None)) == {}
|
||||||
|
|
||||||
|
|
||||||
|
def test_transform_drops_undefined_language():
|
||||||
|
# zz у CORE значит «язык не определён» — в корпус такое класть незачем
|
||||||
|
assert COREParser(proxy_url="").transform(dict(SAMPLE, language={"code": "zz"}))["lang"] is None
|
||||||
|
assert COREParser(proxy_url="").transform(dict(SAMPLE, language=None))["lang"] is None
|
||||||
|
|
||||||
|
|
||||||
|
def _parser_with_pages(pages):
|
||||||
|
"""Парсер, у которого выдача подменена заранее заготовленными страницами."""
|
||||||
|
p = COREParser(api_key="test", proxy_url="")
|
||||||
|
calls = []
|
||||||
|
|
||||||
|
def fake_page(part, query, offset):
|
||||||
|
calls.append((part, offset))
|
||||||
|
return pages.pop(0) if pages else []
|
||||||
|
|
||||||
|
p._page = fake_page # type: ignore[method-assign]
|
||||||
|
p.calls = calls # type: ignore[attr-defined]
|
||||||
|
return p
|
||||||
|
|
||||||
|
|
||||||
|
def test_fetch_drops_core_duplicates():
|
||||||
|
# Одна и та же статья под разными id — CORE так отдаёт всегда
|
||||||
|
page = [
|
||||||
|
{"id": 1, "fullText": LONG},
|
||||||
|
{"id": 2, "fullText": LONG},
|
||||||
|
{"id": 3, "fullText": LONG + "иное"},
|
||||||
|
{"id": 4, "fullText": "коротко"},
|
||||||
|
]
|
||||||
|
p = _parser_with_pages([page])
|
||||||
|
got = list(p.fetch(limit=10, query="repositories.id:1298"))
|
||||||
|
assert [w["id"] for w in got] == [1, 3]
|
||||||
|
|
||||||
|
|
||||||
|
def test_fetch_position_points_at_own_page():
|
||||||
|
# Токен обязан указывать на страницу, откуда пришла статья: пачка может
|
||||||
|
# прерваться на середине, и следующий прогон перечитает её целиком
|
||||||
|
pages = [[{"id": 1, "fullText": LONG}], [{"id": 2, "fullText": LONG + "два"}]]
|
||||||
|
p = _parser_with_pages(pages)
|
||||||
|
got = list(p.fetch(limit=10, query="repositories.id:1298"))
|
||||||
|
assert [w["resume_token"] for w in got] == ["1298:0", "1298:100"]
|
||||||
|
|
||||||
|
|
||||||
|
def test_fetch_walks_archives_and_resumes():
|
||||||
|
p = _parser_with_pages([[], []])
|
||||||
|
list(p.fetch(limit=10, query="(repositories.id:1298 OR repositories.id:949)",
|
||||||
|
resume_token="1298:300"))
|
||||||
|
# Начали с архива из токена и его смещения, пустой архив — переход к следующему
|
||||||
|
assert p.calls == [("1298", 300), ("949", 0)]
|
||||||
|
|
||||||
|
|
||||||
|
def test_fetch_ignores_token_of_unknown_archive():
|
||||||
|
p = _parser_with_pages([[]])
|
||||||
|
list(p.fetch(limit=10, query="repositories.id:1298", resume_token="99999:500"))
|
||||||
|
assert p.calls == [("1298", 0)]
|
||||||
|
|
||||||
|
|
||||||
|
def test_fetch_filters_years_on_our_side():
|
||||||
|
# Годы API не фильтрует, поэтому отсекаем сами — и только если заданы
|
||||||
|
page = [
|
||||||
|
{"id": 1, "fullText": LONG, "yearPublished": 2004},
|
||||||
|
{"id": 2, "fullText": LONG + "два", "yearPublished": 2015},
|
||||||
|
{"id": 3, "fullText": LONG + "три", "yearPublished": 2030},
|
||||||
|
]
|
||||||
|
p = _parser_with_pages([list(page)])
|
||||||
|
got = list(p.fetch(limit=10, query="repositories.id:1298", year_from=2010, year_to=2026))
|
||||||
|
assert [w["id"] for w in got] == [2]
|
||||||
|
|
||||||
|
p2 = _parser_with_pages([list(page)])
|
||||||
|
got2 = list(p2.fetch(limit=10, query="repositories.id:1298"))
|
||||||
|
assert [w["id"] for w in got2] == [1, 2, 3]
|
||||||
|
|
||||||
|
|
||||||
|
def test_min_chars_threshold_is_meaningful():
|
||||||
|
assert MIN_CHARS >= 1000
|
||||||
|
|
||||||
|
|
||||||
|
def test_proxy_choice(monkeypatch):
|
||||||
|
# CORE не отвечает на ключ с российских адресов — по умолчанию идём через
|
||||||
|
# sing-box, но вне compose-сети прокси отключают пустым значением
|
||||||
|
monkeypatch.delenv("CORE_PROXY_URL", raising=False)
|
||||||
|
assert COREParser._proxy(None) == DEFAULT_PROXY
|
||||||
|
monkeypatch.setenv("CORE_PROXY_URL", "socks5://иной:1080")
|
||||||
|
assert COREParser._proxy(None) == "socks5://иной:1080"
|
||||||
|
assert COREParser._proxy("") == "" # явное «без прокси» сильнее окружения
|
||||||
@@ -589,6 +589,17 @@ def _parser_for(source_type: str, cfg: dict[str, Any]) -> tuple[Any, dict[str, A
|
|||||||
"dump_path": query or None, # query = путь к дампу, если задан
|
"dump_path": query or None, # query = путь к дампу, если задан
|
||||||
"start_offset": int(cfg.get("resume_token") or 0),
|
"start_offset": int(cfg.get("resume_token") or 0),
|
||||||
}
|
}
|
||||||
|
elif source_type == "core":
|
||||||
|
# Массовый источник: позиция — "год:смещение". Выборка режется по
|
||||||
|
# годам, потому что offset у CORE упирается в 100 000 (см. core.py)
|
||||||
|
from core import COREParser as P
|
||||||
|
kwargs = {
|
||||||
|
"limit": limit,
|
||||||
|
"query": query,
|
||||||
|
"year_from": cfg.get("year_from"),
|
||||||
|
"year_to": cfg.get("year_to"),
|
||||||
|
"resume_token": cfg.get("resume_token"),
|
||||||
|
}
|
||||||
elif source_type == "pmc_bulk":
|
elif source_type == "pmc_bulk":
|
||||||
from pmc_bulk import PMCBulkParser as P
|
from pmc_bulk import PMCBulkParser as P
|
||||||
kwargs = {
|
kwargs = {
|
||||||
|
|||||||
@@ -12,3 +12,4 @@ langdetect==1.0.9
|
|||||||
beautifulsoup4==4.12.3 # парсер CyberLeninka: детали статьи (fetch_article_details)
|
beautifulsoup4==4.12.3 # парсер CyberLeninka: детали статьи (fetch_article_details)
|
||||||
pydantic-settings==2.2.1
|
pydantic-settings==2.2.1
|
||||||
httpx==0.27.0
|
httpx==0.27.0
|
||||||
|
socksio==1.0.0 # SOCKS5-прокси для httpx (CORE блокирует запросы из РФ, см. scripts/parsers/core.py)
|
||||||
|
|||||||
Reference in New Issue
Block a user