From f1a8d07cb3c6a8de6a04f125628d07a71fd61b30 Mon Sep 17 00:00:00 2001 From: jze9 Date: Thu, 27 Aug 2026 17:54:42 +0500 Subject: [PATCH] =?UTF-8?q?fix(openalex):=20=D0=BF=D0=B0=D1=83=D0=B7=D0=B0?= =?UTF-8?q?=20=D0=BC=D0=B5=D0=B6=D0=B4=D1=83=20=D1=81=D1=82=D1=80=D0=B0?= =?UTF-8?q?=D0=BD=D0=B8=D1=86=D0=B0=D0=BC=D0=B8=20=D0=B8=20=D0=B6=D0=B8?= =?UTF-8?q?=D0=B2=D0=BE=D0=B9=20backoff=20=E2=80=94=20=D0=B7=D0=B0=D0=BB?= =?UTF-8?q?=D0=B8=D0=B2=D0=BA=D0=B0=20=D1=82=D0=BE=D0=BD=D1=83=D0=BB=D0=B0?= =?UTF-8?q?=20=D0=B2=20429?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Массовый запуск показал: лимит вежливого пула OpenAlex (10 req/s) общий на mailto, а не на процесс. Четыре воркера с паузой 0.1с получали сплошные 429, и каждый прогон уходил в 900с бесполезного backoff, не забрав ничего. - RATE_LIMIT_DELAY 0.1 → 1.0с (≈4 req/s на четырёх воркерах); - backoff спит кусками по 5с и отчитывается через progress_cb: прогон больше не выглядит зависшим в админке и отменяется во время ожидания, а не после. Co-Authored-By: Claude Opus 5 --- README.md | 4 +-- scripts/parsers/openalex.py | 38 +++++++++++++++++++++-- scripts/parsers/tests/test_progress_cb.py | 25 +++++++++++++++ 3 files changed, 62 insertions(+), 5 deletions(-) diff --git a/README.md b/README.md index 47980e0..be095e4 100644 --- a/README.md +++ b/README.md @@ -232,7 +232,7 @@ docker compose -f docker-compose.prod.yml --profile observability up -d promethe 1. **Линт** — `ruff` (весь Python) + `mypy` (чистая доменная логика). Конфиги: [`ruff.toml`](ruff.toml), [`mypy.ini`](mypy.ini). -2. **Юнит-тесты** — `pytest` по сервисам: 130 тестов на ядро детекции, скоринга, +2. **Юнит-тесты** — `pytest` по сервисам: 132 теста на ядро детекции, скоринга, парсеров, форматирования, OAuth и прогресса заливки, без внешней инфры (БД/Redis/GPU/Ollama замоканы либо не нужны). @@ -262,7 +262,7 @@ make test-one SVC=worker-gost # тесты одного сервиса | Итоговый % плагиата + цитаты | `worker-gpu/app/scoring.py` | 18 | | ГОСТ 7.1 / 7.0.5 | `worker-gost/app/formatters/` | 17 | | Список литературы | `worker-gost/app/bibliography.py` | 7 | -| Парсеры источников (CyberLeninka, PMC, прогресс-колбэк) | `scripts/parsers/` | 17 | +| Парсеры источников (CyberLeninka, PMC, прогресс-колбэк) | `scripts/parsers/` | 19 | | OAuth-ссылки (Google/Яндекс) | `api/app/core/oauth.py` | 6 | | Прогресс заливки (счётчики, бюджет) | `worker-indexer/app/progress.py` | 7 | | Шкала загрузки источников | `api/app/core/progress.py` | 7 | diff --git a/scripts/parsers/openalex.py b/scripts/parsers/openalex.py index 7c1edc4..e71a8d6 100644 --- a/scripts/parsers/openalex.py +++ b/scripts/parsers/openalex.py @@ -21,6 +21,35 @@ logger = logging.getLogger(__name__) OPENALEX_API = "https://api.openalex.org" DEFAULT_EMAIL = "noreply@jze9.ru" # Для вежливого агента +# Пауза между страницами. Лимит вежливого пула — 10 запросов/сек НА ВЕСЬ ключ +# (mailto), а не на процесс: четыре воркера, качающие разные источники, делят +# его между собой. С прежними 0.1с массовая заливка утыкалась в сплошные 429 и +# каждый прогон уходил в 900с бесполезного backoff. 1с × 4 воркера ≈ 4 req/s. +RATE_LIMIT_DELAY = 1.0 +# Backoff режем на куски: во время сна парсер обязан отчитываться о жизни, +# иначе прогон выглядит зависшим и его нельзя отменить из админки. +BACKOFF_TICK_S = 5.0 + + +def _sleep_alive( + seconds: float, progress_cb: ProgressCallback | None, fetched: int +) -> bool: + """Поспать, отчитываясь о жизни; False — попросили остановиться. + + Минуты сна в backoff нельзя проводить молча: для админки такой прогон + неотличим от зависшего, а отмена не сработает до конца ожидания. + """ + if progress_cb is None: + time.sleep(seconds) + return True + + left = seconds + while left > 0: + time.sleep(min(BACKOFF_TICK_S, left)) + left -= BACKOFF_TICK_S + if not progress_cb(fetched): + return False + return True class OpenAlexParser(BaseParser): @@ -74,6 +103,7 @@ class OpenAlexParser(BaseParser): year_to=year_to, type_filter=type_filter, open_access_only=open_access_only, + progress_cb=progress_cb, ): results.extend(page) if progress_cb and not progress_cb(min(len(results), limit)): @@ -93,6 +123,7 @@ class OpenAlexParser(BaseParser): year_to: int | None, type_filter: str, open_access_only: bool = False, + progress_cb: ProgressCallback | None = None, ) -> Generator[list[dict], None, None]: """Cursor-based пагинация OpenAlex.""" cursor = "*" @@ -155,8 +186,7 @@ class OpenAlexParser(BaseParser): if not cursor: break - # Rate limiting: 10 запросов/сек без ключа - time.sleep(0.1) + time.sleep(RATE_LIMIT_DELAY) except httpx.HTTPStatusError as e: logger.error(f"OpenAlex HTTP ошибка: {e.response.status_code}") @@ -173,7 +203,9 @@ class OpenAlexParser(BaseParser): # лимиту освободиться, каждый воркер продлевал блокировку сам. wait = 60 * (2 ** (rate_limit_retries - 1)) logger.warning(f"Rate limit! Попытка {rate_limit_retries}/{MAX_RATE_LIMIT_RETRIES}, ждём {wait}с...") - time.sleep(wait) + if not _sleep_alive(wait, progress_cb, total_fetched): + logger.info("OpenAlex: ожидание прервано по запросу") + break continue break except Exception as e: diff --git a/scripts/parsers/tests/test_progress_cb.py b/scripts/parsers/tests/test_progress_cb.py index 01a5052..93a925b 100644 --- a/scripts/parsers/tests/test_progress_cb.py +++ b/scripts/parsers/tests/test_progress_cb.py @@ -7,7 +7,9 @@ from typing import Any +import openalex from cyberleninka import CyberLeninkaParser +from openalex import _sleep_alive class FakeResponse: @@ -66,3 +68,26 @@ def test_fetch_works_without_callback(monkeypatch): """Скрипты заливки зовут парсеры без progress_cb — поведение прежнее.""" p = _parser(monkeypatch) assert len(p.fetch(query="x", limit=20)) == 20 + + +def test_backoff_sleep_reports_life_and_can_be_interrupted(monkeypatch): + """Минуты ожидания в backoff OpenAlex — не молчание: тик и шанс остановиться.""" + slept: list[float] = [] + monkeypatch.setattr(openalex.time, "sleep", lambda s: slept.append(s)) + monkeypatch.setattr(openalex, "BACKOFF_TICK_S", 5.0) + + ticks: list[int] = [] + assert _sleep_alive(20, lambda n: ticks.append(n) or True, fetched=7) is True + assert sum(slept) == 20 and ticks == [7, 7, 7, 7] + + slept.clear() + # Останавливаемся на первом же тике — не досыпая оставшиеся 900с + assert _sleep_alive(900, lambda n: False, fetched=7) is False + assert sum(slept) == 5 + + +def test_backoff_sleep_without_callback_just_sleeps(monkeypatch): + slept: list[float] = [] + monkeypatch.setattr(openalex.time, "sleep", lambda s: slept.append(s)) + assert _sleep_alive(60, None, fetched=0) is True + assert slept == [60]