diff --git a/tradein-mvp/backend/app/services/product_handlers.py b/tradein-mvp/backend/app/services/product_handlers.py index d6c70869..5503852b 100644 --- a/tradein-mvp/backend/app/services/product_handlers.py +++ b/tradein-mvp/backend/app/services/product_handlers.py @@ -113,6 +113,16 @@ async def _job_asking_to_sold_ratio( await loop.run_in_executor(None, recompute_asking_to_sold_ratios, db, run_id) +# ── deal_city_price_bands_refresh — sync tier-aware re-derive в executor ────── +async def _job_deal_city_price_bands_refresh( + db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext +) -> None: + from app.tasks.deal_city_price_bands_refresh import refresh_deal_city_price_bands + + loop = asyncio.get_event_loop() + await loop.run_in_executor(None, refresh_deal_city_price_bands, db, run_id) + + # ── refresh_search_matview — REFRESH MATVIEW CONCURRENTLY (own connection) ──── async def _job_refresh_search_matview( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext @@ -382,7 +392,7 @@ def build_product_handlers(ctx: SchedulerContext) -> dict[str, Handler]: """Реестр НЕ-sweep продуктовых source→Handler для kit build_registry. Kit-native sweeps (avito/yandex/cian/domclick city/full-load/newbuilding) НЕ здесь — - их даёт build_registry(_default_kit_handlers). Здесь — 18 именованных + 1 wildcard + их даёт build_registry(_default_kit_handlers). Здесь — 19 именованных + 1 wildcard (deactivate_stale_*), покрывающие каждый НЕ-sweep source боевого scheduler-dispatch. `ctx` — принят для симметрии контракта; сами Handler-job'ы получают ctx во время @@ -399,6 +409,9 @@ def build_product_handlers(ctx: SchedulerContext) -> dict[str, Handler]: "asking_to_sold_ratio_refresh": Handler( _job_asking_to_sold_ratio, "asking_to_sold_ratio_refresh" ), + "deal_city_price_bands_refresh": Handler( + _job_deal_city_price_bands_refresh, "deal_city_price_bands_refresh" + ), "refresh_search_matview": Handler(_job_refresh_search_matview, "refresh_search_matview"), "yandex_address_backfill": Handler(_job_yandex_address_backfill, "yandex_address_backfill"), "sber_index_pull": Handler(_job_sber_index_pull, "sber_index_pull"), diff --git a/tradein-mvp/backend/app/tasks/deal_city_price_bands_refresh.py b/tradein-mvp/backend/app/tasks/deal_city_price_bands_refresh.py new file mode 100644 index 00000000..5e2322f9 --- /dev/null +++ b/tradein-mvp/backend/app/tasks/deal_city_price_bands_refresh.py @@ -0,0 +1,161 @@ +"""Daily recompute of per-city ppm² plausible-deal guard-bands (#2576 Stage B). + +ПРОБЛЕМА: deal_city_price_bands (migration 178, tier-схема — migration 194) +засеяна ON CONFLICT DO UPDATE derivation-запросом. По мере ночного импорта новых +ДКП-сделок (rosreestr_dkp_import) города переходят между tier ('region_fallback' +N<10 → 'rough' N 10-29 → 'full' N>=30), а перцентили внутри tier дрейфуют — нужен +периодический пересчёт по той же derivation. + +Задача синхронная (DB-only, никаких внешних HTTP-вызовов) — запускается +kit-scheduler'ом через product_handlers._job_deal_city_price_bands_refresh +(run_in_executor), по образцу asking_to_sold_ratio.py / snapshot_listing_sources. + +Окно расписания 07:00-08:00 UTC — ПОСЛЕ rosreestr_dkp_import (04:00-06:00 UTC) И +asking_to_sold_ratio_refresh (06:00-07:00 UTC), чтобы бэнды считались по тому же +свежему срезу deals, что и ratio-таблица того же дня. + +SQL derivation ниже — БАЙТ-В-БАЙТ та же логика, что seed в +data/sql/194_deal_city_price_bands_tiers.sql (region_stats / city_stats / tiered: +трёхуровневая схема full N>=30 / rough N 10-29 / region_fallback N 1-9, см. +комментарий в 194 для полного обоснования тиров и hard floor/ceiling клампов). + +Нет DELETE перед re-derive (в отличие от asking_to_sold_ratio.py true-mirror +паттерна) — множество городов монотонно растёт (rosreestr_dkp_import только +INSERT/ON CONFLICT DO UPDATE, никогда не удаляет сделки), поэтому merge-по-city +(ON CONFLICT DO UPDATE) достаточен: город, перешедший в другой tier, просто +перезаписывается на следующем refresh. Екатеринбург НЕ включён (WHERE city <> +'Екатеринбург') — estimator.py fallback на глобальные DEAL_MIN_PPM2/DEAL_MAX_PPM2 +для ЕКБ остаётся byte-identical (invariant из 178/194 сохранён). +""" + +from __future__ import annotations + +import logging + +from sqlalchemy import text +from sqlalchemy.orm import Session + +from app.services import scrape_runs as runs_mod + +logger = logging.getLogger(__name__) + +# ── Derivation + re-seed (БАЙТ-В-БАЙТ из 194) ───────────────────────────────── +_REDERIVE_SQL = text( + """ + WITH region_stats AS ( + SELECT GREATEST( + round(percentile_cont(0.01) WITHIN GROUP (ORDER BY price_per_m2))::int, + 8000 + ) AS region_ppm2_min + FROM deals + WHERE source = 'rosreestr' + AND price_per_m2 IS NOT NULL + AND city IS NOT NULL + AND city <> 'Екатеринбург' + ), + city_stats AS ( + SELECT + city, + GREATEST(round(percentile_cont(0.01) WITHIN GROUP (ORDER BY price_per_m2))::int, 8000) + AS ppm2_p1, + LEAST(round(percentile_cont(0.99) WITHIN GROUP (ORDER BY price_per_m2))::int, 800000) + AS ppm2_p99, + count(*) AS n_deals + FROM deals + WHERE source = 'rosreestr' + AND price_per_m2 IS NOT NULL + AND city IS NOT NULL + AND city <> 'Екатеринбург' + GROUP BY city + ), + tiered AS ( + SELECT city, ppm2_p1 AS ppm2_min, ppm2_p99 AS ppm2_max, n_deals, + 'full'::text AS tier + FROM city_stats + WHERE n_deals >= 30 + AND ppm2_p99 >= 8000 + + UNION ALL + + SELECT city, LEAST(ppm2_p1, 700000) AS ppm2_min, 800000 AS ppm2_max, n_deals, + 'rough'::text AS tier + FROM city_stats + WHERE n_deals BETWEEN 10 AND 29 + + UNION ALL + + SELECT c.city, r.region_ppm2_min AS ppm2_min, 800000 AS ppm2_max, c.n_deals, + 'region_fallback'::text AS tier + FROM city_stats c + CROSS JOIN region_stats r + WHERE c.n_deals < 10 + ) + INSERT INTO deal_city_price_bands (city, ppm2_min, ppm2_max, n_deals, tier, refreshed_at) + SELECT city, ppm2_min, ppm2_max, n_deals, tier, now() + FROM tiered + ON CONFLICT (city) DO UPDATE + SET ppm2_min = EXCLUDED.ppm2_min, + ppm2_max = EXCLUDED.ppm2_max, + n_deals = EXCLUDED.n_deals, + tier = EXCLUDED.tier, + refreshed_at = EXCLUDED.refreshed_at + """ +) + +# ── Post-insert counters ────────────────────────────────────────────────────── +_COUNTERS_SQL = text( + """ + SELECT + COUNT(*) AS rows_written, + COUNT(*) FILTER (WHERE tier = 'full') AS full_rows, + COUNT(*) FILTER (WHERE tier = 'rough') AS rough_rows, + COUNT(*) FILTER (WHERE tier = 'region_fallback') AS region_fallback_rows + FROM deal_city_price_bands + """ +) + + +def refresh_deal_city_price_bands(db: Session, run_id: int) -> dict[str, int]: + """Пересчитать deal_city_price_bands (#2576 Stage B — tier-aware refresh). + + Sync (вызывается scheduler-триггером в executor, как recompute_asking_to_sold_ratios). + Одна транзакция: re-derive INSERT ... ON CONFLICT DO UPDATE (нет DELETE — см. + module docstring), затем counters из таблицы, commit, mark_done. + + Финализирует scrape_runs (mark_done / mark_failed) и пишет counters. + + Returns {"rows_written": N, "full_rows": .., "rough_rows": .., "region_fallback_rows": ..}. + """ + counters: dict[str, int] = { + "rows_written": 0, + "full_rows": 0, + "rough_rows": 0, + "region_fallback_rows": 0, + } + try: + db.execute(_REDERIVE_SQL) + + row = db.execute(_COUNTERS_SQL).mappings().first() + if row is not None: + counters["rows_written"] = int(row["rows_written"] or 0) + counters["full_rows"] = int(row["full_rows"] or 0) + counters["rough_rows"] = int(row["rough_rows"] or 0) + counters["region_fallback_rows"] = int(row["region_fallback_rows"] or 0) + + db.commit() + runs_mod.mark_done(db, run_id, counters) + logger.info( + "refresh_deal_city_price_bands run_id=%d done: " + "rows_written=%d full=%d rough=%d region_fallback=%d", + run_id, + counters["rows_written"], + counters["full_rows"], + counters["rough_rows"], + counters["region_fallback_rows"], + ) + return counters + except Exception as exc: + logger.exception("refresh_deal_city_price_bands run_id=%d failed", run_id) + db.rollback() + runs_mod.mark_failed(db, run_id, str(exc)[:1000], counters) + raise diff --git a/tradein-mvp/backend/data/sql/194_deal_city_price_bands_tiers.sql b/tradein-mvp/backend/data/sql/194_deal_city_price_bands_tiers.sql new file mode 100644 index 00000000..2a6e51ad --- /dev/null +++ b/tradein-mvp/backend/data/sql/194_deal_city_price_bands_tiers.sql @@ -0,0 +1,150 @@ +-- 194_deal_city_price_bands_tiers.sql +-- Эпик #2576 Stage B — многоуровневые ценовые бэнды по городам + честный +-- региональный фолбэк вместо ЕКБ-калиброванного порога. +-- +-- ПРОБЛЕМА: +-- Миграция 178 построила deal_city_price_bands РАЗОВО, только для городов +-- с count(*) >= 30 сделок (HAVING count(*) >= 30) на момент прогона. Auto-refresh +-- не был реализован (см. комментарий в 178). Город без строки в таблице +-- попадает на глобальный DEAL_MIN_PPM2=50_000 (estimator.py) — порог, +-- откалиброванный ИСКЛЮЧИТЕЛЬНО по Екатеринбургу. Для малых городов области +-- это не anti-outlier guard, а cut-off легитимного рынка (Североуральск +-- median ≈ 21.7k ₽/м²). +-- +-- Замер по прод-данным deals (2026-07-31, source='rosreestr', city IS NOT NULL, +-- city <> 'Екатеринбург', price_per_m2 IS NOT NULL — 47 253 сделки / 369 городов): +-- N>=30 сделок → 80 городов (45 988 сделок, 97.3%) — уже покрыты 178. +-- N 15-29 → 21 город ( 460 сделок) — падали на global-50k fallback. +-- N 10-14 → 21 город ( 247 сделок) — падали на global-50k fallback. +-- N 1-9 → 247 городов ( 558 сделок) — падали на global-50k fallback, +-- per-city перцентиль на такой выборке статистически бессмысленен +-- (n=1 → «перцентиль» = единственная сделка). +-- Итого 289 городов / 1265 сделок (2.7% выборки, но 78% ДОЛГОГО ХВОСТА городов) +-- получали ЕКБ-калиброванный пол вместо своей реальной цены. +-- +-- РЕШЕНИЕ — трёхуровневая схема (колонка tier), вместо единого порога 30: +-- 'full' N>=30 — own p1/p99 перцентиль (BYTE-IDENTICAL 178-derivation, +-- ЕКБ и существующие 80 городов НЕ меняются). +-- 'rough' 10<=N<30 — own p1 (floor), ceiling ФИКСИРОВАН на 800000 +-- (не деривится из тонкой выборки — p99 на <30 точках +-- нестабилен, одна дорогая сделка исказит потолок). +-- 'region_fallback' 1<=N<10 — own-данные города СЛИШКОМ тонкие даже для floor +-- (единичная сделка = 100% перцентиля недостоверна). +-- Используем ПУЛ по всей области (region_stats CTE, +-- p1 по 47k+ не-ЕКБ сделкам = 15 263 ₽/м² на момент +-- замера) вместо DEAL_MIN_PPM2=50000 (ЕКБ-калибровка). +-- Честнее: 15k отражает реальный низ рынка обл.66, +-- а не искусственно завышенный екб-порог. +-- +-- Екатеринбург по-прежнему НЕ включён (estimator.py fallback на глобальные +-- DEAL_MIN_PPM2/DEAL_MAX_PPM2 остаётся единственным путём для ЕКБ — invariant +-- из 178 сохранён). После этой миграции ЕВСЕ 369 не-ЕКБ городов, встречающихся +-- в deals, получают строку — Python-fallback в estimator.py (COALESCE(b.ppm2_min, +-- :ppm_min)) отныне срабатывает практически только для ЕКБ (плюс узкое окно +-- между refresh-циклами для только что появившегося города). +-- +-- IDEMPOTENCY: ADD COLUMN IF NOT EXISTS + DO-блок guard на CHECK constraint +-- (PG 16 не поддерживает ADD CONSTRAINT IF NOT EXISTS). INSERT ... ON CONFLICT +-- DO UPDATE — повторный прогон рефрешит бэнды под свежие сделки (та же +-- семантика, что и 178). Без DELETE — множество городов монотонно растёт +-- (rosreestr_dkp_import только INSERT/UPDATE, никогда не удаляет), поэтому +-- merge-по-ключу достаточен (см. app/tasks/deal_city_price_bands_refresh.py — +-- периодический refresh, та же derivation байт-в-байт). +-- +-- Dependencies: 177_deals_city_region.sql (deals.city), 178_deal_city_price_bands.sql +-- (таблица + PK(city)). +-- Apply after: -- +-- Deploy order: эта миграция ПЕРЕД деплоем backend-кода, который регистрирует +-- scheduler-source 'deal_city_price_bands_refresh' (product_handlers.py) — +-- см. 195_scrape_schedules_seed_deal_city_price_bands_refresh.sql (deploy after +-- backend-код задеплоен, тот же порядок, что 088). + +BEGIN; + +ALTER TABLE deal_city_price_bands + ADD COLUMN IF NOT EXISTS tier text NOT NULL DEFAULT 'full'; + +DO $$ +BEGIN + IF NOT EXISTS ( + SELECT 1 FROM pg_constraint WHERE conname = 'deal_city_price_bands_tier_check' + ) THEN + ALTER TABLE deal_city_price_bands + ADD CONSTRAINT deal_city_price_bands_tier_check + CHECK (tier IN ('full', 'rough', 'region_fallback')); + END IF; +END $$; + +COMMENT ON COLUMN deal_city_price_bands.tier IS + 'full: N>=30 сделок, own p1/p99 band (миграция 178, unchanged). ' + 'rough: 10<=N<30, own p1 floor + фиксированный 800000 ceiling (миграция 194). ' + 'region_fallback: 1<=N<10, pooled Свердловская-обл. p1 floor (region_stats, ' + 'все не-ЕКБ сделки) + фиксированный 800000 ceiling — вместо ' + 'ЕКБ-калиброванного DEAL_MIN_PPM2=50000 (estimator.py).'; + +WITH region_stats AS ( + -- Пул по ВСЕЙ области (не-ЕКБ) — честный фолбэк для городов, где own-выборка + -- (N<10) слишком тонкая для собственного перцентиля. + SELECT GREATEST( + round(percentile_cont(0.01) WITHIN GROUP (ORDER BY price_per_m2))::int, + 8000 + ) AS region_ppm2_min + FROM deals + WHERE source = 'rosreestr' + AND price_per_m2 IS NOT NULL + AND city IS NOT NULL + AND city <> 'Екатеринбург' +), +city_stats AS ( + SELECT + city, + GREATEST(round(percentile_cont(0.01) WITHIN GROUP (ORDER BY price_per_m2))::int, 8000) + AS ppm2_p1, + LEAST(round(percentile_cont(0.99) WITHIN GROUP (ORDER BY price_per_m2))::int, 800000) + AS ppm2_p99, + count(*) AS n_deals + FROM deals + WHERE source = 'rosreestr' + AND price_per_m2 IS NOT NULL + AND city IS NOT NULL + AND city <> 'Екатеринбург' + GROUP BY city +), +tiered AS ( + -- full — байт-в-байт исходная 178-derivation (own p1/p99), плюс тот же + -- анти-мусорный инвариант (p99 < 8000 → город не матчил бы ни одну сделку). + SELECT city, ppm2_p1 AS ppm2_min, ppm2_p99 AS ppm2_max, n_deals, + 'full'::text AS tier + FROM city_stats + WHERE n_deals >= 30 + AND ppm2_p99 >= 8000 + + UNION ALL + + -- rough — собственный p1 (floor), ceiling НЕ деривится (тонкая выборка). + SELECT city, LEAST(ppm2_p1, 700000) AS ppm2_min, 800000 AS ppm2_max, n_deals, + 'rough'::text AS tier + FROM city_stats + WHERE n_deals BETWEEN 10 AND 29 + + UNION ALL + + -- region_fallback — собственных данных недостаточно даже для floor, берём + -- пул по области целиком. + SELECT c.city, r.region_ppm2_min AS ppm2_min, 800000 AS ppm2_max, c.n_deals, + 'region_fallback'::text AS tier + FROM city_stats c + CROSS JOIN region_stats r + WHERE c.n_deals < 10 +) +INSERT INTO deal_city_price_bands (city, ppm2_min, ppm2_max, n_deals, tier, refreshed_at) +SELECT city, ppm2_min, ppm2_max, n_deals, tier, now() +FROM tiered +ON CONFLICT (city) DO UPDATE + SET ppm2_min = EXCLUDED.ppm2_min, + ppm2_max = EXCLUDED.ppm2_max, + n_deals = EXCLUDED.n_deals, + tier = EXCLUDED.tier, + refreshed_at = EXCLUDED.refreshed_at; + +COMMIT; diff --git a/tradein-mvp/backend/data/sql/195_scrape_schedules_seed_deal_city_price_bands_refresh.sql b/tradein-mvp/backend/data/sql/195_scrape_schedules_seed_deal_city_price_bands_refresh.sql new file mode 100644 index 00000000..a043fb34 --- /dev/null +++ b/tradein-mvp/backend/data/sql/195_scrape_schedules_seed_deal_city_price_bands_refresh.sql @@ -0,0 +1,62 @@ +-- 195_scrape_schedules_seed_deal_city_price_bands_refresh.sql +-- Эпик #2576 Stage B — seed scrape_schedules row для daily-рефреша +-- deal_city_price_bands (миграция 194). +-- +-- ПРОБЛЕМА: 178/194 заполняют deal_city_price_bands на момент прогона миграции. +-- По мере ночного импорта новых ДКП-сделок (rosreestr_dkp_import, 04:00-06:00 UTC) +-- бэнды (own p1/p99, tier-границы N) устаревают — города переходят между tier +-- ('region_fallback' → 'rough' → 'full') по мере накопления сделок, а сами +-- перцентили внутри tier дрейфуют. Auto-refresh отсутствовал (см. follow-up +-- в 178) — эта миграция закрывает разрыв. +-- +-- Задача (app/tasks/deal_city_price_bands_refresh.py, byte-identical derivation +-- 194) — pure-internal DB re-derivation, никаких внешних HTTP-вызовов. Запускается +-- kit-scheduler'ом через product_handlers._job_deal_city_price_bands_refresh +-- (run_in_executor, по образцу _job_asking_to_sold_ratio). +-- +-- enabled = true — БЕЗОПАСНО включать сразу (тот же аргумент, что 082/088: pure DB, +-- без анти-бота). +-- Окно 07:00-08:00 UTC — ПОСЛЕ rosreestr_dkp_import (04:00-06:00, см. 072) И +-- asking_to_sold_ratio_refresh (06:00-07:00, см. 082), чтобы бэнды считались по +-- тому же свежему срезу deals, что и ratio-таблица того же дня. +-- next_run_at = завтрашнее наступление окна (tomorrow + 07:00 UTC) — тот же паттерн, +-- что 078/079/082/088 (иначе get_due_schedules() выстрелит сразу после деплоя). +-- +-- ЗАВИСИМОСТИ: 052_scrape_schedules.sql (таблица + UNIQUE(source)), +-- 194_deal_city_price_bands_tiers.sql (tier-колонка, которую переиспользует refresh). +-- Idempotent: ON CONFLICT (source) DO NOTHING — безопасно запускать повторно. +-- Deploy order: применять ПОСЛЕ деплоя backend-кода, регистрирующего +-- 'deal_city_price_bands_refresh' в product_handlers.build_product_handlers() +-- (тот же порядок, что 088 relative к scheduler.py) — иначе kit-scheduler не +-- найдёт Handler для нового source и упадёт в "unknown source" на первом due-run +-- (не раньше завтрашнего окна — не блокирует деплой). + +BEGIN; + +INSERT INTO scrape_schedules ( + source, + enabled, + window_start_hour, + window_end_hour, + next_run_at, + default_params +) +VALUES +( + 'deal_city_price_bands_refresh', + true, -- SAFE: pure internal DB, no external calls + 7, + 8, + ((CURRENT_DATE + INTERVAL '1 day') + make_interval(hours => 7)) AT TIME ZONE 'UTC', + '{}'::jsonb +) +ON CONFLICT (source) DO NOTHING; + +COMMENT ON TABLE scrape_schedules IS + 'In-app scheduler config (заменяет cron-script setup). ' + 'Sources: avito_city_sweep, yandex_city_sweep (dormant, #561), ' + 'cian_history_backfill, rosreestr_dkp_import, listing_source_snapshot (#570), ' + 'asking_to_sold_ratio_refresh (#648), refresh_search_matview (#769), ' + 'deal_city_price_bands_refresh (#2576 Stage B).'; + +COMMIT;