gendesign/backend/app/services/site_finder/poi_loader.py
lekss361 dbf46228fb
All checks were successful
Deploy / changes (push) Successful in 9s
Deploy / build-frontend (push) Has been skipped
Deploy / deploy-caddy (push) Has been skipped
Deploy / build-backend (push) Successful in 2m29s
Deploy / build-worker (push) Successful in 4m1s
Deploy / deploy (push) Successful in 1m28s
Deploy / deploy-status (push) Successful in 2s
Deploy / perimeter-smoke (push) Successful in 1m44s
Дробить плитку Overpass только при перегрузке, а не при отказе сети (#3528)
2026-09-15 07:33:56 +00:00

473 lines
26 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""Загрузчик OSM POI из Overpass API для site-finder.
Запускается раз в неделю через Celery beat (регион по умолчанию — ЕКБ, см. DEFAULT_REGION).
Поддерживает фильтр "не старше 2 лет" (требование Максима) — last_osm_edit_date.
Параметризован регионом (REGION_BBOX) — sync_poi_to_db(region=...) может грузить любой
зарегистрированный bbox, не только ЕКБ; большие bbox автоматически режутся на тайлы
(_bbox_tiles), чтобы не упереться в лимиты одного Overpass-запроса.
"""
import asyncio
import json
import logging
import math
from datetime import date, datetime, timedelta
import httpx
from sqlalchemy import text
from app.core.db import SessionLocal
logger = logging.getLogger(__name__)
OVERPASS_URL = "https://overpass-api.de/api/interpreter"
EKB_BBOX = (56.7, 60.5, 56.95, 60.75) # (south, west, north, east)
# Прямоугольники продуктовых ядер по региону — источник загрузки POI больше не зашит
# в одну константу (было: только ЕКБ, блок «что рядом» молчал для остальных регионов).
# Site Finder — независимая половина монорепо со своим окружением/БД и НЕ импортирует
# tradein-mvp (у того свой реестр `app/services/regions.py`), поэтому bbox для Москвы
# продублирован явно, а не через кросс-импорт. Значение — bbox_product_core региона 77
# (tradein-mvp/backend/app/services/regions.py REGIONS[77]), пересчитанное в тот же
# (south, west, north, east) порядок, что EKB_BBOX выше. "ekb" остаётся значением по
# умолчанию ВЕЗДЕ (sync_poi_to_db / fetch_overpass) — существующее weekly-расписание
# (tasks.poi_sync.sync_osm_poi_ekb) не передаёт region и не должно молча сменить город.
REGION_BBOX: dict[str, tuple[float, float, float, float]] = {
"ekb": EKB_BBOX,
"msk": (55.55, 37.30, 55.95, 37.90), # (south, west, north, east)
}
DEFAULT_REGION = "ekb"
# Максимальный размер стороны ОДНОГО Overpass-запроса в градусах. У ЕКБ обе стороны
# bbox — ровно 0.25° (проверенный на практике размер: per-category запрос укладывается
# в timeout:30 без 504). Для региона с большей стороной bbox запрос режется на грид
# тайлов такого же порядка вместо одного большого — иначе на плотном городе (Москва на
# порядок плотнее ЕКБ по числу POI) Overpass либо отдаёт 504, либо (хуже) частично
# посчитанный ответ без явной ошибки, и загрузка молча обрежется. Для ЕКБ (0.25×0.25)
# тайлинг даёт РОВНО один тайл, совпадающий с EKB_BBOX бит-в-бит — поведение дефолтного
# региона не меняется.
MAX_TILE_SIDE_DEG = 0.25
# Маппинг набора OSM-тегов (все теги в кортеже должны совпасть — AND) → нормализованная
# category. Каждая запись — один per-category Overpass-запрос (см. _build_overpass_query);
# несколько записей с ОДИНАКОВЫМ значением category (как у metro_stop ниже) — это "ИЛИ" на
# уровне отдельных HTTP-запросов: элемент, подходящий под любую из альтернативных схем
# разметки, попадёт в категорию.
OSM_CATEGORIES: dict[tuple[tuple[str, str], ...], str] = {
# amenity tags — школы расширены (school/college/university)
(("amenity", "school"),): "school",
(("amenity", "college"),): "school",
(("amenity", "university"),): "school",
(("amenity", "kindergarten"),): "kindergarten",
(("amenity", "pharmacy"),): "pharmacy",
(("amenity", "hospital"),): "hospital",
(("amenity", "clinic"),): "hospital",
# shop tags — supermarket расширен
(("shop", "mall"),): "shop_mall",
(("shop", "supermarket"),): "shop_supermarket",
(("shop", "hypermarket"),): "shop_supermarket",
(("shop", "convenience"),): "shop_small",
(("shop", "bakery"),): "shop_small",
# leisure
(("leisure", "park"),): "park",
# transit
(("railway", "tram_stop"),): "tram_stop",
(("highway", "bus_stop"),): "bus_stop",
# Метро ЕКБ (9 станций, одна линия). Fix (location-index rework): фильтр раньше ловил
# ТОЛЬКО station=subway и подтягивал лишь 5/9 станций — часть станций в OSM размечена
# без ключа "station" вовсе, комбинацией railway=station + subway=yes (альтернативная,
# но распространённая схема разметки метро). Обе схемы — отдельными записями ниже, чтобы
# не терять станции, размеченные любой из них.
(("station", "subway"),): "metro_stop",
(("railway", "station"), ("subway", "yes")): "metro_stop",
}
def _bbox_tiles(
bbox: tuple[float, float, float, float], max_side_deg: float = MAX_TILE_SIDE_DEG
) -> list[tuple[float, float, float, float]]:
"""Разбить bbox на равномерный грид тайлов со стороной ≤ max_side_deg.
(south, west, north, east) → список тайлов того же формата. Для bbox, у которого обе
стороны уже ≤ max_side_deg (текущий EKB_BBOX: 0.25×0.25), возвращает список ровно из
ОДНОГО тайла, идентичного входному bbox — тайлинг не меняет поведение для ЕКБ.
"""
south, west, north, east = bbox
rows = max(1, math.ceil(round((north - south) / max_side_deg, 6)))
cols = max(1, math.ceil(round((east - west) / max_side_deg, 6)))
lat_step = (north - south) / rows
lon_step = (east - west) / cols
tiles = []
for r in range(rows):
for c in range(cols):
tiles.append(
(
south + r * lat_step,
west + c * lon_step,
south + (r + 1) * lat_step,
west + (c + 1) * lon_step,
)
)
return tiles
def _split_bbox_quadrants(
bbox: tuple[float, float, float, float],
) -> list[tuple[float, float, float, float]]:
"""Разбить bbox на 4 равные четверти (2×2) — используется адаптивным ретраем
_fetch_category, когда сам тайл всё равно оказался слишком тяжёлым для Overpass."""
south, west, north, east = bbox
mid_lat = (south + north) / 2
mid_lon = (west + east) / 2
return [
(south, west, mid_lat, mid_lon),
(south, mid_lon, mid_lat, east),
(mid_lat, west, north, mid_lon),
(mid_lat, mid_lon, north, east),
]
# Живой замер 2026-09-13: uniform-тайл 0.25×0.25 (размер ЕКБ) для category=bus_stop в
# Москве отдал 504 Gateway Timeout ОДНИМ тайлом (12 181 bus_stop во всём продуктовом
# ядре — уже больше, чем ВСЕ 14 категорий ЕКБ вместе, 4 850). Единый "правильный" размер
# тайла под все 14 категорий Москвы заранее не подобрать — плотность по городу сильно
# неравномерна (плотный центр / разреженная периферия), а у Overpass нет заголовка с
# "это частичный ответ" — единственный надёжный сигнал перегруза — HTTP-ошибка/таймаут.
# Поэтому вместо фиксированного маленького тайла — АДАПТИВНОЕ дробление: тайл, на
# котором per-category запрос дважды падает, дробится на 4 четверти и каждая
# перезапрашивается рекурсивно (до RECURSIVE_SPLIT_MAX_DEPTH). Для ЕКБ recursion
# НИКОГДА не срабатывает (единственный тайл исторически всегда отвечал 200) — поведение
# дефолтного региона не меняется.
RECURSIVE_SPLIT_MAX_DEPTH = 3
# Дробить тайл имеет смысл ТОЛЬКО когда сервер отказал из-за тяжести запроса:
# 504/429/503 и таймаут чтения — «я не успел посчитать», четверть посчитается.
# Отказ на уровне транспорта (connection refused / network unreachable) про размер
# запроса не говорит ВООБЩЕ: хост нас не принимает, и дробление превращает один
# отказ в 4, 16, 64 повторных стука. Живой случай 15.09.2026: загрузка Москвы
# поймала блокировку overpass-api.de по IP и за три минуты выдала 58 отказов на
# 4 успеха — ровно этот механизм.
_SPLIT_WORTHY_STATUS = frozenset({429, 503, 504})
# Подряд идущие транспортные отказы = хост нас не принимает. Продолжать прогон
# бессмысленно и вредно (углубляем блокировку), поэтому после порога — стоп всего
# прогона с явной ошибкой, а не тихий пропуск категорий.
_MAX_CONSECUTIVE_TRANSPORT_ERRORS = 5
class OverpassUnreachableError(RuntimeError):
"""Overpass отказывает на уровне соединения подряд — прогон остановлен."""
class _RunState:
"""Счётчик подряд идущих транспортных отказов в рамках одного fetch_overpass."""
__slots__ = ("consecutive_transport_errors",)
def __init__(self) -> None:
self.consecutive_transport_errors = 0
def _is_overload(exc: Exception) -> bool:
"""True, если отказ говорит «запрос слишком тяжёлый» (есть смысл дробить)."""
if isinstance(exc, httpx.TimeoutException):
return True
if isinstance(exc, httpx.HTTPStatusError):
return exc.response.status_code in _SPLIT_WORTHY_STATUS
return False
def _is_transport_error(exc: Exception) -> bool:
"""True для отказа на уровне соединения (хост не принимает), не про размер запроса."""
return isinstance(exc, httpx.TransportError) and not isinstance(exc, httpx.TimeoutException)
def _build_overpass_query(
tag_filters: tuple[tuple[str, str], ...], bbox: tuple[float, float, float, float]
) -> str:
"""Запрос для ОДНОЙ комбинации tag=value (обычно один тег, иногда несколько — все AND)
в ОДНОМ тайле bbox (south, west, north, east) — см. _bbox_tiles.
Раньше делали один большой запрос на все 14 категорий — Overpass возвращал
504 Gateway Timeout (запрос слишком тяжёлый). Сплит на per-category даёт
быстрые запросы вместо одного 60+ секундного.
"""
south, west, north, east = bbox
bbox_str = f"({south},{west},{north},{east})"
filt = "".join(f'["{k}"="{v}"]' for k, v in tag_filters)
return f"[out:json][timeout:30];(node{filt}{bbox_str};way{filt}{bbox_str};);out center meta;"
def _classify(tags: dict[str, str]) -> str | None:
"""Определить category из OSM-тегов. None если не соответствует ни одной."""
for tag_filters, cat in OSM_CATEGORIES.items():
if all(tags.get(k) == v for k, v in tag_filters):
return cat
return None
def _tag_filters_desc(tag_filters: tuple[tuple[str, str], ...]) -> str:
return ",".join(f"{k}={v}" for k, v in tag_filters)
async def _fetch_category(
client: httpx.AsyncClient,
tag_filters: tuple[tuple[str, str], ...],
category: str,
bbox: tuple[float, float, float, float],
state: _RunState,
depth: int = 0,
) -> list[dict]:
"""Один per-category Overpass-запрос (для ОДНОГО тайла bbox) с ОДНИМ повтором при
транзиентной ошибке; если тайл падает оба раза — адаптивно дробится на 4 четверти
(см. RECURSIVE_SPLIT_MAX_DEPTH) и перезапрашивается рекурсивно, вместо того чтобы
тихо потерять весь тайл.
Fix (location-index rework, "не потерялись крупные категории"): раньше единственная
неудача (таймаут / 504) на всю неделю обнуляла категорию целиком (следующая попытка —
только на следующем weekly run). Один retry с паузой снимает большую часть транзиентных
сбоев без риска зациклиться (Overpass rate-limit — max 2 concurrent, поэтому не более
2 попыток на тайл ДО дробления).
"""
tag_desc = _tag_filters_desc(tag_filters)
query = _build_overpass_query(tag_filters, bbox)
for attempt in (1, 2):
try:
r = await client.post(OVERPASS_URL, data={"data": query})
r.raise_for_status()
elements: list[dict] = r.json().get("elements", [])
logger.info(
"Overpass: %s (%s) bbox=%s depth=%d%d [attempt %d]",
tag_desc,
category,
bbox,
depth,
len(elements),
attempt,
)
# Привязываем category именно к тому per-category запросу, под который
# элемент реально пришёл. Элемент с двумя целевыми тегами (например
# amenity=pharmacy + shop=supermarket) приходит дважды — каждая копия
# несёт свою category. Иначе _classify по dict-порядку молча терял бы
# вторую категорию при UPSERT по UNIQUE(osm_type, osm_id, category). См. #1372.
state.consecutive_transport_errors = 0
for el in elements:
el["_gd_category"] = category
return elements
except Exception as e:
if _is_transport_error(e):
state.consecutive_transport_errors += 1
if state.consecutive_transport_errors >= _MAX_CONSECUTIVE_TRANSPORT_ERRORS:
raise OverpassUnreachableError(
f"Overpass отказывает на уровне соединения "
f"{state.consecutive_transport_errors} раз подряд ({e}) — прогон "
f"остановлен, чтобы не стучаться в блокирующий хост"
) from e
logger.warning(
"Overpass transport error for %s bbox=%s (подряд %d) — тайл пропущен "
"без дробления: %s",
tag_desc,
bbox,
state.consecutive_transport_errors,
e,
)
return []
state.consecutive_transport_errors = 0
if attempt == 1:
logger.warning("Overpass failed for %s (attempt 1, retrying): %s", tag_desc, e)
await asyncio.sleep(3.0)
continue
if _is_overload(e) and depth < RECURSIVE_SPLIT_MAX_DEPTH:
logger.warning(
"Overpass failed for %s bbox=%s twice — splitting into 4 quadrants "
"(depth %d%d) instead of dropping the tile: %s",
tag_desc,
bbox,
depth,
depth + 1,
e,
)
combined: list[dict] = []
for quadrant in _split_bbox_quadrants(bbox):
combined.extend(
await _fetch_category(
client, tag_filters, category, quadrant, state, depth + 1
)
)
await asyncio.sleep(1.0)
return combined
logger.warning(
"Overpass failed for %s bbox=%s at max split depth %d — tile skipped this run: %s",
tag_desc,
bbox,
depth,
e,
)
return []
async def fetch_overpass(region: str = DEFAULT_REGION) -> list[dict]:
"""Запросить Overpass API per category × per tile, вернуть combined список elements.
Делаем отдельные запросы вместо одного гигантского — большой запрос
отдаёт 504 Gateway Timeout. Между запросами sleep 1с (Overpass usage
policy: max 2 concurrent, лучше 1 req/s). bbox региона режется на тайлы
≤ MAX_TILE_SIDE_DEG (_bbox_tiles) — для "ekb" это ровно один тайл (без
изменения поведения), для регионов с большим bbox (напр. "msk") — несколько,
чтобы не поймать 504 или тихо обрезанный ответ на плотном городе.
Overpass блокирует default `python-httpx/*` User-Agent (406) — поэтому
явный UA с контактом проекта.
"""
bbox = REGION_BBOX[region]
tiles = _bbox_tiles(bbox)
state = _RunState()
headers = {
"User-Agent": "GenDesign-SiteFinder/1.0 (+https://gendsgn.ru)",
"Accept": "application/json",
}
all_elements: list[dict] = []
async with httpx.AsyncClient(timeout=60, headers=headers) as client:
for tag_filters, category in OSM_CATEGORIES.items():
for tile in tiles:
elements = await _fetch_category(client, tag_filters, category, tile, state)
all_elements.extend(elements)
await asyncio.sleep(1.0)
logger.info(
"Overpass region=%s: total %d elements across %d category-queries × %d tiles",
region,
len(all_elements),
len(OSM_CATEGORIES),
len(tiles),
)
return all_elements
def sync_poi_to_db(region: str = DEFAULT_REGION) -> dict[str, int]:
"""Синхронизирует POI из Overpass в osm_poi_ekb для одного региона.
region — ключ REGION_BBOX ("ekb" по умолчанию, сохраняет старое поведение
weekly-расписания). Имя таблицы osm_poi_ekb — историческое (изначально
ЕКБ-only); таблица читается ещё в двух местах вне Site Finder (FDW-таблица
gendesign_osm_poi_ekb + локальное зеркало osm_poi_ekb_local в tradein), поэтому
НЕ переименована: переименование потянуло бы миграции в обеих половинах
монорепо (FDW-объект + зеркало + их индексы) ради косметики.
UPSERT по UNIQUE(osm_type, osm_id, category).
Returns: counters {fetched, inserted, updated, skipped_old}.
"""
elements = asyncio.run(fetch_overpass(region))
# 730 дней ≈ 2 года: избегаем ValueError 29 февраля (year-2 не високосный → нет 29.02).
# Точность ±1 день несущественна для фильтра "не старше 2 лет" (требование Максима). См. #1232.
two_years_ago = date.today() - timedelta(days=730)
inserted = 0
updated = 0
skipped_old = 0
skipped = 0
fetched = len(elements)
db = SessionLocal()
try:
for el in elements:
tags: dict[str, str] = el.get("tags") or {}
# category проставлена в fetch_overpass под тот per-category запрос, под
# который элемент пришёл (multi-tag элемент приходит несколько раз, каждая
# копия со своей category). Fallback на _classify для прямых вызовов. См. #1372.
category = el.get("_gd_category") or _classify(tags)
if not category:
continue
osm_id: int = el["id"]
osm_type: str = el["type"] # 'node' | 'way' | 'relation'
# node имеет lat/lon напрямую, way — через center
if osm_type == "node":
lat = el.get("lat")
lon = el.get("lon")
else:
center: dict = el.get("center") or {}
lat = center.get("lat")
lon = center.get("lon")
if lat is None or lon is None:
continue
# timestamp последней правки элемента
ts: str | None = el.get("timestamp")
last_edit: date | None = None
if ts:
try:
last_edit = datetime.fromisoformat(ts.replace("Z", "+00:00")).date()
except Exception:
last_edit = None
# Мягкий фильтр "не старше 2 лет" — данные сохраняем, помечаем для API
if last_edit and last_edit < two_years_ago:
skipped_old += 1
try:
with db.begin_nested(): # SAVEPOINT — откат только этой записи
result = db.execute(
text("""
INSERT INTO osm_poi_ekb
(osm_id, osm_type, category, name, lat, lon, geom,
last_osm_edit_date, tags, fetched_at)
VALUES (:osm_id, :osm_type, :category, :name, :lat, :lon,
ST_SetSRID(ST_MakePoint(:lon, :lat), 4326),
:last_edit, CAST(:tags AS jsonb), NOW())
ON CONFLICT (osm_type, osm_id, category) DO UPDATE
SET name = EXCLUDED.name,
lat = EXCLUDED.lat,
lon = EXCLUDED.lon,
geom = EXCLUDED.geom,
last_osm_edit_date = EXCLUDED.last_osm_edit_date,
tags = EXCLUDED.tags,
fetched_at = NOW()
RETURNING (xmax = 0) AS is_insert
"""),
{
"osm_id": osm_id,
"osm_type": osm_type,
"category": category,
"name": tags.get("name"),
"lat": lat,
"lon": lon,
"last_edit": last_edit,
"tags": json.dumps(tags, ensure_ascii=False),
},
).scalar()
if result:
inserted += 1
else:
updated += 1
except Exception as e:
# Дефектный OSM-элемент (битый tags-jsonb, нарушение constraint, PostGIS-
# ошибка у way/relation без корректного center) не должен валить весь
# weekly-sync. SAVEPOINT откатывает только эту строку, продолжаем
# с остальными. См. #1343 и .claude/rules/backend.md (SAVEPOINT pattern).
logger.warning(
"poi_sync insert failed for %s/%s (category=%s): %s",
osm_type,
osm_id,
category,
e,
)
skipped += 1
db.commit()
except Exception as e:
db.rollback()
logger.exception("poi_sync: unexpected error, outer tx rolled back: %s", e)
raise
finally:
db.close()
return {
"fetched": fetched,
"inserted": inserted,
"updated": updated,
"skipped_old": skipped_old,
"skipped": skipped,
}