"""Batch geocoding для listings с NULL lat/lon. Запускается: - Manual через POST /admin/scrape/geocode-missing-listings - Scheduled: nightly via scrape_schedules (source='geocode_missing_listings', migration 110) — wired into in-app scheduler, window 06:00-09:00 UTC. Pattern: dedup по паре (address, city) — 1 уникальная пара → 1 geocode call → UPDATE всех listings с этим address+city (#2594 шаг 2/3: listings.city теперь заполняется скрапером из контекста развёртки — один и тот же текст адреса в разных городах («ул. Победы, 30» в ЕКБ и в Нижнем Тагиле) должен получать РАЗНЫЕ координаты, а не схлопываться в один geocode-вызов и один UPDATE по тексту адреса). Rate limit: Nominatim 1 req/sec (#2593: Yandex Geocoder tier удалён из geocoder). Отличие от /admin/geocode-missing (per-ID): - Этот модуль группирует по (address, city) → меньше API calls (dedup), но не схлопывает разные города с одинаковым текстом адреса. - Поддерживает all sources включая Avito (после PR #487 убрали jitter). - Возвращает GeocodeBackfillResult с детальными counters. - Loop-safe: SELECT фильтрует geocode_tried_at IS NULL OR tried_at < 7 days; при geocode failure помечает tried_at=NOW() → пара (address, city) не переотбирается в этом же run. """ from __future__ import annotations import logging import time from dataclasses import dataclass, field from sqlalchemy import text from sqlalchemy.orm import Session from app.services import scrape_runs as runs_mod from app.services.estimator import _geocode_is_coarse from app.services.geocoder import geocode logger = logging.getLogger(__name__) @dataclass class GeocodeBackfillResult: addresses_total: int = 0 # unique addresses pending geocode в этом batch addresses_processed: int = 0 # фактически обработано addresses_geocoded: int = 0 # успешно получили coords addresses_failed: int = 0 # geocoder вернул None listings_updated: int = 0 # total listings затронуто (1 address → N listings) cache_hits: int = 0 # из geocode_cache (instant) cache_misses: int = 0 # реальные geocoder calls duration_sec: float = field(default=0.0) async def geocode_missing_listings( db: Session, *, batch_size: int = 200, dry_run: bool = False, ) -> GeocodeBackfillResult: """Geocode listings с NULL coords (любой source). Steps: 1. SELECT address, city FROM listings WHERE lat IS NULL AND address IS NOT NULL GROUP BY address, city ORDER BY COUNT(*) DESC LIMIT batch_size (приоритет парам address+city с большим числом listings — больший ROI per geocode call; группировка по паре, НЕ только по address — #2594 шаг 2/3: один и тот же текст адреса в разных городах — разные записи) 2. Для каждой пары (address, city): - geocode(address, db, city_hint=city) — auto-cache (hit или miss) - Если есть результат: UPDATE listings SET lat, lon WHERE address = :addr AND city IS NOT DISTINCT FROM :city AND lat IS NULL (IS NOT DISTINCT FROM, а не `=` — стандартная SQL NULL-семантика: `city = NULL` никогда не true, поэтому обычным `=` группа с city IS NULL не обновилась бы вообще ни для одной строки; `IS NOT DISTINCT FROM` трактует NULL=NULL как совпадение, оставаясь строгим при непустом city — нужная нам симметрия) - PostGIS trigger (listings_set_geom_trg) автоматически обновит geom 3. Log progress каждые 50 addresses. Args: batch_size: max addresses to process per call (default 200 ≈ 3.5 min Nominatim) dry_run: только показать что бы сделалось, без UPDATE Returns: GeocodeBackfillResult с counters. """ start = time.monotonic() result = GeocodeBackfillResult() # 1. Найти top-N пар (address, city) с NULL coords (DESC by occurrence count). # Группировка по паре, а не только по address (#2594 шаг 2/3) — один и тот же # текст адреса в разных городах (напр. «ул. Победы, 30» в ЕКБ и в Нижнем Тагиле) # это разные записи с разными координатами, их нельзя схлопывать в один # geocode-вызов. GROUP BY address, city трактует NULL city как отдельную # группу (стандартная SQL-семантика группировки NULL как равных друг другу). # Фильтруем пары, по которым геокодер уже пробовал и не нашёл — они помечены # geocode_tried_at. Повторяем попытку только если tried_at старше 7 дней (возможен # переезд адреса в кэше или смена провайдера), либо tried_at IS NULL (ещё не пробовали). # Это делает функцию loop-safe: при вызове несколько раз в одном прогоне # failed-пары не переотбираются бесконечно. rows = ( db.execute( text( """ SELECT address, city, COUNT(*) AS listings_count FROM listings WHERE lat IS NULL AND address IS NOT NULL AND length(trim(address)) >= 5 AND (geocode_tried_at IS NULL OR geocode_tried_at < NOW() - INTERVAL '7 days') GROUP BY address, city ORDER BY listings_count DESC, address ASC, city ASC NULLS FIRST LIMIT :limit """ ), {"limit": batch_size}, ) .mappings() .all() ) result.addresses_total = len(rows) if not rows: logger.info("geocode_missing: 0 pending addresses — nothing to do") result.duration_sec = time.monotonic() - start return result logger.info( "geocode_missing: starting batch=%d total_pending_addresses=%d (top by listings count)", batch_size, result.addresses_total, ) for idx, row in enumerate(rows): address: str = row["address"] city: str | None = row.get("city") listings_count: int = row["listings_count"] result.addresses_processed += 1 try: geo = await geocode(address, db, city_hint=city) except Exception as exc: logger.warning("geocode_missing: geocode raised for '%s': %s", address[:60], exc) result.addresses_failed += 1 if not dry_run: # Пометить tried_at чтобы пара (address, city) не переотбиралась # в следующих batch'ах этого же прогона (loop-safe backoff 7 дней). # IS NOT DISTINCT FROM — city=NULL это отдельная группа, обычное # `=` не поймает NULL-город и не должно задеть другой город с тем # же текстом адреса. db.execute( text( "UPDATE listings SET geocode_tried_at = NOW()" " WHERE address = :addr AND city IS NOT DISTINCT FROM :city" " AND lat IS NULL" ), {"addr": address, "city": city}, ) db.commit() continue if geo is None: result.addresses_failed += 1 logger.info( "geocode_missing: NOT FOUND '%s' city=%r (used in %d listings)", address[:60], city, listings_count, ) if not dry_run: # Пометить tried_at — geocoder не нашёл адрес, backoff 7 дней. db.execute( text( "UPDATE listings SET geocode_tried_at = NOW()" " WHERE address = :addr AND city IS NOT DISTINCT FROM :city" " AND lat IS NULL" ), {"addr": address, "city": city}, ) db.commit() continue if geo.provider == "cache": result.cache_hits += 1 else: result.cache_misses += 1 result.addresses_geocoded += 1 if dry_run: logger.info( "geocode_missing[dry]: '%s' → (%.5f, %.5f) provider=%s would update %d listings", address[:60], geo.lat, geo.lon, geo.provider, listings_count, ) continue # Определяем точность геокода: city-centroid (нет номера дома) → 'city'. # _geocode_is_coarse() проверяет confidence='locality' ИЛИ отсутствие # house-number токена (1-3 цифры) в full_address — оба случая означают # что геокодер не дошёл до дома и вернул центр НП/города (#769 Part E). precision: str | None = "city" if _geocode_is_coarse(geo) else None # UPDATE listings — PostGIS trigger (listings_set_geom_trg) обновит geom автоматически. # geo_precision и geocode_tried_at проставляются одновременно с координатами. # city IS NOT DISTINCT FROM :city — обновляем ТОЛЬКО пару (address, city), из # которой был geocode-запрос; иначе тот же текст адреса в другом городе # (city IS NULL или другой явный город) перезаписался бы чужими координатами. update_result = db.execute( text( """ UPDATE listings SET lat = :lat, lon = :lon, geo_precision = :precision, geocode_tried_at = NOW() WHERE address = :addr AND city IS NOT DISTINCT FROM :city AND lat IS NULL """ ), { "lat": geo.lat, "lon": geo.lon, "precision": precision, "addr": address, "city": city, }, ) db.commit() result.listings_updated += update_result.rowcount if (idx + 1) % 50 == 0: elapsed = time.monotonic() - start rate = (idx + 1) / elapsed if elapsed > 0 else 0 logger.info( "geocode_missing: progress %d/%d " "(geocoded=%d failed=%d cache_hits=%d listings_updated=%d) rate=%.1f addr/s", idx + 1, len(rows), result.addresses_geocoded, result.addresses_failed, result.cache_hits, result.listings_updated, rate, ) result.duration_sec = time.monotonic() - start logger.info( "geocode_missing: DONE batch=%d processed=%d geocoded=%d failed=%d " "cache=(hit=%d miss=%d) listings_updated=%d duration=%.1fs", batch_size, result.addresses_processed, result.addresses_geocoded, result.addresses_failed, result.cache_hits, result.cache_misses, result.listings_updated, result.duration_sec, ) return result async def run_geocode_missing_listings( db: Session, *, run_id: int, params: dict, ) -> GeocodeBackfillResult: """Run-lifecycle wrapper: batch loop с wall-clock budget. Запускается планировщиком (source='geocode_missing_listings') или вручную. Гоняет geocode_missing_listings() в цикле до тех пор пока: - res.addresses_total == 0 (нет pending адресов) - res.addresses_total < batch_size (дренаж — последний batch меньше полного) - истёк budget_sec Params (из default_params jsonb в scrape_schedules): batch_size: int — адресов за один вызов geocode_missing_listings (default 200). budget_sec: float — максимальное время прогона в секундах (default 1800 = 30 мин). Lifecycle: update_heartbeat перед loop → аккумуляция counters → mark_done / mark_failed. """ batch_size = int(params.get("batch_size", 200)) budget_sec = float(params.get("budget_sec", 1800)) total = GeocodeBackfillResult() counters: dict[str, int] = { "checked": 0, "saved": 0, "skipped": 0, } try: runs_mod.update_heartbeat(db, run_id, counters) start = time.monotonic() while True: res = await geocode_missing_listings(db, batch_size=batch_size, dry_run=False) # Аккумулируем counters total.addresses_total += res.addresses_total total.addresses_processed += res.addresses_processed total.addresses_geocoded += res.addresses_geocoded total.addresses_failed += res.addresses_failed total.listings_updated += res.listings_updated total.cache_hits += res.cache_hits total.cache_misses += res.cache_misses counters = { "checked": total.addresses_processed, "saved": total.listings_updated, "skipped": total.addresses_failed, } runs_mod.update_heartbeat(db, run_id, counters) elapsed = time.monotonic() - start if res.addresses_total == 0: logger.info( "run_geocode_missing_listings: run_id=%d — нет pending адресов, завершаем", run_id, ) break if res.addresses_total < batch_size: logger.info( "run_geocode_missing_listings: run_id=%d — дренаж " "(addresses_total=%d < batch_size=%d), завершаем", run_id, res.addresses_total, batch_size, ) break if elapsed > budget_sec: logger.info( "run_geocode_missing_listings: run_id=%d — бюджет %.0fs исчерпан " "(elapsed=%.1fs), завершаем", run_id, budget_sec, elapsed, ) break total.duration_sec = time.monotonic() - start runs_mod.mark_done(db, run_id, counters) logger.info( "run_geocode_missing_listings: run_id=%d DONE — " "processed=%d geocoded=%d failed=%d listings_updated=%d duration=%.1fs", run_id, total.addresses_processed, total.addresses_geocoded, total.addresses_failed, total.listings_updated, total.duration_sec, ) return total except Exception as exc: total.duration_sec = time.monotonic() - start logger.exception( "run_geocode_missing_listings: run_id=%d FAILED after %.1fs", run_id, total.duration_sec, ) runs_mod.mark_failed(db, run_id, str(exc)[:1000], counters) raise