gendesign/tradein-mvp/backend/app/services/dtp_stat_loader.py
lekss361 50f0674977
All checks were successful
Deploy Trade-In / changes (push) Successful in 11s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 4m9s
Deploy Trade-In / build-backend (push) Successful in 1m54s
Deploy Trade-In / deploy (push) Successful in 1m51s
Deploy Trade-In / deploy-status (push) Successful in 1s
Deploy Trade-In / perimeter-smoke (push) Successful in 11s
feat(tradein): слой ДТП из dtp-stat.ru в PostGIS + радиусные агрегаты (#3428)
2026-09-08 22:10:56 +00:00

281 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.

"""dtp-stat.ru loader: слой ДТП по Свердловской обл. в `dtp_incidents` (#3410).
CONTEXT: карточка района МЕРЫ хочет "фактор безопасности" — сколько ДТП/тяжёлых/
погибших/раненых было рядом за последние годы. dtp-stat.ru агрегирует открытые данные
ГИБДД в GeoJSON по регионам, без auth.
ИСТОЧНИК: https://dtp-stat.ru/media/opendata/sverdlovskaia-oblast.geojson.zip — ZIP
(~5 МБ) с одним файлом `sverdlovskaia-oblast.geojson` (~71 МБ распакованного,
FeatureCollection, ~31k Point-фич). Проверено живьём 08.09.2026: HTTP 200,
5 047 685 байт. Домен обычный (не RU-gov минцифровский сертификат) — verify=False
НЕ ставим, лишний.
МОНИТОРИНГ ПРОТУХАНИЯ (не реализован в этом PR — зафиксировано, чтобы следующий не
считал пустую дельту багом): на момент проверки 08.09.2026 дамп источника заморожен,
Last-Modified 26.02.2026. TRUNCATE+INSERT переливает те же ~31k строк каждый рефреш —
это ожидаемо, а не признак сломанного парсинга. Если нужен freshness-монитор — см.
паттерн `recon_macro.md` (СберИндекс), здесь не сделан осознанно (вне скоупа #3410).
СТРИМИНГ (обязательно, файл большой): `json.load()` целиком держал бы ~71 МБ + объектный
граф в памяти разом. Используем `ijson.items(stream, "features.item")` поверх
файлового объекта из `zipfile.ZipFile.open(name)` — постоянная память вне зависимости
от размера файла, каждая feature обрабатывается и отбрасывается по одной.
ПДн — НЕ СОХРАНЯЕМ (осознанное решение, зафиксировано и в 294_dtp_incidents.sql):
`properties.vehicles[].participants[]` источника несёт пол/роль/нарушения физлиц —
это персональные данные конкретных людей, продукту не нужны и не разрешены. Парсер
(`_parse_feature`) НИКОГДА не читает ключи `vehicles`/`participants` — ни в raw-jsonb,
ни отдельной колонкой. Из `properties` берём только неперсональные агрегаты и атрибуты
происшествия (see DtpIncidentRow).
Идемпотентность: TRUNCATE + bulk INSERT в одной транзакции (тот же паттерн, что
`app/tasks/osm_poi_ekb_refresh.py` — полная замена, не upsert; повторный прогон с тем
же дампом даёт тот же результат).
psycopg v3: SQL через `text(...)` использует `CAST(:x AS type)`, НИКОГДА `:x::type`.
Массивы (`weather`, `nearby`, тип `text[]`) — psycopg v3 адаптирует Python `list[str]`
в `text[]` напрямую через `CAST(:x AS text[])` (тот же идиом, что
`app/tasks/deactivate_stale_avito.py`).
"""
from __future__ import annotations
import io
import logging
import zipfile
from collections.abc import Iterator
from dataclasses import dataclass
from datetime import datetime
from pathlib import Path
from typing import IO, Any
import httpx
import ijson
from sqlalchemy import text
from sqlalchemy.orm import Session
logger = logging.getLogger(__name__)
SRC_URL = "https://dtp-stat.ru/media/opendata/sverdlovskaia-oblast.geojson.zip"
DOWNLOAD_TIMEOUT_SEC = 180
INSERT_CHUNK_SIZE = 2000
# properties.datetime — "2015-01-01 03:00:00", без явной таймзоны в источнике.
# Храним как naive → колонка timestamptz получит его в timezone сессии БД (обычно UTC).
# Источник не публикует TZ, поэтому точная привязка к МСК не гарантирована — для
# радиусных агрегатов по годам (dtp_index.py) это несущественно.
_DT_FMT = "%Y-%m-%d %H:%M:%S"
@dataclass(slots=True)
class DtpIncidentRow:
"""Одна строка ДТП — ТОЛЬКО неперсональные атрибуты (см. докстринг модуля)."""
source_id: str
dtp_at: datetime | None
category: str | None
severity: str | None
dead: int | None
injured: int | None
address: str | None
light: str | None
weather: list[str] | None
nearby: list[str] | None
lat: float
lon: float
def _parse_datetime(raw: Any) -> datetime | None:
if not isinstance(raw, str) or not raw.strip():
return None
try:
return datetime.strptime(raw.strip(), _DT_FMT)
except ValueError:
logger.warning("dtp_stat_loader: не распарсен datetime %r", raw)
return None
def _parse_str_list(raw: Any) -> list[str] | None:
if not isinstance(raw, list) or not raw:
return None
return [str(item) for item in raw if item is not None]
def _parse_int(raw: Any) -> int | None:
if isinstance(raw, bool):
return None
if isinstance(raw, int):
return raw
return None
# Санити-рамка координат. Выгрузка dtp-stat по Свердловской области содержит битые
# точки: на реальном дампе 2026-02-26 из 31 253 записей 74 лежат вне широт региона,
# в том числе 10 ровно в (0.0, 0.0) и такие значения, как (1.0, 1.0) и (41.0, 47.0).
# Пускать их в geo-таблицу нельзя: слой кормит радиусные агрегаты локации.
# Рамка намеренно шире официальных границ области (~56.0-62.3 N, ~57.2-66.2 E),
# чтобы не срезать легитимные приграничные точки.
REGION_LAT_RANGE = (55.0, 63.5)
REGION_LON_RANGE = (56.0, 67.5)
def _coords_plausible(lat: float, lon: float) -> bool:
"""Точка похожа на реальное ДТП в регионе выгрузки, а не на мусор источника."""
if lat == 0.0 and lon == 0.0:
return False
if not (REGION_LAT_RANGE[0] <= lat <= REGION_LAT_RANGE[1]):
return False
return REGION_LON_RANGE[0] <= lon <= REGION_LON_RANGE[1]
def _parse_feature(feature: dict[str, Any]) -> DtpIncidentRow | None:
"""Чистая функция: одна GeoJSON Feature -> DtpIncidentRow, либо None (пропуск).
НЕ читает `properties.vehicles` / `properties.participants` — намеренно (ПДн,
см. докстринг модуля). Пропускает фичи без валидной Point-геометрии/id, а также
с координатами вне санити-рамки региона (см. `_coords_plausible`).
"""
geometry = feature.get("geometry") or {}
if geometry.get("type") != "Point":
return None
coords = geometry.get("coordinates")
if not isinstance(coords, list) or len(coords) < 2:
return None
try:
lon, lat = float(coords[0]), float(coords[1])
except (TypeError, ValueError):
return None
if not _coords_plausible(lat, lon):
return None
props = feature.get("properties") or {}
source_id = props.get("id")
if source_id is None:
return None
return DtpIncidentRow(
source_id=str(source_id),
dtp_at=_parse_datetime(props.get("datetime")),
category=props.get("category") if isinstance(props.get("category"), str) else None,
severity=props.get("severity") if isinstance(props.get("severity"), str) else None,
dead=_parse_int(props.get("dead_count")),
injured=_parse_int(props.get("injured_count")),
address=props.get("address") if isinstance(props.get("address"), str) else None,
light=props.get("light") if isinstance(props.get("light"), str) else None,
weather=_parse_str_list(props.get("weather")),
nearby=_parse_str_list(props.get("nearby")),
lat=lat,
lon=lon,
)
def iter_incidents(stream: IO[bytes]) -> Iterator[DtpIncidentRow]:
"""Потоковый парс GeoJSON FeatureCollection -> DtpIncidentRow, по одной фиче за раз.
`stream` — бинарный файловый объект (из `zipfile.ZipFile.open()` или обычного
`open(path, "rb")`). Использует `ijson.items(..., "features.item")` — постоянная
память, НЕ читает файл целиком.
"""
for feature in ijson.items(stream, "features.item"):
row = _parse_feature(feature)
if row is not None:
yield row
def _open_geojson_in_zip(zf: zipfile.ZipFile) -> IO[bytes]:
names = [n for n in zf.namelist() if n.lower().endswith(".geojson")]
if not names:
raise ValueError(f"zip не содержит .geojson записей: {zf.namelist()!r}")
return zf.open(names[0])
def download_dtp_zip(*, client: httpx.Client) -> bytes:
"""Скачать ZIP dtp-stat.ru. Открытые данные, без auth/PII — стандартный verify."""
resp = client.get(SRC_URL, timeout=DOWNLOAD_TIMEOUT_SEC)
resp.raise_for_status()
return resp.content
_TRUNCATE_SQL = text("TRUNCATE dtp_incidents")
_INSERT_SQL = text(
"""
INSERT INTO dtp_incidents
(source_id, dtp_at, category, severity, dead, injured, address, light,
weather, nearby, lat, lon, geom)
VALUES
(CAST(:source_id AS text), CAST(:dtp_at AS timestamptz),
CAST(:category AS text), CAST(:severity AS text),
CAST(:dead AS integer), CAST(:injured AS integer),
CAST(:address AS text), CAST(:light AS text),
CAST(:weather AS text[]), CAST(:nearby AS text[]),
CAST(:lat AS double precision), CAST(:lon AS double precision),
ST_SetSRID(ST_MakePoint(CAST(:lon AS double precision), CAST(:lat AS double precision)),
4326))
ON CONFLICT (source_id) DO NOTHING
"""
)
def _row_params(row: DtpIncidentRow) -> dict[str, Any]:
return {
"source_id": row.source_id,
"dtp_at": row.dtp_at,
"category": row.category,
"severity": row.severity,
"dead": row.dead,
"injured": row.injured,
"address": row.address,
"light": row.light,
"weather": row.weather,
"nearby": row.nearby,
"lat": row.lat,
"lon": row.lon,
}
def load_dtp_incidents(
db: Session,
*,
src_path: str | Path | None = None,
client: httpx.Client | None = None,
chunk_size: int = INSERT_CHUNK_SIZE,
dry_run: bool = False,
) -> dict[str, int]:
"""TRUNCATE dtp_incidents; потоковый парс + bulk INSERT из ZIP dtp-stat.ru.
`src_path` — локальный ZIP (тесты/ручной прогон); если None — качаем `SRC_URL`
через `client` (обязателен параметром, чтобы был чем подменить в тестах;
создаётся временный `httpx.Client()` если не передан).
`dry_run=True` — парсит и считает строки, НЕ трогает БД (ни TRUNCATE, ни INSERT).
Не коммитит — коммитит caller (app/tasks/dtp_stat_refresh.py).
"""
owns_client = client is None
client = client or httpx.Client()
try:
if src_path is not None:
with zipfile.ZipFile(src_path) as zf, _open_geojson_in_zip(zf) as stream:
rows = list(iter_incidents(stream))
else:
data = download_dtp_zip(client=client)
with zipfile.ZipFile(io.BytesIO(data)) as zf, _open_geojson_in_zip(zf) as stream:
rows = list(iter_incidents(stream))
finally:
if owns_client:
client.close()
parsed = len(rows)
inserted = 0
if not dry_run:
db.execute(_TRUNCATE_SQL)
for i in range(0, len(rows), chunk_size):
chunk = rows[i : i + chunk_size]
result = db.execute(_INSERT_SQL, [_row_params(r) for r in chunk])
inserted += result.rowcount or 0
result_counts = {"parsed": parsed, "inserted": inserted}
logger.info("dtp_stat_loader load DONE (dry_run=%s): %s", dry_run, result_counts)
return result_counts