Merge pull request 'feat(tradein/db): авто-refresh ценовых бэндов по городам + порог для малых городов (#2576)' (#2579) from feat/tradein-price-bands-refresh into main
All checks were successful
Deploy Trade-In / changes (push) Successful in 12s
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 2m6s
Deploy Trade-In / build-backend (push) Successful in 1m16s
Deploy Trade-In / deploy (push) Successful in 5m19s

This commit is contained in:
lekss361 2026-07-31 14:01:34 +00:00
commit e1935b609f
5 changed files with 388 additions and 1 deletions

View file

@ -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"),

View file

@ -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

View file

@ -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;

View file

@ -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;

View file

@ -57,6 +57,7 @@ _PRODUCT_SOURCES: set[str] = {
"rosreestr_dkp_import",
"listing_source_snapshot",
"asking_to_sold_ratio_refresh",
"deal_city_price_bands_refresh",
"refresh_search_matview",
"yandex_address_backfill",
"deactivate_stale_avito",