"""Daily recompute of the asking→sold correction ratios (#648 Stage 4). ПРОБЛЕМА: asking_to_sold_ratios (migration 080) засеяна один раз derivation-CTE (медиана SOLD ДКП за 12 мес vs медиана ASKING активных listings, per-rooms + global -1 fallback). Оценщик (estimator.py, Stage 3) читает её per-estimate (cache 300s) и домножает asking-median прогноз на sold/asking. По мере ночного импорта новых ДКП-сделок (rosreestr_dkp_import) ratio устаревает — нужен периодический пересчёт по тому же derivation. TRUE-MIRROR REFRESH (флаг database-expert): 080-seed использовал ON CONFLICT DO UPDATE, который при повторном прогоне ОБНОВЛЯЕТ существующие строки, но НЕ удаляет per-rooms строки бакетов, упавших ниже порога 30/30 на более позднем прогоне (StaLE rows). Поэтому refresh сначала DELETE FROM asking_to_sold_ratios WHERE district = '' (все строки #648), затем заново гоняет ту же 080-derivation INSERT...SELECT. DELETE+INSERT в ОДНОЙ транзакции (атомарно — таблица никогда не пуста mid-refresh). После commit таблица == свежий re-seed. Задача синхронная (DB-only, никаких внешних HTTP-вызовов) — запускается in-app scheduler'ом через trigger_asking_to_sold_ratio_run() (scheduler.py), по образцу snapshot_listing_sources / import_rosreestr_dkp (sync task в run_in_executor). Окно расписания 06:00-07:00 UTC — ПОСЛЕ rosreestr_dkp_import (04:00-06:00 UTC), чтобы refresh потреблял свежие ДКП-сделки того же дня. SQL derivation ниже — БАЙТ-В-БАЙТ та же логика, что seed в data/sql/080_asking_to_sold_ratios.sql (deal_side / ask_side / per_bucket + deal_global / ask_global / global_row: трейлинг-12мес окно, ppm²-полоса [_PPM2_MIN, settings.asking_ratio_ppm2_max] (default [30000,1200000]), бакет LEAST(GREATEST(rooms,0),4), порог n_deals>=30 AND n_listings>=30 для per_rooms, global -1 строка всегда). ON CONFLICT убран — DELETE идёт первым, конфликтов нет (повторный прогон в одной tx невозможен, refresh = re-seed по семантике). """ from __future__ import annotations import logging from sqlalchemy import text from sqlalchemy.orm import Session from app.core.config import settings from app.services import scrape_runs as runs_mod # Нижняя граница ppm² — отсекает нежилые/технические сделки; не меняется. _PPM2_MIN: int = 30_000 # Верхняя граница берётся из settings.asking_ratio_ppm2_max (default 1_200_000). # QA-note: точное значение сверить с `SELECT max(price_per_m2) FROM deals # WHERE source='rosreestr'` на проде — ceiling должен быть > max(ppm²) premium-сделок. logger = logging.getLogger(__name__) # ── True-mirror cleanup: drop all #648 rows before re-derivation ────────────── # district = '' — все строки #648 (district зарезервирован под #647, пока всегда ''). # Удаляем ПЕРЕД re-derive, чтобы бакеты, упавшие ниже порога 30/30, не оставались # stale (ON CONFLICT DO UPDATE такие строки бы не тронул). В одной транзакции с INSERT. _DELETE_SQL = text( """ DELETE FROM asking_to_sold_ratios WHERE district = '' """ ) # ── Derivation + re-seed (БАЙТ-В-БАЙТ из 080, ON CONFLICT убран — DELETE идёт первым) ── # deal_side / ask_side / per_bucket + deal_global / ask_global / global_row: # sold_median = percentile_cont(0.5) по deals.price_per_m2 (source='rosreestr', # ppm² ∈ [_PPM2_MIN, settings.asking_ratio_ppm2_max], deal_date >= CURRENT_DATE − 12 months), # бакет LEAST(GREATEST(rooms,0),4). # ask_median = percentile_cont(0.5) по listings.price_per_m2 # (is_active, та же ppm²-полоса [_PPM2_MIN, asking_ratio_ppm2_max]). # per_rooms строки — только при n_deals>=30 AND n_listings>=30 AND ask>0 AND sold>0. # global -1 строка (basis='global_fallback') — всегда (если ask>0 AND sold>0). window_months=12. # Порог/окно — литералы; ppm²-полоса передаётся bind-параметрами :ppm2_min/:ppm2_max # (безопасно от SQL-инъекций; CAST не нужен — psycopg v3 передаёт int напрямую). _REDERIVE_SQL = text( """ WITH -- SOLD медианы по бакетам комнат за трейлинг-12мес (ДКП Росреестра). deal_side AS ( SELECT LEAST(GREATEST(rooms, 0), 4) AS rooms_bucket, percentile_cont(0.5) WITHIN GROUP (ORDER BY price_per_m2) AS sold_median, COUNT(*) AS n_deals FROM deals WHERE source = 'rosreestr' AND rooms IS NOT NULL AND price_per_m2 BETWEEN :ppm2_min AND :ppm2_max AND deal_date >= CURRENT_DATE - INTERVAL '12 months' GROUP BY LEAST(GREATEST(rooms, 0), 4) ), -- ASKING медианы по бакетам комнат среди ТЕКУЩИХ активных объявлений. ask_side AS ( SELECT LEAST(GREATEST(rooms, 0), 4) AS rooms_bucket, percentile_cont(0.5) WITHIN GROUP (ORDER BY price_per_m2) AS ask_median, COUNT(*) AS n_listings FROM listings WHERE is_active AND rooms IS NOT NULL AND price_per_m2 BETWEEN :ppm2_min AND :ppm2_max GROUP BY LEAST(GREATEST(rooms, 0), 4) ), -- Per-rooms строки: только бакеты с обеими сторонами, прошедшие порог 30/30 и ask>0. -- Тонкие бакеты (n<30) сюда НЕ попадают → estimator делает fallback на -1. per_bucket AS ( SELECT d.rooms_bucket, ''::text AS district, (d.sold_median / a.ask_median)::numeric AS ratio, round(d.sold_median)::bigint AS sold_median, round(a.ask_median)::bigint AS ask_median, d.n_deals::int AS n_deals, a.n_listings::int AS n_listings, 12 AS window_months, 'per_rooms'::text AS basis FROM deal_side d JOIN ask_side a USING (rooms_bucket) WHERE d.n_deals >= 30 AND a.n_listings >= 30 AND a.ask_median IS NOT NULL AND a.ask_median > 0 AND d.sold_median IS NOT NULL -- divide-safety (порог n_deals>=30 гарантирует) AND d.sold_median > 0 ), -- SOLD медиана по ВСЕМ комнатам (без бакет-фильтра) за трейлинг-12мес — для global row. deal_global AS ( SELECT percentile_cont(0.5) WITHIN GROUP (ORDER BY price_per_m2) AS sold_median, COUNT(*) AS n_deals FROM deals WHERE source = 'rosreestr' AND rooms IS NOT NULL AND price_per_m2 BETWEEN :ppm2_min AND :ppm2_max AND deal_date >= CURRENT_DATE - INTERVAL '12 months' ), -- ASKING медиана по ВСЕМ активным listings (без бакет-фильтра) — для global row. ask_global AS ( SELECT percentile_cont(0.5) WITHIN GROUP (ORDER BY price_per_m2) AS ask_median, COUNT(*) AS n_listings FROM listings WHERE is_active AND rooms IS NOT NULL AND price_per_m2 BETWEEN :ppm2_min AND :ppm2_max ), -- Global fallback строка rooms_bucket=-1 (пишется всегда, если ask>0). global_row AS ( SELECT -1 AS rooms_bucket, ''::text AS district, (d.sold_median / a.ask_median)::numeric AS ratio, round(d.sold_median)::bigint AS sold_median, round(a.ask_median)::bigint AS ask_median, d.n_deals::int AS n_deals, a.n_listings::int AS n_listings, 12 AS window_months, 'global_fallback'::text AS basis FROM deal_global d CROSS JOIN ask_global a WHERE a.ask_median IS NOT NULL AND a.ask_median > 0 AND d.sold_median IS NOT NULL -- защита от пустого окна сделок (иначе ratio=NULL) AND d.sold_median > 0 ) INSERT INTO asking_to_sold_ratios ( rooms_bucket, district, ratio, sold_median, ask_median, n_deals, n_listings, window_months, basis ) SELECT rooms_bucket, district, ratio, sold_median, ask_median, n_deals, n_listings, window_months, basis FROM global_row UNION ALL SELECT rooms_bucket, district, ratio, sold_median, ask_median, n_deals, n_listings, window_months, basis FROM per_bucket """ ) # ── Post-insert counters ────────────────────────────────────────────────────── # Считываем итог из таблицы (всё ещё в той же транзакции — до commit): сколько строк # записано всего, сколько per_rooms, был ли использован global -1 fallback. _COUNTERS_SQL = text( """ SELECT COUNT(*) AS rows_written, COUNT(*) FILTER (WHERE basis = 'per_rooms') AS per_rooms_rows, COUNT(*) FILTER (WHERE rooms_bucket = -1) AS used_global_fallback FROM asking_to_sold_ratios WHERE district = '' """ ) def recompute_asking_to_sold_ratios(db: Session, run_id: int) -> dict[str, int]: """Пересчитать asking_to_sold_ratios (TRUE-MIRROR refresh, #648 Stage 4). Sync (вызывается scheduler-триггером в executor, как snapshot_listing_sources). В ОДНОЙ транзакции (атомарно — таблица никогда не пуста mid-refresh): 1. DELETE FROM asking_to_sold_ratios WHERE district = '' — снести stale-строки. 2. Заново прогнать 080-derivation INSERT...SELECT (per_rooms при 30/30 + global -1). Затем counters из таблицы, commit, mark_done. Семантика == re-seed миграции 080. Финализирует scrape_runs (mark_done / mark_failed) и пишет counters. Returns {"rows_written": N, "per_rooms_rows": M, "used_global_fallback": 0|1}. """ counters: dict[str, int] = { "rows_written": 0, "per_rooms_rows": 0, "used_global_fallback": 0, } try: # DELETE + re-derive INSERT в одной транзакции (НЕ коммитим между ними — # таблица не должна остаться пустой, если INSERT упадёт). db.execute(_DELETE_SQL) db.execute( _REDERIVE_SQL, {"ppm2_min": _PPM2_MIN, "ppm2_max": settings.asking_ratio_ppm2_max}, ) row = db.execute(_COUNTERS_SQL).mappings().first() if row is not None: counters["rows_written"] = int(row["rows_written"] or 0) counters["per_rooms_rows"] = int(row["per_rooms_rows"] or 0) counters["used_global_fallback"] = int(row["used_global_fallback"] or 0) db.commit() runs_mod.mark_done(db, run_id, counters) logger.info( "recompute_asking_to_sold_ratios run_id=%d done: " "rows_written=%d per_rooms_rows=%d used_global_fallback=%d", run_id, counters["rows_written"], counters["per_rooms_rows"], counters["used_global_fallback"], ) return counters except Exception as exc: logger.exception("recompute_asking_to_sold_ratios run_id=%d failed", run_id) db.rollback() runs_mod.mark_failed(db, run_id, str(exc)[:1000], counters) raise