feat(tradein-sql): bootstrap houses + link 98.8% of listings (Phase A+B) #527
3 changed files with 337 additions and 1 deletions
|
|
@ -46,6 +46,20 @@ class AnalogLot(BaseModel):
|
||||||
distance_m: int | None = None # расстояние до целевой квартиры в метрах
|
distance_m: int | None = None # расстояние до целевой квартиры в метрах
|
||||||
|
|
||||||
|
|
||||||
|
class CianValuationSummary(BaseModel):
|
||||||
|
"""Cian Valuation Calculator данные для UI.
|
||||||
|
|
||||||
|
Источник: external_valuations table (source='cian_valuation').
|
||||||
|
Заполняется только при successful Cian Calculator call в estimator.py.
|
||||||
|
"""
|
||||||
|
|
||||||
|
sale_price_rub: int | None = None
|
||||||
|
rent_price_rub: int | None = None
|
||||||
|
chart: list[dict[str, Any]] = Field(default_factory=list)
|
||||||
|
chart_change_pct: float | None = None
|
||||||
|
chart_change_direction: Literal["increase", "decrease", "neutral"] | None = None
|
||||||
|
|
||||||
|
|
||||||
class AggregatedEstimate(BaseModel):
|
class AggregatedEstimate(BaseModel):
|
||||||
estimate_id: UUID
|
estimate_id: UUID
|
||||||
median_price_rub: int
|
median_price_rub: int
|
||||||
|
|
@ -66,6 +80,7 @@ class AggregatedEstimate(BaseModel):
|
||||||
sources_used: list[str] = Field(default_factory=list) # ['avito', 'cian', 'rosreestr']
|
sources_used: list[str] = Field(default_factory=list) # ['avito', 'cian', 'rosreestr']
|
||||||
data_freshness_minutes: int | None = None # сколько минут назад был самый свежий парсинг
|
data_freshness_minutes: int | None = None # сколько минут назад был самый свежий парсинг
|
||||||
est_days_on_market: int | None = None # прогноз срока продажи (медиана по аналогам)
|
est_days_on_market: int | None = None # прогноз срока продажи (медиана по аналогам)
|
||||||
|
cian_valuation: CianValuationSummary | None = None
|
||||||
# ── Параметры оценённой квартиры — нужны, чтобы восстановить карточку
|
# ── Параметры оценённой квартиры — нужны, чтобы восстановить карточку
|
||||||
# при открытии оценки по ссылке (?id=), когда формы-инпута уже нет ──
|
# при открытии оценки по ссылке (?id=), когда формы-инпута уже нет ──
|
||||||
area_m2: float | None = None
|
area_m2: float | None = None
|
||||||
|
|
|
||||||
|
|
@ -31,7 +31,12 @@ from uuid import uuid4
|
||||||
from sqlalchemy import text
|
from sqlalchemy import text
|
||||||
from sqlalchemy.orm import Session
|
from sqlalchemy.orm import Session
|
||||||
|
|
||||||
from app.schemas.trade_in import AggregatedEstimate, AnalogLot, TradeInEstimateInput
|
from app.schemas.trade_in import (
|
||||||
|
AggregatedEstimate,
|
||||||
|
AnalogLot,
|
||||||
|
CianValuationSummary,
|
||||||
|
TradeInEstimateInput,
|
||||||
|
)
|
||||||
from app.services.geocoder import GeocodeResult, geocode
|
from app.services.geocoder import GeocodeResult, geocode
|
||||||
from app.services.house_metadata import get_house_metadata
|
from app.services.house_metadata import get_house_metadata
|
||||||
from app.services.scrapers.avito_imv import (
|
from app.services.scrapers.avito_imv import (
|
||||||
|
|
@ -865,6 +870,17 @@ async def estimate_quality(
|
||||||
sources_used=sources_used,
|
sources_used=sources_used,
|
||||||
data_freshness_minutes=freshness_min,
|
data_freshness_minutes=freshness_min,
|
||||||
est_days_on_market=_estimate_days_on_market(listings_clean, deals),
|
est_days_on_market=_estimate_days_on_market(listings_clean, deals),
|
||||||
|
cian_valuation=(
|
||||||
|
CianValuationSummary(
|
||||||
|
sale_price_rub=int(cian_val.sale_price_rub) if cian_val.sale_price_rub else None,
|
||||||
|
rent_price_rub=int(cian_val.rent_price_rub) if cian_val.rent_price_rub else None,
|
||||||
|
chart=list(cian_val.chart or []),
|
||||||
|
chart_change_pct=cian_val.chart_change_pct,
|
||||||
|
chart_change_direction=cian_val.chart_change_direction,
|
||||||
|
)
|
||||||
|
if cian_val is not None
|
||||||
|
else None
|
||||||
|
),
|
||||||
area_m2=payload.area_m2,
|
area_m2=payload.area_m2,
|
||||||
rooms=payload.rooms,
|
rooms=payload.rooms,
|
||||||
floor=payload.floor,
|
floor=payload.floor,
|
||||||
|
|
@ -1598,4 +1614,5 @@ def _empty_estimate(
|
||||||
analogs=[],
|
analogs=[],
|
||||||
actual_deals=[],
|
actual_deals=[],
|
||||||
expires_at=expires_at,
|
expires_at=expires_at,
|
||||||
|
cian_valuation=None,
|
||||||
)
|
)
|
||||||
|
|
|
||||||
|
|
@ -0,0 +1,304 @@
|
||||||
|
-- 063_backfill_houses_and_link_listings.sql
|
||||||
|
-- Bootstrap houses table from existing listings data (Phase A+B).
|
||||||
|
--
|
||||||
|
-- Context:
|
||||||
|
-- listings: 18,256 rows (all have addresses)
|
||||||
|
-- houses: 0 rows (empty)
|
||||||
|
-- listings.house_id_fk: 0 linked
|
||||||
|
-- house_placement_history: 160 orphans, all source=yandex_valuation, raw_payload has no address
|
||||||
|
--
|
||||||
|
-- Steps:
|
||||||
|
-- 1. CREATE OR REPLACE FUNCTION tradein_normalize_short_addr
|
||||||
|
-- 2. INSERT Avito/CIAN-sourced houses from listings WHERE house_source+house_ext_id NOT NULL
|
||||||
|
-- 3. INSERT derived houses for listings WITHOUT ext_house_id (grouped by normalized address)
|
||||||
|
-- DEDUP: skips addresses already covered by step 2 (prevents split house rows)
|
||||||
|
-- 4. UPDATE listings.house_id_fk via source+ext_house_id match (direct path)
|
||||||
|
-- 5. UPDATE listings.house_id_fk via normalized address match (fuzzy path)
|
||||||
|
-- DEDUP: COALESCE prefers avito-sourced house over derived at same address
|
||||||
|
-- 6. Best-effort relink house_placement_history via raw_payload->>'address' (likely 0 rows)
|
||||||
|
-- 7. Dedup sanity check: WARN if any normalized address maps to >1 houses rows
|
||||||
|
-- 8. RAISE NOTICE with final counters
|
||||||
|
--
|
||||||
|
-- Dependencies:
|
||||||
|
-- - pgcrypto extension (for digest()) -- verified present
|
||||||
|
-- - listings table, houses table, house_placement_history table
|
||||||
|
-- - houses UNIQUE constraint on (source, ext_house_id)
|
||||||
|
-- - house_placement_history UNIQUE constraint on (source, ext_item_id)
|
||||||
|
--
|
||||||
|
-- Idempotency:
|
||||||
|
-- - INSERT ... ON CONFLICT DO NOTHING
|
||||||
|
-- - UPDATE ... WHERE house_id_fk IS NULL
|
||||||
|
-- - CREATE OR REPLACE FUNCTION
|
||||||
|
--
|
||||||
|
-- Deploy order: after 062_clean_avito_addresses.sql
|
||||||
|
|
||||||
|
BEGIN;
|
||||||
|
|
||||||
|
-- -------------------------------------------------------------------------
|
||||||
|
-- Step 1: Address normalizer function
|
||||||
|
-- Strips regional/city prefix and apartment/corpus suffix.
|
||||||
|
-- Handles addresses like:
|
||||||
|
-- "ул. Репина, 75/2 стр." → "ул. Репина, 75/2 стр."
|
||||||
|
-- "Свердловская обл., Екатеринбург, ул. Большакова, 17" → "ул. Большакова, 17"
|
||||||
|
-- "улица Яскина, 12 · р-н Октябрьский" → "улица Яскина, 12"
|
||||||
|
-- "Азина, 4.1.1" → "Азина, 4.1.1"
|
||||||
|
-- -------------------------------------------------------------------------
|
||||||
|
CREATE OR REPLACE FUNCTION tradein_normalize_short_addr(addr text)
|
||||||
|
RETURNS text
|
||||||
|
LANGUAGE sql
|
||||||
|
IMMUTABLE
|
||||||
|
PARALLEL SAFE
|
||||||
|
AS $$
|
||||||
|
SELECT trim(both ' ,.' FROM
|
||||||
|
regexp_replace(
|
||||||
|
regexp_replace(
|
||||||
|
regexp_replace(
|
||||||
|
regexp_replace(
|
||||||
|
regexp_replace(
|
||||||
|
-- Strip leading "Россия / РФ / Российская Федерация, "
|
||||||
|
regexp_replace(addr, '^\s*(?:Россия|РФ|Российская\s+Федерация)\s*,\s*', '', 'i'),
|
||||||
|
-- Strip leading region "Свердловская обл., " / "Свердловская область, "
|
||||||
|
'^[А-ЯЁа-яё][А-ЯЁа-яё\s-]+\s+обл(?:асть|\.)\s*,\s*', '', 'i'
|
||||||
|
),
|
||||||
|
-- Strip leading city "Екатеринбург, " / "г. Екатеринбург, " / "г Екатеринбург, "
|
||||||
|
'^г\.?\s*[А-ЯЁ][А-ЯЁа-яё-]+\s*,\s*', '', 'i'
|
||||||
|
),
|
||||||
|
-- Strip district suffix " · р-н ..." or " · district"
|
||||||
|
'\s*·\s*.+$', '', 'i'
|
||||||
|
),
|
||||||
|
-- Strip apartment/corpus/office suffix ", кв./корп./оф./пом./подъезд N..."
|
||||||
|
',\s*(?:кв\.?|корп\.?|к\.?|оф\.?|пом\.?|подъезд)\s*\d+.*$', '', 'i'
|
||||||
|
),
|
||||||
|
-- Strip trailing whitespace artefacts
|
||||||
|
'\s{2,}', ' ', 'g'
|
||||||
|
)
|
||||||
|
);
|
||||||
|
$$;
|
||||||
|
|
||||||
|
-- -------------------------------------------------------------------------
|
||||||
|
-- Step 2: Insert Avito/CIAN-sourced houses from listings
|
||||||
|
-- Uses DISTINCT ON to pick the most recently scraped data per (source, ext_house_id).
|
||||||
|
-- -------------------------------------------------------------------------
|
||||||
|
INSERT INTO houses (
|
||||||
|
source,
|
||||||
|
ext_house_id,
|
||||||
|
url,
|
||||||
|
address,
|
||||||
|
lat,
|
||||||
|
lon,
|
||||||
|
geom,
|
||||||
|
year_built,
|
||||||
|
house_type,
|
||||||
|
total_floors,
|
||||||
|
full_address,
|
||||||
|
first_seen_at,
|
||||||
|
last_scraped_at
|
||||||
|
)
|
||||||
|
SELECT DISTINCT ON (l.house_source, l.house_ext_id)
|
||||||
|
l.house_source AS source,
|
||||||
|
l.house_ext_id AS ext_house_id,
|
||||||
|
COALESCE(l.house_url, '') AS url,
|
||||||
|
l.address,
|
||||||
|
l.lat,
|
||||||
|
l.lon,
|
||||||
|
l.geom,
|
||||||
|
l.year_built,
|
||||||
|
l.house_type,
|
||||||
|
l.total_floors,
|
||||||
|
l.address AS full_address,
|
||||||
|
NOW() AS first_seen_at,
|
||||||
|
NOW() AS last_scraped_at
|
||||||
|
FROM listings l
|
||||||
|
WHERE l.house_source IS NOT NULL
|
||||||
|
AND l.house_ext_id IS NOT NULL
|
||||||
|
AND l.address IS NOT NULL
|
||||||
|
ORDER BY l.house_source, l.house_ext_id, l.scraped_at DESC
|
||||||
|
ON CONFLICT (source, ext_house_id) DO NOTHING;
|
||||||
|
|
||||||
|
-- -------------------------------------------------------------------------
|
||||||
|
-- Step 3: Insert derived houses for listings WITHOUT ext_house_id
|
||||||
|
-- Groups listings by normalized address; each group = one derived house.
|
||||||
|
-- Synthetic ext_house_id = first 32 hex chars of sha256(normalized_addr).
|
||||||
|
-- Uses DISTINCT ON (norm_addr) picking latest scraped_at per group.
|
||||||
|
--
|
||||||
|
-- DEDUP FIX: Excludes normalized addresses already covered by Step 2's
|
||||||
|
-- avito/cian-sourced houses. Prevents the same physical building getting
|
||||||
|
-- TWO rows (one avito-sourced, one derived) when it has listings from
|
||||||
|
-- both Avito (with ext_house_id) and other sources (without ext_house_id).
|
||||||
|
-- -------------------------------------------------------------------------
|
||||||
|
INSERT INTO houses (
|
||||||
|
source,
|
||||||
|
ext_house_id,
|
||||||
|
url,
|
||||||
|
address,
|
||||||
|
lat,
|
||||||
|
lon,
|
||||||
|
geom,
|
||||||
|
year_built,
|
||||||
|
house_type,
|
||||||
|
total_floors,
|
||||||
|
full_address,
|
||||||
|
short_address,
|
||||||
|
first_seen_at,
|
||||||
|
last_scraped_at
|
||||||
|
)
|
||||||
|
SELECT DISTINCT ON (norm_addr)
|
||||||
|
'derived' AS source,
|
||||||
|
substring(encode(digest(norm_addr, 'sha256'), 'hex'), 1, 32) AS ext_house_id,
|
||||||
|
'' AS url,
|
||||||
|
src.address,
|
||||||
|
src.lat,
|
||||||
|
src.lon,
|
||||||
|
src.geom,
|
||||||
|
src.year_built,
|
||||||
|
src.house_type,
|
||||||
|
src.total_floors,
|
||||||
|
src.address AS full_address,
|
||||||
|
norm_addr AS short_address,
|
||||||
|
NOW() AS first_seen_at,
|
||||||
|
NOW() AS last_scraped_at
|
||||||
|
FROM (
|
||||||
|
SELECT
|
||||||
|
l.*,
|
||||||
|
tradein_normalize_short_addr(l.address) AS norm_addr
|
||||||
|
FROM listings l
|
||||||
|
WHERE l.address IS NOT NULL
|
||||||
|
AND (l.house_source IS NULL OR l.house_ext_id IS NULL)
|
||||||
|
) src
|
||||||
|
WHERE src.norm_addr IS NOT NULL
|
||||||
|
AND length(src.norm_addr) >= 5
|
||||||
|
AND src.norm_addr NOT IN (
|
||||||
|
SELECT tradein_normalize_short_addr(address)
|
||||||
|
FROM houses
|
||||||
|
WHERE source != 'derived' AND address IS NOT NULL
|
||||||
|
)
|
||||||
|
ORDER BY norm_addr, src.scraped_at DESC
|
||||||
|
ON CONFLICT (source, ext_house_id) DO NOTHING;
|
||||||
|
|
||||||
|
-- -------------------------------------------------------------------------
|
||||||
|
-- Step 4: Link listings.house_id_fk via direct source+ext_house_id match
|
||||||
|
-- Uses listings_house_ext_id_idx (source, house_ext_id) — sargable.
|
||||||
|
-- -------------------------------------------------------------------------
|
||||||
|
UPDATE listings l
|
||||||
|
SET house_id_fk = h.id
|
||||||
|
FROM houses h
|
||||||
|
WHERE l.house_id_fk IS NULL
|
||||||
|
AND l.house_source IS NOT NULL
|
||||||
|
AND l.house_ext_id IS NOT NULL
|
||||||
|
AND l.house_source = h.source
|
||||||
|
AND l.house_ext_id = h.ext_house_id;
|
||||||
|
|
||||||
|
-- -------------------------------------------------------------------------
|
||||||
|
-- Step 5: Link listings.house_id_fk via normalized address (fuzzy path)
|
||||||
|
-- Covers listings where house_source/house_ext_id is absent.
|
||||||
|
-- houses.short_address was populated above (derived rows); joined by equality.
|
||||||
|
--
|
||||||
|
-- DEDUP FIX: For cian/yandex listings at an address that already has an
|
||||||
|
-- avito-sourced house (from Step 2), we prefer that avito-sourced house
|
||||||
|
-- over any derived row. This consolidates all listings for the same physical
|
||||||
|
-- building under a single houses row regardless of scrape source.
|
||||||
|
--
|
||||||
|
-- COALESCE priority:
|
||||||
|
-- 1. Non-derived house whose address normalizes to the same value
|
||||||
|
-- 2. Derived house whose short_address equals the normalized listing address
|
||||||
|
--
|
||||||
|
-- NOTE: houses_addr_fp_idx is on address_fingerprint; short_address has no
|
||||||
|
-- dedicated index. If this migration is re-run on large data, consider:
|
||||||
|
-- CREATE INDEX IF NOT EXISTS houses_short_addr_idx ON houses(short_address)
|
||||||
|
-- WHERE source = 'derived';
|
||||||
|
-- For a one-shot migration on 18k rows this is acceptable (seq scan ~ms).
|
||||||
|
-- -------------------------------------------------------------------------
|
||||||
|
WITH norm_l AS (
|
||||||
|
SELECT id, tradein_normalize_short_addr(address) AS norm_addr
|
||||||
|
FROM listings
|
||||||
|
WHERE house_id_fk IS NULL AND address IS NOT NULL
|
||||||
|
)
|
||||||
|
UPDATE listings l
|
||||||
|
SET house_id_fk = COALESCE(
|
||||||
|
-- Prefer avito/cian-sourced house at same normalized address (cross-source consolidation)
|
||||||
|
(SELECT h.id FROM houses h
|
||||||
|
WHERE h.source != 'derived'
|
||||||
|
AND tradein_normalize_short_addr(h.address) = nl.norm_addr
|
||||||
|
LIMIT 1),
|
||||||
|
-- Fall back to derived house
|
||||||
|
(SELECT h.id FROM houses h
|
||||||
|
WHERE h.source = 'derived'
|
||||||
|
AND h.short_address = nl.norm_addr
|
||||||
|
LIMIT 1)
|
||||||
|
)
|
||||||
|
FROM norm_l nl
|
||||||
|
WHERE l.id = nl.id AND l.house_id_fk IS NULL;
|
||||||
|
|
||||||
|
-- -------------------------------------------------------------------------
|
||||||
|
-- Step 6: Best-effort relink orphan house_placement_history
|
||||||
|
-- All 160 orphan rows are source=yandex_valuation; raw_payload has no
|
||||||
|
-- 'address' key (verified: has_addr=false for all 160 rows).
|
||||||
|
-- This UPDATE will match 0 rows but is included for completeness/idempotency.
|
||||||
|
-- Phase C (Celery IMV scraper) will assign house_id on new inserts directly.
|
||||||
|
-- -------------------------------------------------------------------------
|
||||||
|
UPDATE house_placement_history hph
|
||||||
|
SET house_id = h.id
|
||||||
|
FROM houses h
|
||||||
|
WHERE hph.house_id IS NULL
|
||||||
|
AND (hph.raw_payload ->> 'address') IS NOT NULL
|
||||||
|
AND h.source = 'derived'
|
||||||
|
AND h.short_address = tradein_normalize_short_addr(hph.raw_payload ->> 'address');
|
||||||
|
|
||||||
|
-- -------------------------------------------------------------------------
|
||||||
|
-- Step 7: Dedup sanity check
|
||||||
|
-- Warn if any normalized address maps to more than one houses row.
|
||||||
|
-- A non-zero count after the dedup fix above means manual review is needed.
|
||||||
|
-- -------------------------------------------------------------------------
|
||||||
|
DO $$
|
||||||
|
DECLARE
|
||||||
|
dup_count int;
|
||||||
|
BEGIN
|
||||||
|
SELECT count(*) INTO dup_count FROM (
|
||||||
|
SELECT tradein_normalize_short_addr(address) AS norm, count(*) AS n
|
||||||
|
FROM houses
|
||||||
|
WHERE address IS NOT NULL
|
||||||
|
GROUP BY tradein_normalize_short_addr(address)
|
||||||
|
HAVING count(*) > 1
|
||||||
|
) dup;
|
||||||
|
IF dup_count > 0 THEN
|
||||||
|
RAISE WARNING 'backfill 063: % normalized addresses still have multiple houses rows', dup_count;
|
||||||
|
ELSE
|
||||||
|
RAISE NOTICE 'backfill 063: dedup OK — 1 house per normalized address';
|
||||||
|
END IF;
|
||||||
|
END $$;
|
||||||
|
|
||||||
|
-- -------------------------------------------------------------------------
|
||||||
|
-- Step 8: Final counters via RAISE NOTICE
|
||||||
|
-- -------------------------------------------------------------------------
|
||||||
|
DO $$
|
||||||
|
DECLARE
|
||||||
|
v_houses_total bigint;
|
||||||
|
v_houses_derived bigint;
|
||||||
|
v_houses_direct bigint;
|
||||||
|
v_listings_total bigint;
|
||||||
|
v_listings_linked bigint;
|
||||||
|
v_listings_addr bigint;
|
||||||
|
v_history_total bigint;
|
||||||
|
v_history_linked bigint;
|
||||||
|
BEGIN
|
||||||
|
SELECT count(*) INTO v_houses_total FROM houses;
|
||||||
|
SELECT count(*) FILTER (WHERE source = 'derived') INTO v_houses_derived FROM houses;
|
||||||
|
SELECT count(*) FILTER (WHERE source != 'derived') INTO v_houses_direct FROM houses;
|
||||||
|
|
||||||
|
SELECT count(*) INTO v_listings_total FROM listings;
|
||||||
|
SELECT count(*) FILTER (WHERE house_id_fk IS NOT NULL) INTO v_listings_linked FROM listings;
|
||||||
|
SELECT count(*) FILTER (WHERE address IS NOT NULL) INTO v_listings_addr FROM listings;
|
||||||
|
|
||||||
|
SELECT count(*) INTO v_history_total FROM house_placement_history;
|
||||||
|
SELECT count(*) FILTER (WHERE house_id IS NOT NULL) INTO v_history_linked FROM house_placement_history;
|
||||||
|
|
||||||
|
RAISE NOTICE
|
||||||
|
'backfill 063 final: houses=% (direct=%, derived=%)'
|
||||||
|
' | listings_linked=%/% (addr_coverage=%)'
|
||||||
|
' | history_linked=%/%',
|
||||||
|
v_houses_total, v_houses_direct, v_houses_derived,
|
||||||
|
v_listings_linked, v_listings_total, v_listings_addr,
|
||||||
|
v_history_linked, v_history_total;
|
||||||
|
END $$;
|
||||||
|
|
||||||
|
COMMIT;
|
||||||
Loading…
Add table
Reference in a new issue