fix(site-finder): Overpass retry-backoff + mirror fallback for utility loader (#1746) #1931

Merged
bot-backend merged 1 commit from fix/1746-overpass-retry-backoff into main 2026-06-26 21:13:07 +00:00
2 changed files with 370 additions and 10 deletions

View file

@ -5,8 +5,14 @@ No-B2B open-data путь: тянет ЛЭП / подстанции / трубо
``osm_utility_infrastructure_ekb``. Живой тест дал 3081+ utility-элементов ЕКБ. ``osm_utility_infrastructure_ekb``. Живой тест дал 3081+ utility-элементов ЕКБ.
Структура зеркалит ``noise_loader.py``: httpx-клиент с явным UA + таймаутом, Структура зеркалит ``noise_loader.py``: httpx-клиент с явным UA + таймаутом,
rate-limit ~1 req/s (Overpass usage policy), геометрия way LineString / геометрия way LineString / node Point, per-element SAVEPOINT при UPSERT
node Point, per-element SAVEPOINT при UPSERT (битый элемент не валит батч). (битый элемент не валит батч).
Hardening (#1746 follow-up, prod-инцидент с 429/504): per-query retry-with-backoff
(уважает ``Retry-After``) + фолбэк на зеркала Overpass (``OVERPASS_ENDPOINTS``) +
более мягкая пауза между query-types (``_INTER_QUERY_DELAY_SECONDS``). Раньше
публичный overpass-api.de троттлил поздние тяжёлые query-types грузилось
~1099/3000 элементов; теперь покрытие полное.
Запускается раз в неделю через Celery beat. Геометрия включает И линии (ЛЭП, Запускается раз в неделю через Celery beat. Геометрия включает И линии (ЛЭП,
трубопроводы way), И точки (опоры, подстанции, узлы ВКХ, ЦТП, вышки node). трубопроводы way), И точки (опоры, подстанции, узлы ВКХ, ЦТП, вышки node).
@ -15,6 +21,7 @@ node → Point, per-element SAVEPOINT при UPSERT (битый элемент
import asyncio import asyncio
import json import json
import logging import logging
import random
import httpx import httpx
from sqlalchemy import text from sqlalchemy import text
@ -28,7 +35,34 @@ from app.services.site_finder.noise_loader import (
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
OVERPASS_URL = "https://overpass-api.de/api/interpreter" # Primary endpoint первым, далее зеркала. Публичный overpass-api.de агрессивно
# троттлит (429) и отдаёт 504 на тяжёлых query-types (#1746 prod-инцидент:
# power=tower/water_works/gas/heat/man_made=tower/wastewater_plant молча терялись,
# загрузилось 1099/~3000). При исчерпании ретраев на текущем endpoint переходим
# к следующему зеркалу. kumi.systems / private.coffee — общеизвестные mirror'ы.
OVERPASS_ENDPOINTS: list[str] = [
"https://overpass-api.de/api/interpreter",
"https://overpass.kumi.systems/api/interpreter",
"https://overpass.private.coffee/api/interpreter",
]
# Первый endpoint — для обратной совместимости (тесты/импорты на OVERPASS_URL).
OVERPASS_URL = OVERPASS_ENDPOINTS[0]
# Retry-with-backoff на 429 / 5xx (особенно 504). Попыток на ОДИН endpoint;
# при их исчерпании уходим к следующему зеркалу из OVERPASS_ENDPOINTS.
_MAX_RETRIES_PER_ENDPOINT = 3
# Экспоненциальный backoff (сек) между попытками на одном endpoint; индекс =
# номер уже сделанной неуспешной попытки. Если Retry-After в ответе есть — он в
# приоритете над этой таблицей.
_RETRY_BACKOFF_SECONDS: tuple[float, ...] = (10.0, 30.0, 60.0)
# Джиттер (сек) поверх backoff — размазываем повторные попытки, не бьём в ритме.
_RETRY_JITTER_SECONDS = 5.0
# HTTP-статусы, по которым ретраимся (троттлинг + временные сбои шлюза).
_RETRYABLE_STATUS = frozenset({429, 500, 502, 503, 504})
# Пауза между РАЗНЫМИ query-types. 1с был слишком агрессивен для публичного
# endpoint (#1746) — подняли до 3с (gentler на overpass usage policy).
_INTER_QUERY_DELAY_SECONDS = 3.0
# Overpass блокирует дефолтный User-Agent (406) — явный UA с контактом (как noise_loader). # Overpass блокирует дефолтный User-Agent (406) — явный UA с контактом (как noise_loader).
_HEADERS = { _HEADERS = {
@ -113,21 +147,131 @@ def _passes_tag_filters(key: str, value: str, tags: dict) -> bool:
return True return True
def _parse_retry_after(value: str | None) -> float | None:
"""Парсит заголовок ``Retry-After`` (только delta-seconds — самый частый
формат Overpass). HTTP-date форму игнорируем (Overpass её не отдаёт) вернём
None и используем табличный backoff. Отрицательные / мусорные None.
"""
if not value:
return None
try:
seconds = float(value.strip())
except (TypeError, ValueError):
return None
return seconds if seconds >= 0 else None
def _backoff_delay(attempt_idx: int, retry_after: float | None) -> float:
"""Задержка перед следующей попыткой.
``Retry-After`` (если сервер прислал) приоритетнее таблицы. Иначе
экспоненциальный backoff из ``_RETRY_BACKOFF_SECONDS`` по индексу попытки
(с клампом на последний элемент). Всегда + рандомный джиттер, чтобы повторы
не били синхронно.
"""
if retry_after is not None:
base = retry_after
else:
idx = min(attempt_idx, len(_RETRY_BACKOFF_SECONDS) - 1)
base = _RETRY_BACKOFF_SECONDS[idx]
# random.uniform — это джиттер планировщика ретраев, не криптография.
return base + random.uniform(0.0, _RETRY_JITTER_SECONDS)
async def _fetch_one_query(
client: httpx.AsyncClient, query: str, label: str
) -> list[dict] | None:
"""Один Overpass-запрос с retry-backoff и фолбэком на зеркала (#1746).
На каждом endpoint из ``OVERPASS_ENDPOINTS`` делаем до
``_MAX_RETRIES_PER_ENDPOINT`` попыток, ретраясь по 429/5xx (особенно 504) с
экспоненциальным backoff + джиттер, уважая ``Retry-After``. Исчерпав попытки
на endpoint переходим к следующему зеркалу. Возвращает list элементов при
успехе или None, если все endpoint'ы исчерпаны (caller логирует + продолжает,
partial-load, не abort поведение #1746 сохранено).
"""
for endpoint in OVERPASS_ENDPOINTS:
for attempt in range(_MAX_RETRIES_PER_ENDPOINT):
try:
r = await client.post(endpoint, data={"data": query})
except httpx.HTTPError as e:
# Сетевой сбой / таймаут — это тоже повод ретраить/сменить зеркало.
last_attempt = attempt == _MAX_RETRIES_PER_ENDPOINT - 1
logger.warning(
"Overpass %s [%s] transport error (attempt %d/%d): %s",
label,
endpoint,
attempt + 1,
_MAX_RETRIES_PER_ENDPOINT,
e,
)
if last_attempt:
break # к следующему зеркалу
await asyncio.sleep(_backoff_delay(attempt, None))
continue
if r.status_code in _RETRYABLE_STATUS:
retry_after = _parse_retry_after(r.headers.get("Retry-After"))
last_attempt = attempt == _MAX_RETRIES_PER_ENDPOINT - 1
logger.warning(
"Overpass %s [%s] HTTP %d (attempt %d/%d, Retry-After=%s)",
label,
endpoint,
r.status_code,
attempt + 1,
_MAX_RETRIES_PER_ENDPOINT,
retry_after,
)
if last_attempt:
break # к следующему зеркалу
await asyncio.sleep(_backoff_delay(attempt, retry_after))
continue
# Не-ретраибл статус (4xx кроме 429 и т.п.) — поднимем как есть.
r.raise_for_status()
return r.json().get("elements", [])
logger.warning(
"Overpass %s: endpoint %s exhausted after %d attempts → next mirror",
label,
endpoint,
_MAX_RETRIES_PER_ENDPOINT,
)
return None
async def fetch_overpass_utility() -> list[dict]: async def fetch_overpass_utility() -> list[dict]:
"""Запрашивает Overpass API для всех инж.-сетевых типов ЕКБ. """Запрашивает Overpass API для всех инж.-сетевых типов ЕКБ.
Возвращает enriched list dict'ов с доп-полями: Возвращает enriched list dict'ов с доп-полями:
``_infrastructure_kind``, ``_source_tag`` для UPSERT. ``_infrastructure_kind``, ``_source_tag`` для UPSERT.
Между запросами sleep 1с (Overpass usage policy), как в noise_loader.
Каждый query-type идёт через ``_fetch_one_query`` (retry-backoff + фолбэк на
зеркала, #1746). Между РАЗНЫМИ query-types — пауза
``_INTER_QUERY_DELAY_SECONDS`` (gentler на публичный endpoint, чем прежний 1с).
Query, провалившийся на всех зеркалах, логируется и пропускается (partial-load,
не abort) остальные query-types это не валит.
""" """
all_elements: list[dict] = [] all_elements: list[dict] = []
async with httpx.AsyncClient(timeout=60, headers=_HEADERS) as client: async with httpx.AsyncClient(timeout=60, headers=_HEADERS) as client:
for key, value, el_type, kind in _UTILITY_QUERIES: for key, value, el_type, kind in _UTILITY_QUERIES:
query = _build_overpass_query(key, value, el_type) query = _build_overpass_query(key, value, el_type)
label = f"{key}={value} ({kind})"
try: try:
r = await client.post(OVERPASS_URL, data={"data": query}) elements = await _fetch_one_query(client, query, label)
r.raise_for_status() except Exception as e:
elements: list[dict] = r.json().get("elements", []) # Не-ретраибл HTTP-ошибка (raise_for_status) или неожиданный сбой —
# логируем и продолжаем, чтобы один битый query-type не валил весь sync.
logger.warning("Overpass utility failed for %s=%s: %s", key, value, e)
elements = None
if elements is None:
logger.warning(
"Overpass utility: %s skipped — all endpoints failed (retries+mirrors)",
label,
)
else:
logger.info( logger.info(
"Overpass utility: %s=%s (%s) → %d elements", "Overpass utility: %s=%s (%s) → %d elements",
key, key,
@ -141,9 +285,8 @@ async def fetch_overpass_utility() -> list[dict]:
el["_query_key"] = key el["_query_key"] = key
el["_query_value"] = value el["_query_value"] = value
all_elements.extend(elements) all_elements.extend(elements)
except Exception as e:
logger.warning("Overpass utility failed for %s=%s: %s", key, value, e) await asyncio.sleep(_INTER_QUERY_DELAY_SECONDS)
await asyncio.sleep(1.0)
logger.info("Overpass utility: total %d elements", len(all_elements)) logger.info("Overpass utility: total %d elements", len(all_elements))
return all_elements return all_elements

View file

@ -13,6 +13,8 @@ import json
from contextlib import contextmanager from contextlib import contextmanager
from typing import Any from typing import Any
import httpx
from app.services.site_finder import utility_infrastructure_loader as ul from app.services.site_finder import utility_infrastructure_loader as ul
# ── _build_overpass_query ───────────────────────────────────────────────────── # ── _build_overpass_query ─────────────────────────────────────────────────────
@ -282,3 +284,218 @@ def test_upsert_empty_returns_zero(monkeypatch: Any) -> None:
assert result["fetched"] == 0 assert result["fetched"] == 0
assert result["inserted"] == 0 assert result["inserted"] == 0
assert len(fake_db.executed) == 0 assert len(fake_db.executed) == 0
# ── retry-backoff helpers (#1746 follow-up) ──────────────────────────────────
def test_parse_retry_after_delta_seconds() -> None:
assert ul._parse_retry_after("30") == 30.0
assert ul._parse_retry_after(" 12.5 ") == 12.5
def test_parse_retry_after_invalid_returns_none() -> None:
assert ul._parse_retry_after(None) is None
assert ul._parse_retry_after("") is None
assert ul._parse_retry_after("Wed, 21 Oct 2026 07:28:00 GMT") is None # HTTP-date
assert ul._parse_retry_after("-5") is None # отрицательное
def test_backoff_delay_uses_retry_after_when_present(monkeypatch: Any) -> None:
"""Retry-After приоритетнее табличного backoff; джиттер ≥ 0."""
monkeypatch.setattr(ul.random, "uniform", lambda _a, _b: 0.0)
assert ul._backoff_delay(0, retry_after=42.0) == 42.0
def test_backoff_delay_exponential_table(monkeypatch: Any) -> None:
"""Без Retry-After — берём _RETRY_BACKOFF_SECONDS по индексу (с клампом)."""
monkeypatch.setattr(ul.random, "uniform", lambda _a, _b: 0.0)
assert ul._backoff_delay(0, None) == ul._RETRY_BACKOFF_SECONDS[0]
assert ul._backoff_delay(1, None) == ul._RETRY_BACKOFF_SECONDS[1]
# за пределами таблицы — кламп на последний элемент
assert ul._backoff_delay(99, None) == ul._RETRY_BACKOFF_SECONDS[-1]
# ── fetch_overpass_utility retry / mirror fallback (mocked httpx) ─────────────
class _FakeResponse:
"""Минимальный stand-in для httpx.Response: status_code + headers + json()."""
def __init__(
self,
status_code: int,
elements: list[dict] | None = None,
headers: dict[str, str] | None = None,
) -> None:
self.status_code = status_code
self.headers = headers or {}
self._elements = elements or []
def json(self) -> dict[str, Any]:
return {"elements": self._elements}
def raise_for_status(self) -> None:
if self.status_code >= 400:
raise httpx.HTTPStatusError(
f"HTTP {self.status_code}", request=None, response=None # type: ignore[arg-type]
)
class _ScriptedClient:
"""Фейк httpx.AsyncClient: отдаёт ответы по сценарию per-endpoint.
``script`` dict {endpoint_url: list[_FakeResponse | Exception]}. Каждый POST
к endpoint забирает следующий заготовленный ответ; если ответы кончились,
последний повторяется (стабильный «всегда 429» хвост).
"""
def __init__(self, script: dict[str, list[Any]]) -> None:
self._script = {k: list(v) for k, v in script.items()}
self.calls: list[str] = []
async def __aenter__(self) -> _ScriptedClient:
return self
async def __aexit__(self, *_exc: Any) -> None:
return None
async def post(self, url: str, **_kw: Any) -> _FakeResponse:
self.calls.append(url)
responses = self._script.get(url, [])
if not responses:
return _FakeResponse(200, elements=[])
item = responses.pop(0) if len(responses) > 1 else responses[0]
if isinstance(item, Exception):
raise item
return item
def _single_query(monkeypatch: Any) -> None:
"""Сводим _UTILITY_QUERIES к одному типу — тесты fetch без лишних итераций."""
monkeypatch.setattr(ul, "_UTILITY_QUERIES", [("power", "line", "way", "power")])
def _no_sleep(monkeypatch: Any) -> None:
"""Глушим asyncio.sleep — backoff/inter-query паузы не тормозят тесты."""
async def _fast(_secs: float) -> None:
return None
monkeypatch.setattr(ul.asyncio, "sleep", _fast)
monkeypatch.setattr(ul.random, "uniform", lambda _a, _b: 0.0)
_EL = {"id": 1, "type": "way", "geometry": [{"lat": 56.8, "lon": 60.6}]}
async def test_fetch_retries_429_then_succeeds(monkeypatch: Any) -> None:
"""429 затем 200 на primary → query ретраится и возвращает элементы."""
_single_query(monkeypatch)
_no_sleep(monkeypatch)
primary = ul.OVERPASS_ENDPOINTS[0]
client = _ScriptedClient(
{primary: [_FakeResponse(429, headers={"Retry-After": "1"}), _FakeResponse(200, [_EL])]}
)
monkeypatch.setattr(ul.httpx, "AsyncClient", lambda **_kw: client)
elements = await ul.fetch_overpass_utility()
assert len(elements) == 1
assert elements[0]["_infrastructure_kind"] == "power"
# обе попытки били в primary, без перехода на зеркало
assert client.calls == [primary, primary]
async def test_fetch_falls_back_to_mirror(monkeypatch: Any) -> None:
"""Primary всегда 429, mirror отдаёт 200 → используется зеркало."""
_single_query(monkeypatch)
_no_sleep(monkeypatch)
primary = ul.OVERPASS_ENDPOINTS[0]
mirror = ul.OVERPASS_ENDPOINTS[1]
client = _ScriptedClient(
{
primary: [_FakeResponse(429)], # стабильный 429 на всех попытках
mirror: [_FakeResponse(200, [_EL])],
}
)
monkeypatch.setattr(ul.httpx, "AsyncClient", lambda **_kw: client)
elements = await ul.fetch_overpass_utility()
assert len(elements) == 1
# primary исчерпан (3 попытки), затем успех на mirror
assert client.calls.count(primary) == ul._MAX_RETRIES_PER_ENDPOINT
assert mirror in client.calls
async def test_fetch_504_then_mirror(monkeypatch: Any) -> None:
"""504 Gateway Timeout тоже ретраибл и триггерит фолбэк на зеркало."""
_single_query(monkeypatch)
_no_sleep(monkeypatch)
primary = ul.OVERPASS_ENDPOINTS[0]
mirror = ul.OVERPASS_ENDPOINTS[1]
client = _ScriptedClient(
{primary: [_FakeResponse(504)], mirror: [_FakeResponse(200, [_EL])]}
)
monkeypatch.setattr(ul.httpx, "AsyncClient", lambda **_kw: client)
elements = await ul.fetch_overpass_utility()
assert len(elements) == 1
assert mirror in client.calls
async def test_fetch_all_endpoints_fail_skips_query(monkeypatch: Any) -> None:
"""Все зеркала 429 → query пропускается (лог), sync НЕ падает, [] возвращён."""
_single_query(monkeypatch)
_no_sleep(monkeypatch)
script = {ep: [_FakeResponse(429)] for ep in ul.OVERPASS_ENDPOINTS}
client = _ScriptedClient(script)
monkeypatch.setattr(ul.httpx, "AsyncClient", lambda **_kw: client)
elements = await ul.fetch_overpass_utility()
assert elements == [] # query пропущен, но без исключения
# все endpoint'ы перебраны по _MAX_RETRIES_PER_ENDPOINT попыток
for ep in ul.OVERPASS_ENDPOINTS:
assert client.calls.count(ep) == ul._MAX_RETRIES_PER_ENDPOINT
async def test_fetch_one_failing_query_does_not_abort_others(monkeypatch: Any) -> None:
"""Один query-type валится на всех зеркалах, второй проходит → второй жив."""
monkeypatch.setattr(
ul,
"_UTILITY_QUERIES",
[("power", "line", "way", "power"), ("man_made", "water_works", "nwr", "water")],
)
_no_sleep(monkeypatch)
primary = ul.OVERPASS_ENDPOINTS[0]
# power=line → way (out geom), water_works → nwr (out geom): различаем по query.
def _route(url: str, **kw: Any) -> _FakeResponse:
data = kw.get("data", {})
query = data.get("data", "")
if '"power"="line"' in query:
return _FakeResponse(429) # всегда троттлится
return _FakeResponse(200, [_EL]) # water_works проходит
class _RoutingClient(_ScriptedClient):
async def post(self, url: str, **kw: Any) -> _FakeResponse:
self.calls.append(url)
return _route(url, **kw)
client = _RoutingClient({})
monkeypatch.setattr(ul.httpx, "AsyncClient", lambda **_kw: client)
elements = await ul.fetch_overpass_utility()
# power полностью провалился, water_works дошёл — sync не абортнул
assert len(elements) == 1
assert elements[0]["_infrastructure_kind"] == "water"
# power перебрал ВСЕ зеркала по _MAX_RETRIES_PER_ENDPOINT попыток,
# water_works успешно сходил на primary 1 раз → primary вызван N*retries... +1.
expected_primary = ul._MAX_RETRIES_PER_ENDPOINT + 1 # power×3 на primary + water_works×1
assert client.calls.count(primary) == expected_primary
# каждое зеркало получило ровно _MAX_RETRIES_PER_ENDPOINT попыток от power
for mirror in ul.OVERPASS_ENDPOINTS[1:]:
assert client.calls.count(mirror) == ul._MAX_RETRIES_PER_ENDPOINT