gendesign/backend/app/services/site_finder/poi_loader.py
lekss361 580be61914
All checks were successful
Deploy / build-backend (push) Successful in 6m30s
Deploy / build-worker (push) Successful in 6m44s
Deploy / changes (push) Successful in 11s
Deploy Trade-In / changes (push) Successful in 15s
Deploy Trade-In / build-backend (push) Successful in 1m19s
Deploy / build-frontend (push) Has been skipped
Deploy / deploy (push) Successful in 2m26s
Deploy Trade-In / deploy (push) Successful in 3m10s
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / build-frontend (push) Successful in 3m6s
Deploy Trade-In / test (push) Successful in 5m5s
fix(tradein/location): заменить сломанный коэффициент локации на калиброванный индекс (#2531)
2026-07-26 21:48:15 +00:00

269 lines
13 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. Поддерживает фильтр
"не старше 2 лет" (требование Максима) — last_osm_edit_date.
"""
import asyncio
import json
import logging
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)
# Маппинг набора 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 _build_overpass_query(tag_filters: tuple[tuple[str, str], ...]) -> str:
"""Запрос для ОДНОЙ комбинации tag=value (обычно один тег, иногда несколько — все AND).
Раньше делали один большой запрос на все 14 категорий — Overpass возвращал
504 Gateway Timeout (запрос слишком тяжёлый). Сплит на per-category даёт
быстрые запросы вместо одного 60+ секундного.
"""
south, west, north, east = EKB_BBOX
bbox = 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};way{filt}{bbox};);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
) -> list[dict]:
"""Один per-category Overpass-запрос с ОДНИМ повтором при транзиентной ошибке.
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)
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) → %d [attempt %d]", tag_desc, category, len(elements), attempt
)
# Привязываем category именно к тому per-category запросу, под который
# элемент реально пришёл. Элемент с двумя целевыми тегами (например
# amenity=pharmacy + shop=supermarket) приходит дважды — каждая копия
# несёт свою category. Иначе _classify по dict-порядку молча терял бы
# вторую категорию при UPSERT по UNIQUE(osm_type, osm_id, category). См. #1372.
for el in elements:
el["_gd_category"] = category
return elements
except Exception as e:
if attempt == 1:
logger.warning("Overpass failed for %s (attempt 1, retrying): %s", tag_desc, e)
await asyncio.sleep(3.0)
continue
logger.warning(
"Overpass failed for %s after retry — category skipped this run: %s", tag_desc, e
)
return []
async def fetch_overpass() -> list[dict]:
"""Запросить Overpass API per category, вернуть combined список elements.
Делаем отдельные запросы вместо одного гигантского — большой запрос
отдаёт 504 Gateway Timeout. Между запросами sleep 1с (Overpass usage
policy: max 2 concurrent, лучше 1 req/s).
Overpass блокирует default `python-httpx/*` User-Agent (406) — поэтому
явный UA с контактом проекта.
"""
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():
elements = await _fetch_category(client, tag_filters, category)
all_elements.extend(elements)
await asyncio.sleep(1.0)
logger.info(
"Overpass: total %d elements across %d category-queries",
len(all_elements),
len(OSM_CATEGORIES),
)
return all_elements
def sync_poi_to_db() -> dict[str, int]:
"""Синхронизирует POI из Overpass в osm_poi_ekb.
UPSERT по UNIQUE(osm_type, osm_id, category).
Returns: counters {fetched, inserted, updated, skipped_old}.
"""
elements = asyncio.run(fetch_overpass())
# 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,
}