fix(openalex): пауза между страницами и живой backoff — заливка тонула в 429
Массовый запуск показал: лимит вежливого пула 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 <noreply@anthropic.com>
This commit is contained in:
@@ -232,7 +232,7 @@ docker compose -f docker-compose.prod.yml --profile observability up -d promethe
|
|||||||
|
|
||||||
1. **Линт** — `ruff` (весь Python) + `mypy` (чистая доменная логика).
|
1. **Линт** — `ruff` (весь Python) + `mypy` (чистая доменная логика).
|
||||||
Конфиги: [`ruff.toml`](ruff.toml), [`mypy.ini`](mypy.ini).
|
Конфиги: [`ruff.toml`](ruff.toml), [`mypy.ini`](mypy.ini).
|
||||||
2. **Юнит-тесты** — `pytest` по сервисам: 130 тестов на ядро детекции, скоринга,
|
2. **Юнит-тесты** — `pytest` по сервисам: 132 теста на ядро детекции, скоринга,
|
||||||
парсеров, форматирования, OAuth и прогресса заливки, без внешней инфры
|
парсеров, форматирования, OAuth и прогресса заливки, без внешней инфры
|
||||||
(БД/Redis/GPU/Ollama замоканы либо не нужны).
|
(БД/Redis/GPU/Ollama замоканы либо не нужны).
|
||||||
|
|
||||||
@@ -262,7 +262,7 @@ make test-one SVC=worker-gost # тесты одного сервиса
|
|||||||
| Итоговый % плагиата + цитаты | `worker-gpu/app/scoring.py` | 18 |
|
| Итоговый % плагиата + цитаты | `worker-gpu/app/scoring.py` | 18 |
|
||||||
| ГОСТ 7.1 / 7.0.5 | `worker-gost/app/formatters/` | 17 |
|
| ГОСТ 7.1 / 7.0.5 | `worker-gost/app/formatters/` | 17 |
|
||||||
| Список литературы | `worker-gost/app/bibliography.py` | 7 |
|
| Список литературы | `worker-gost/app/bibliography.py` | 7 |
|
||||||
| Парсеры источников (CyberLeninka, PMC, прогресс-колбэк) | `scripts/parsers/` | 17 |
|
| Парсеры источников (CyberLeninka, PMC, прогресс-колбэк) | `scripts/parsers/` | 19 |
|
||||||
| OAuth-ссылки (Google/Яндекс) | `api/app/core/oauth.py` | 6 |
|
| OAuth-ссылки (Google/Яндекс) | `api/app/core/oauth.py` | 6 |
|
||||||
| Прогресс заливки (счётчики, бюджет) | `worker-indexer/app/progress.py` | 7 |
|
| Прогресс заливки (счётчики, бюджет) | `worker-indexer/app/progress.py` | 7 |
|
||||||
| Шкала загрузки источников | `api/app/core/progress.py` | 7 |
|
| Шкала загрузки источников | `api/app/core/progress.py` | 7 |
|
||||||
|
|||||||
@@ -21,6 +21,35 @@ logger = logging.getLogger(__name__)
|
|||||||
|
|
||||||
OPENALEX_API = "https://api.openalex.org"
|
OPENALEX_API = "https://api.openalex.org"
|
||||||
DEFAULT_EMAIL = "noreply@jze9.ru" # Для вежливого агента
|
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):
|
class OpenAlexParser(BaseParser):
|
||||||
@@ -74,6 +103,7 @@ class OpenAlexParser(BaseParser):
|
|||||||
year_to=year_to,
|
year_to=year_to,
|
||||||
type_filter=type_filter,
|
type_filter=type_filter,
|
||||||
open_access_only=open_access_only,
|
open_access_only=open_access_only,
|
||||||
|
progress_cb=progress_cb,
|
||||||
):
|
):
|
||||||
results.extend(page)
|
results.extend(page)
|
||||||
if progress_cb and not progress_cb(min(len(results), limit)):
|
if progress_cb and not progress_cb(min(len(results), limit)):
|
||||||
@@ -93,6 +123,7 @@ class OpenAlexParser(BaseParser):
|
|||||||
year_to: int | None,
|
year_to: int | None,
|
||||||
type_filter: str,
|
type_filter: str,
|
||||||
open_access_only: bool = False,
|
open_access_only: bool = False,
|
||||||
|
progress_cb: ProgressCallback | None = None,
|
||||||
) -> Generator[list[dict], None, None]:
|
) -> Generator[list[dict], None, None]:
|
||||||
"""Cursor-based пагинация OpenAlex."""
|
"""Cursor-based пагинация OpenAlex."""
|
||||||
cursor = "*"
|
cursor = "*"
|
||||||
@@ -155,8 +186,7 @@ class OpenAlexParser(BaseParser):
|
|||||||
if not cursor:
|
if not cursor:
|
||||||
break
|
break
|
||||||
|
|
||||||
# Rate limiting: 10 запросов/сек без ключа
|
time.sleep(RATE_LIMIT_DELAY)
|
||||||
time.sleep(0.1)
|
|
||||||
|
|
||||||
except httpx.HTTPStatusError as e:
|
except httpx.HTTPStatusError as e:
|
||||||
logger.error(f"OpenAlex HTTP ошибка: {e.response.status_code}")
|
logger.error(f"OpenAlex HTTP ошибка: {e.response.status_code}")
|
||||||
@@ -173,7 +203,9 @@ class OpenAlexParser(BaseParser):
|
|||||||
# лимиту освободиться, каждый воркер продлевал блокировку сам.
|
# лимиту освободиться, каждый воркер продлевал блокировку сам.
|
||||||
wait = 60 * (2 ** (rate_limit_retries - 1))
|
wait = 60 * (2 ** (rate_limit_retries - 1))
|
||||||
logger.warning(f"Rate limit! Попытка {rate_limit_retries}/{MAX_RATE_LIMIT_RETRIES}, ждём {wait}с...")
|
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
|
continue
|
||||||
break
|
break
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
|
|||||||
@@ -7,7 +7,9 @@
|
|||||||
|
|
||||||
from typing import Any
|
from typing import Any
|
||||||
|
|
||||||
|
import openalex
|
||||||
from cyberleninka import CyberLeninkaParser
|
from cyberleninka import CyberLeninkaParser
|
||||||
|
from openalex import _sleep_alive
|
||||||
|
|
||||||
|
|
||||||
class FakeResponse:
|
class FakeResponse:
|
||||||
@@ -66,3 +68,26 @@ def test_fetch_works_without_callback(monkeypatch):
|
|||||||
"""Скрипты заливки зовут парсеры без progress_cb — поведение прежнее."""
|
"""Скрипты заливки зовут парсеры без progress_cb — поведение прежнее."""
|
||||||
p = _parser(monkeypatch)
|
p = _parser(monkeypatch)
|
||||||
assert len(p.fetch(query="x", limit=20)) == 20
|
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]
|
||||||
|
|||||||
Reference in New Issue
Block a user