- SQL migrations 053 (scraper_settings table) + 054 (seed global + per-source rows)
- scraper_settings.py: get_scraper_delay() returns max(per_source, global);
in-memory cache TTL=60s; invalidate_cache() for immediate effect on PUT
- avito/cian/n1 scrapers load delay from DB in __init__ (mirrors yandex pattern)
- Admin API: GET /scraper-settings (list all), PUT /scraper-settings/{source}
with cache invalidation on update; CAST(:d AS numeric) per psycopg v3 rules
- 10 unit tests: global>per_source, per_source>global, global=0, DB error fallback,
cache invalidation, API endpoint smoke
130 lines
4.8 KiB
Python
130 lines
4.8 KiB
Python
"""scraper_settings.py — загрузка задержек парсеров из таблицы scraper_settings.
|
||
|
||
Кеш в памяти с TTL 60 секунд — не нагружает БД при каждом запросе.
|
||
invalidate_cache() вызывается из admin API PUT для немедленного применения.
|
||
|
||
Ключевая логика:
|
||
get_scraper_delay(source) = max(per_source_delay, global_delay)
|
||
Строка source='global' задаёт нижнюю планку для ВСЕХ парсеров.
|
||
Если global=0 — используется только per-source значение.
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import logging
|
||
import threading
|
||
import time
|
||
from collections.abc import Generator
|
||
from contextlib import contextmanager
|
||
|
||
from sqlalchemy import text
|
||
from sqlalchemy.orm import Session
|
||
|
||
from app.core.db import SessionLocal
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
# Специальный ключ глобальной задержки (строка в scraper_settings с этим source).
|
||
_GLOBAL_KEY = "global"
|
||
|
||
# TTL кеша в секундах — после истечения перечитывается из БД.
|
||
_CACHE_TTL_SEC = 60.0
|
||
|
||
# Дефолтные задержки если строка в БД отсутствует.
|
||
_DEFAULT_DELAY_BY_SOURCE: dict[str, float] = {
|
||
"avito": 7.0,
|
||
"cian": 5.0,
|
||
"n1": 5.0,
|
||
"yandex": 5.0,
|
||
"domrf": 5.0,
|
||
"rosreestr": 5.0,
|
||
}
|
||
|
||
# Fallback для неизвестного источника.
|
||
_GLOBAL_DEFAULT_DELAY = 5.0
|
||
|
||
# Алиасы источников: ключ → canonical source (можно расширять без изменения вызывающего кода).
|
||
_KEY_ALIASES: dict[str, str] = {}
|
||
|
||
# Кеш: source → (delay_value, timestamp).
|
||
_CACHE: dict[str, tuple[float, float]] = {}
|
||
_CACHE_LOCK = threading.Lock()
|
||
|
||
|
||
@contextmanager
|
||
def _open_session() -> Generator[Session, None, None]:
|
||
"""Открыть сессию БД с гарантированным закрытием."""
|
||
db = SessionLocal()
|
||
try:
|
||
yield db
|
||
finally:
|
||
db.close()
|
||
|
||
|
||
def _get_setting_cached(key: str) -> float:
|
||
"""Кеш-aware чтение одного source из scraper_settings.
|
||
|
||
Возвращает request_delay_sec из БД (с кешем TTL=60s).
|
||
При отсутствии строки — фолбек на _DEFAULT_DELAY_BY_SOURCE или 0.0 для global.
|
||
"""
|
||
now = time.time()
|
||
with _CACHE_LOCK:
|
||
cached = _CACHE.get(key)
|
||
if cached is not None and now - cached[1] < _CACHE_TTL_SEC:
|
||
return cached[0]
|
||
|
||
try:
|
||
with _open_session() as db:
|
||
row = db.execute(
|
||
text("SELECT request_delay_sec FROM scraper_settings WHERE source = :s"),
|
||
{"s": key},
|
||
).first()
|
||
if row is not None:
|
||
value = float(row[0])
|
||
elif key == _GLOBAL_KEY:
|
||
# Глобальная строка ещё не создана — не применяем floor.
|
||
value = 0.0
|
||
else:
|
||
value = _DEFAULT_DELAY_BY_SOURCE.get(key, _GLOBAL_DEFAULT_DELAY)
|
||
except Exception as exc:
|
||
logger.warning(
|
||
"scraper_settings: load failed for %s — using default: %s", key, exc
|
||
)
|
||
value = 0.0 if key == _GLOBAL_KEY else _DEFAULT_DELAY_BY_SOURCE.get(
|
||
key, _GLOBAL_DEFAULT_DELAY
|
||
)
|
||
|
||
with _CACHE_LOCK:
|
||
_CACHE[key] = (value, now)
|
||
return value
|
||
|
||
|
||
def get_scraper_delay(source: str) -> float:
|
||
"""Вернуть эффективную задержку для source = max(per_source, global).
|
||
|
||
Строка source='global' задаёт нижнюю планку для всех парсеров.
|
||
Если global=0 — используется только per-source значение.
|
||
"""
|
||
key = _KEY_ALIASES.get(source, source)
|
||
per_source = _get_setting_cached(key)
|
||
global_value = _get_setting_cached(_GLOBAL_KEY)
|
||
return max(per_source, global_value)
|
||
|
||
|
||
def invalidate_cache(source: str | None = None) -> None:
|
||
"""Сбросить кеш для source (или весь кеш если source=None).
|
||
|
||
Вызывается из admin API PUT /scraper-settings/{source} для немедленного
|
||
применения нового значения без ожидания TTL.
|
||
"""
|
||
with _CACHE_LOCK:
|
||
if source is None:
|
||
_CACHE.clear()
|
||
logger.info("scraper_settings: cache cleared (all sources)")
|
||
else:
|
||
_CACHE.pop(source, None)
|
||
# Если source — алиас, сбрасываем canonical key тоже.
|
||
canonical = _KEY_ALIASES.get(source)
|
||
if canonical:
|
||
_CACHE.pop(canonical, None)
|
||
logger.info("scraper_settings: cache cleared for source=%s", source)
|