feat(tradein/rosreestr): параметризовать импорт ДКП по региону, deals.doc_type (#3051)
Трек 2 подготовки Mera к Москве. import_rosreestr_dkp принимает region_code из params (default 66 — байт-в-байт прежнее поведение), валидирует его через app.services.regions.REGIONS. Регион с canonical_city (77 — Москва, Росреестр отдаёт округ/поселение вместо города) подставляет city/address через одну SQL-ветку на bind-параметре :canonical_city, а не Python if/else на код региона; city IS NOT NULL не фильтруется для такого региона (иначе теряется ~10% строк), исходные city/okato/quarter_cad_number/district уходят в raw_payload. Чекпоинт курсора (_resume_dkp_cursor) стал per-region: source для поиска предыдущего прогона строится через _dkp_source_for_region (66 сохраняет легаси-имя 'rosreestr_dkp_import', остальные — суффикс кода) — иначе прогон по 77 либо никогда не резюмился бы (source-литерал не матчил), либо, при более наивном фиксе, унёс бы курсор чужого региона. product_handlers регистрирует wildcard rosreestr_dkp_import_* (по образцу deactivate_stale_*/avito_city_sweep_*), deploy/import-rosreestr.sh получил REGION_CODE env (bash-путь не region-generic — city-override только в Python). Migration 288: deals.doc_type + backfill 'ДКП' для source=rosreestr, foreign table gendesign_rosreestr_deals расширена okato/quarter_cad_number/district (проверено live на прод-БД), выключенный seed rosreestr_dkp_import_77.
This commit is contained in:
parent
278f8055a4
commit
84ee8e5990
8 changed files with 527 additions and 35 deletions
|
|
@ -750,6 +750,13 @@ def build_product_handlers(ctx: SchedulerContext) -> dict[str, Handler]:
|
||||||
post_claim=reschedule_after_minutes(param="interval_minutes", default=360),
|
post_claim=reschedule_after_minutes(param="interval_minutes", default=360),
|
||||||
),
|
),
|
||||||
"rosreestr_dkp_import": Handler(_job_rosreestr_dkp, "rosreestr_dkp_import"),
|
"rosreestr_dkp_import": Handler(_job_rosreestr_dkp, "rosreestr_dkp_import"),
|
||||||
|
# Wildcard (#3051 п.3): rosreestr_dkp_import_77 (Москва, миграция 288) и любой
|
||||||
|
# будущий region_code-суффикс из той же семьи резолвятся сюда через
|
||||||
|
# resolve_handler по префиксу (тот же механизм, что deactivate_stale_* /
|
||||||
|
# avito_city_sweep_* — см. scraper_kit.orchestration.scheduler.resolve_handler).
|
||||||
|
# region_code берётся из default_params строки расписания (import_rosreestr_dkp
|
||||||
|
# сам валидирует его через app.services.regions.REGIONS), Handler-тело общее.
|
||||||
|
"rosreestr_dkp_import_*": Handler(_job_rosreestr_dkp, "rosreestr_dkp_import_*"),
|
||||||
"listing_source_snapshot": Handler(_job_listing_source_snapshot, "listing_source_snapshot"),
|
"listing_source_snapshot": Handler(_job_listing_source_snapshot, "listing_source_snapshot"),
|
||||||
"asking_to_sold_ratio_refresh": Handler(
|
"asking_to_sold_ratio_refresh": Handler(
|
||||||
_job_asking_to_sold_ratio, "asking_to_sold_ratio_refresh"
|
_job_asking_to_sold_ratio, "asking_to_sold_ratio_refresh"
|
||||||
|
|
|
||||||
|
|
@ -47,6 +47,15 @@ class Region:
|
||||||
Регион без тира должен деградировать ЯВНО (потребитель
|
Регион без тира должен деградировать ЯВНО (потребитель
|
||||||
спрашивает unsupported_tier_reason и логирует/маркирует),
|
спрашивает unsupported_tier_reason и логирует/маркирует),
|
||||||
а не молча считать дальше без источника.
|
а не молча считать дальше без источника.
|
||||||
|
canonical_city — #3051: имя города, которым ПЕРЕЗАПИСЫВАЕТСЯ `city`
|
||||||
|
строк, приходящих из источника без надёжного city-поля
|
||||||
|
(Росреестр по Москве отдаёт муниципальный округ/поселение
|
||||||
|
вместо города — «Раменки», «Сосенское» — а не «Москва»).
|
||||||
|
None — источник несёт свой city как есть, без override
|
||||||
|
(регион 66: byte-for-byte прежнее поведение). Not-None —
|
||||||
|
потребитель (import_rosreestr_dkp) подставляет это имя
|
||||||
|
вместо city источника и НЕ фильтрует по city IS NOT NULL
|
||||||
|
(иначе на 77 теряется ~10% строк с пустым city).
|
||||||
"""
|
"""
|
||||||
|
|
||||||
code: int
|
code: int
|
||||||
|
|
@ -58,6 +67,7 @@ class Region:
|
||||||
city_token: str
|
city_token: str
|
||||||
cities: frozenset[str]
|
cities: frozenset[str]
|
||||||
enrichment_tiers: frozenset[str]
|
enrichment_tiers: frozenset[str]
|
||||||
|
canonical_city: str | None = None
|
||||||
|
|
||||||
|
|
||||||
def is_within_bbox(lat: float, lon: float, bbox: BBox) -> bool:
|
def is_within_bbox(lat: float, lon: float, bbox: BBox) -> bool:
|
||||||
|
|
@ -144,6 +154,10 @@ REGIONS: dict[int, Region] = {
|
||||||
# sber_index покрывают регион 66. Пустое множество здесь — не заглушка,
|
# sber_index покрывают регион 66. Пустое множество здесь — не заглушка,
|
||||||
# а ФАКТ, который потребители обязаны озвучивать (см. класс-докстринг).
|
# а ФАКТ, который потребители обязаны озвучивать (см. класс-докстринг).
|
||||||
enrichment_tiers=frozenset(),
|
enrichment_tiers=frozenset(),
|
||||||
|
# #3051: Росреестр по Москве отдаёт в city муниципальный округ/поселение
|
||||||
|
# ("муниципальный округ Раменки", "поселение Сосенское"), не сам город —
|
||||||
|
# import_rosreestr_dkp подставляет каноничное имя вместо city источника.
|
||||||
|
canonical_city="Москва",
|
||||||
),
|
),
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -25,6 +25,7 @@ Zombie-reap, advisory-lock claim и tick-loop теперь целиком в
|
||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import json
|
||||||
import logging
|
import logging
|
||||||
from typing import Any
|
from typing import Any
|
||||||
|
|
||||||
|
|
@ -50,6 +51,7 @@ from sqlalchemy.orm import Session
|
||||||
|
|
||||||
from app.core.shutdown import shutdown_requested
|
from app.core.shutdown import shutdown_requested
|
||||||
from app.services import scrape_runs as runs_mod
|
from app.services import scrape_runs as runs_mod
|
||||||
|
from app.services.regions import REGIONS
|
||||||
|
|
||||||
__all__ = ["compute_next_run_at", "has_running_run"]
|
__all__ = ["compute_next_run_at", "has_running_run"]
|
||||||
|
|
||||||
|
|
@ -238,6 +240,24 @@ async def _execute_cian_backfill(
|
||||||
|
|
||||||
_DKP_SOURCE = "rosreestr_dkp_import"
|
_DKP_SOURCE = "rosreestr_dkp_import"
|
||||||
|
|
||||||
|
|
||||||
|
def _dkp_source_for_region(region_code: int) -> str:
|
||||||
|
"""Имя scrape_runs.source для чекпоинта данного региона (#3051 п.3).
|
||||||
|
|
||||||
|
66 — байт-в-байт прежнее имя ('rosreestr_dkp_import'), под которым годами
|
||||||
|
писались scrape_runs. Остальные регионы получают суффикс кода — тот же
|
||||||
|
формат, что и у строки scrape_schedules ('rosreestr_dkp_import_77',
|
||||||
|
seed — миграция 288), которую резолвит wildcard 'rosreestr_dkp_import_*'
|
||||||
|
в product_handlers.py. Изоляция чекпоинтов между регионами держится именно
|
||||||
|
на разных source: _resume_dkp_cursor ищет ПРЕДЫДУЩИЙ прогон с ТЕМ ЖЕ source,
|
||||||
|
поэтому курсор региона 77 никогда не подхватит last_id региона 66 (и
|
||||||
|
наоборот) — они просто разные строки в scrape_runs.source.
|
||||||
|
"""
|
||||||
|
if region_code == 66:
|
||||||
|
return _DKP_SOURCE
|
||||||
|
return f"{_DKP_SOURCE}_{region_code}"
|
||||||
|
|
||||||
|
|
||||||
# Потолок возраста чекпоинта: старше — last_id прошлого прогона не подхватываем, прогон
|
# Потолок возраста чекпоинта: старше — last_id прошлого прогона не подхватываем, прогон
|
||||||
# стартует с id=0 (issue #3168). У предиката `id > last_id` нет протухания в смысле
|
# стартует с id=0 (issue #3168). У предиката `id > last_id` нет протухания в смысле
|
||||||
# свипов (он остаётся корректным сколь угодно долго), но апстрим
|
# свипов (он остаётся корректным сколь угодно долго), но апстрим
|
||||||
|
|
@ -263,7 +283,9 @@ _DKP_RESUME_CANDIDATE_SQL = text("""
|
||||||
""")
|
""")
|
||||||
|
|
||||||
|
|
||||||
def _resume_dkp_cursor(db: Session, run_id: int) -> tuple[int, dict[str, Any]]:
|
def _resume_dkp_cursor(
|
||||||
|
db: Session, run_id: int, source: str = _DKP_SOURCE
|
||||||
|
) -> tuple[int, dict[str, Any]]:
|
||||||
"""Продолжить last_id прошлого прогона или начать с 0 — решение + explain (issue #3168).
|
"""Продолжить last_id прошлого прогона или начать с 0 — решение + explain (issue #3168).
|
||||||
|
|
||||||
last_id раньше жил только в памяти процесса (init на 0 при каждом запуске): heartbeat
|
last_id раньше жил только в памяти процесса (init на 0 при каждом запуске): heartbeat
|
||||||
|
|
@ -271,7 +293,14 @@ def _resume_dkp_cursor(db: Session, run_id: int) -> tuple[int, dict[str, Any]]:
|
||||||
обратно. Обрыв (деплой/OOM/рестарт хоста) откатывал прогресс на 0 и заставлял
|
обратно. Обрыв (деплой/OOM/рестарт хоста) откатывал прогресс на 0 и заставлял
|
||||||
пере-сканировать источник с начала.
|
пере-сканировать источник с начала.
|
||||||
|
|
||||||
Кандидат — ПОСЛЕДНИЙ прогон этого source (тот же принцип, что и
|
`source` (#3051 п.3) — per-region ключ чекпоинта (см. _dkp_source_for_region):
|
||||||
|
дефолт _DKP_SOURCE сохраняет прежнее поведение вызовов без явного аргумента
|
||||||
|
(регион 66). Кандидат ищется СТРОГО по этому source — прогон региона 77
|
||||||
|
(source='rosreestr_dkp_import_77') никогда не видит last_id региона 66
|
||||||
|
(source='rosreestr_dkp_import') и наоборот: разные регионы физически не
|
||||||
|
матчат друг друга в WHERE source = :source ниже.
|
||||||
|
|
||||||
|
Кандидат — ПОСЛЕДНИЙ прогон ЭТОГО source (тот же принцип, что и
|
||||||
scraper_kit.orchestration.scheduler._pick_resume, локальная копия ладдера — контракт
|
scraper_kit.orchestration.scheduler._pick_resume, локальная копия ладдера — контракт
|
||||||
другой: нет params/interval_days, курсор числовой, а не bucket-set):
|
другой: нет params/interval_days, курсор числовой, а не bucket-set):
|
||||||
- 'running' / 'zombie' — прогон, которого не завершили штатно.
|
- 'running' / 'zombie' — прогон, которого не завершили штатно.
|
||||||
|
|
@ -285,7 +314,7 @@ def _resume_dkp_cursor(db: Session, run_id: int) -> tuple[int, dict[str, Any]]:
|
||||||
кодом (kit_runs.update_heartbeat — merge, не замена), чтобы решение было видно в
|
кодом (kit_runs.update_heartbeat — merge, не замена), чтобы решение было видно в
|
||||||
scrape_runs, а не только в логе.
|
scrape_runs, а не только в логе.
|
||||||
"""
|
"""
|
||||||
row = db.execute(_DKP_RESUME_CANDIDATE_SQL, {"source": _DKP_SOURCE, "rid": run_id}).fetchone()
|
row = db.execute(_DKP_RESUME_CANDIDATE_SQL, {"source": source, "rid": run_id}).fetchone()
|
||||||
verdict: dict[str, Any] = {"resume_from": None}
|
verdict: dict[str, Any] = {"resume_from": None}
|
||||||
|
|
||||||
if row is None:
|
if row is None:
|
||||||
|
|
@ -322,22 +351,35 @@ def import_rosreestr_dkp(
|
||||||
) -> None:
|
) -> None:
|
||||||
"""Import ДКП-сделок из gendesign rosreestr_deals через postgres_fdw.
|
"""Import ДКП-сделок из gendesign rosreestr_deals через postgres_fdw.
|
||||||
|
|
||||||
Python-порт import-rosreestr.sh (Variant C из #563).
|
Python-порт import-rosreestr.sh (Variant C из #563). #3051 п.3: параметризовано
|
||||||
|
по региону (params["region_code"], реестр — app.services.regions.REGIONS) —
|
||||||
|
было хардкод region_code=66.
|
||||||
|
|
||||||
Источник: foreign table gendesign_rosreestr_deals (создана в migration 072).
|
Источник: foreign table gendesign_rosreestr_deals (создана в migration 072,
|
||||||
|
okato/quarter_cad_number/district добавлены миграцией 288).
|
||||||
SERVER gendesign_remote настроен в 060_postgres_fdw_extension.sql.
|
SERVER gendesign_remote настроен в 060_postgres_fdw_extension.sql.
|
||||||
USER MAPPING создаётся при startup в core/fdw.py (tradein_fdw_reader).
|
USER MAPPING создаётся при startup в core/fdw.py (tradein_fdw_reader).
|
||||||
|
|
||||||
Область покрытия: вся Свердловская область (region_code=66), не только Екатеринбург —
|
Область покрытия: region_code из params (default 66 — вся Свердловская область,
|
||||||
прежний ILIKE-фильтр по подстроке города (ограничивавший импорт одним Екатеринбургом)
|
не только Екатеринбург). Неизвестный код региона (нет в REGIONS) — ValueError,
|
||||||
снят (Mera trade-in расширяется на весь регион, unlocks +47183 сделок вне ЕКБ уже
|
прогон падает явно, а не молча импортирует мусор с чужим region_code.
|
||||||
сидящих в source foreign table). address и deals.city строятся из реального city
|
|
||||||
источника (не хардкод "Екатеринбург"), deals.region_code заполняется из строки
|
region_code=66 (регион БЕЗ canonical_city в реестре) — поведение байт-в-байт
|
||||||
источника (= 66 при текущем фильтре).
|
прежнее: city/address строятся из city источника, обязателен фильтр
|
||||||
|
city IS NOT NULL AND trim(city) != ''.
|
||||||
|
|
||||||
|
Регион С canonical_city (77 — Москва): Росреестр отдаёт в city муниципальный
|
||||||
|
округ/поселение ("муниципальный округ Раменки", "поселение Сосенское"), НЕ
|
||||||
|
город — city/address подставляют region.canonical_city, а не city источника;
|
||||||
|
фильтр city IS NOT NULL НЕ применяется (иначе теряется ~10% строк с пустым
|
||||||
|
city источника). Исходные city/okato/quarter_cad_number/district уходят в
|
||||||
|
raw_payload (jsonb) — единственная ветка SQL решает это через bind-параметр
|
||||||
|
:canonical_city (CASE WHEN ... IS NOT NULL), а не отдельный Python if/else на
|
||||||
|
конкретный код региона.
|
||||||
|
|
||||||
Фильтры (совпадают с import-rosreestr.sh + Fix_Rosreestr_Dkp_Filter_May24):
|
Фильтры (совпадают с import-rosreestr.sh + Fix_Rosreestr_Dkp_Filter_May24):
|
||||||
- region_code = 66 (вся Свердловская область, все города)
|
- region_code = :region_code (параметризовано, было хардкод 66)
|
||||||
- city IS NOT NULL AND trim(city) != '' (непустой город → корректный address)
|
- city IS NOT NULL AND trim(city) != '' — ТОЛЬКО если у региона нет canonical_city
|
||||||
- realestate_type_code = '002001003000' (квартира)
|
- realestate_type_code = '002001003000' (квартира)
|
||||||
- area BETWEEN 18 AND 200
|
- area BETWEEN 18 AND 200
|
||||||
- deal_price BETWEEN 1000000 AND 100000000
|
- deal_price BETWEEN 1000000 AND 100000000
|
||||||
|
|
@ -354,17 +396,33 @@ def import_rosreestr_dkp(
|
||||||
(WHERE id > last_id ORDER BY id). Heartbeat обновляется каждый батч (= checkpoint),
|
(WHERE id > last_id ORDER BY id). Heartbeat обновляется каждый батч (= checkpoint),
|
||||||
мержем (kit_runs.update_heartbeat), а не заменой. На старте _resume_dkp_cursor решает
|
мержем (kit_runs.update_heartbeat), а не заменой. На старте _resume_dkp_cursor решает
|
||||||
продолжить с last_id прошлого прогона или начать с 0 — чекпоинт переживает рестарт
|
продолжить с last_id прошлого прогона или начать с 0 — чекпоинт переживает рестарт
|
||||||
процесса (деплой/OOM/SIGTERM), пока не старше суток (issue #3168).
|
процесса (деплой/OOM/SIGTERM), пока не старше суток (issue #3168). Курсор — ПЕР
|
||||||
|
РЕГИОН (#3051 п.3): _resume_dkp_cursor вызывается с source=_dkp_source_for_region
|
||||||
|
(region_code), поэтому last_id региона 77 никогда не подхватывает last_id региона
|
||||||
|
66 — они разные scrape_runs.source ('rosreestr_dkp_import' vs
|
||||||
|
'rosreestr_dkp_import_77'), см. докстринг _dkp_source_for_region.
|
||||||
SAVEPOINT per row — один сбойный row не откатывает батч.
|
SAVEPOINT per row — один сбойный row не откатывает батч.
|
||||||
|
|
||||||
Координаты: NULL после импорта — геокодинг остаётся follow-up (geocode-deals).
|
Координаты: NULL после импорта — геокодинг остаётся follow-up (geocode-deals).
|
||||||
TODO (follow-up): запустить geocode backfill после import.
|
TODO (follow-up): запустить geocode backfill после import.
|
||||||
|
|
||||||
Cleanup: удаляет legacy строки address='Екатеринбург, реальная сделка' (pre-#549).
|
Cleanup: удаляет legacy строки address='Екатеринбург, реальная сделка' (pre-#549,
|
||||||
|
region-agnostic — синтетические строки существовали только для ЕКБ).
|
||||||
"""
|
"""
|
||||||
since: str = str(params.get("since", "2024-01-01"))
|
since: str = str(params.get("since", "2024-01-01"))
|
||||||
batch_size: int = int(params.get("batch_size", 2000))
|
batch_size: int = int(params.get("batch_size", 2000))
|
||||||
|
|
||||||
|
region_code: int = int(params.get("region_code", 66))
|
||||||
|
region = REGIONS.get(region_code)
|
||||||
|
if region is None:
|
||||||
|
raise ValueError(
|
||||||
|
f"rosreestr_dkp_import: region_code={region_code} не найден в "
|
||||||
|
f"app.services.regions.REGIONS (известны: {sorted(REGIONS)}) — "
|
||||||
|
"прогон остановлен, чтобы не импортировать сделки с неизвестным "
|
||||||
|
"региональным контекстом (city/address-правила для него не определены)"
|
||||||
|
)
|
||||||
|
dkp_source = _dkp_source_for_region(region_code)
|
||||||
|
|
||||||
counters: dict[str, int] = {
|
counters: dict[str, int] = {
|
||||||
"rows_fetched": 0,
|
"rows_fetched": 0,
|
||||||
"rows_inserted": 0,
|
"rows_inserted": 0,
|
||||||
|
|
@ -400,7 +458,7 @@ def import_rosreestr_dkp(
|
||||||
)
|
)
|
||||||
db.rollback()
|
db.rollback()
|
||||||
|
|
||||||
last_id, resume_verdict = _resume_dkp_cursor(db, run_id)
|
last_id, resume_verdict = _resume_dkp_cursor(db, run_id, source=dkp_source)
|
||||||
total_batches = 0
|
total_batches = 0
|
||||||
kit_runs.update_heartbeat(db, run_id, resume_verdict)
|
kit_runs.update_heartbeat(db, run_id, resume_verdict)
|
||||||
logger.info(
|
logger.info(
|
||||||
|
|
@ -450,9 +508,20 @@ def import_rosreestr_dkp(
|
||||||
id,
|
id,
|
||||||
id AS source_id_src,
|
id AS source_id_src,
|
||||||
'ros:dkp:' || CAST(id AS text) AS dedup_hash,
|
'ros:dkp:' || CAST(id AS text) AS dedup_hash,
|
||||||
trim(city) || ', ' || trim(street) AS address,
|
-- #3051: регион с canonical_city (Москва) подставляет его вместо
|
||||||
|
-- city источника (округ/поселение, не город) — CASE на bind-параметре,
|
||||||
|
-- не Python if/else на код региона.
|
||||||
|
CASE
|
||||||
|
WHEN CAST(:canonical_city AS text) IS NOT NULL
|
||||||
|
THEN CAST(:canonical_city AS text) || ', ' || trim(street)
|
||||||
|
ELSE trim(city) || ', ' || trim(street)
|
||||||
|
END AS address,
|
||||||
region_code,
|
region_code,
|
||||||
trim(city) AS city,
|
CASE
|
||||||
|
WHEN CAST(:canonical_city AS text) IS NOT NULL
|
||||||
|
THEN CAST(:canonical_city AS text)
|
||||||
|
ELSE trim(city)
|
||||||
|
END AS city,
|
||||||
CASE
|
CASE
|
||||||
WHEN area < 30 THEN 0
|
WHEN area < 30 THEN 0
|
||||||
WHEN area < 44 THEN 1
|
WHEN area < 44 THEN 1
|
||||||
|
|
@ -472,10 +541,27 @@ def import_rosreestr_dkp(
|
||||||
year_build AS year_built,
|
year_build AS year_built,
|
||||||
round(deal_price)::bigint AS price_rub,
|
round(deal_price)::bigint AS price_rub,
|
||||||
round(price_per_sqm)::int AS price_per_m2,
|
round(price_per_sqm)::int AS price_per_m2,
|
||||||
period_start_date AS deal_date
|
period_start_date AS deal_date,
|
||||||
|
doc_type,
|
||||||
|
-- Исходный city/okato/quarter_cad_number/district — ТОЛЬКО когда
|
||||||
|
-- city перезаписан canonical_city выше (иначе NULL, регион 66
|
||||||
|
-- byte-for-byte прежний: raw_payload не заполнялся и не заполняется).
|
||||||
|
CASE
|
||||||
|
WHEN CAST(:canonical_city AS text) IS NOT NULL THEN
|
||||||
|
jsonb_build_object(
|
||||||
|
'src_city', city,
|
||||||
|
'okato', okato,
|
||||||
|
'quarter_cad_number', quarter_cad_number,
|
||||||
|
'district', district
|
||||||
|
)
|
||||||
|
ELSE NULL
|
||||||
|
END AS raw_payload
|
||||||
FROM gendesign_rosreestr_deals
|
FROM gendesign_rosreestr_deals
|
||||||
WHERE region_code = 66
|
WHERE region_code = CAST(:region_code AS int)
|
||||||
AND city IS NOT NULL AND trim(city) <> ''
|
AND (
|
||||||
|
CAST(:canonical_city AS text) IS NOT NULL
|
||||||
|
OR (city IS NOT NULL AND trim(city) <> '')
|
||||||
|
)
|
||||||
AND realestate_type_code = '002001003000'
|
AND realestate_type_code = '002001003000'
|
||||||
AND area BETWEEN 18 AND 200
|
AND area BETWEEN 18 AND 200
|
||||||
AND deal_price BETWEEN 1000000 AND 100000000
|
AND deal_price BETWEEN 1000000 AND 100000000
|
||||||
|
|
@ -486,7 +572,13 @@ def import_rosreestr_dkp(
|
||||||
ORDER BY id
|
ORDER BY id
|
||||||
LIMIT CAST(:batch_size AS int)
|
LIMIT CAST(:batch_size AS int)
|
||||||
"""),
|
"""),
|
||||||
{"since": since, "last_id": last_id, "batch_size": batch_size},
|
{
|
||||||
|
"since": since,
|
||||||
|
"last_id": last_id,
|
||||||
|
"batch_size": batch_size,
|
||||||
|
"region_code": region_code,
|
||||||
|
"canonical_city": region.canonical_city,
|
||||||
|
},
|
||||||
)
|
)
|
||||||
.mappings()
|
.mappings()
|
||||||
.all()
|
.all()
|
||||||
|
|
@ -524,7 +616,7 @@ def import_rosreestr_dkp(
|
||||||
INSERT INTO deals (
|
INSERT INTO deals (
|
||||||
source, dedup_hash, source_id, address, region_code, city,
|
source, dedup_hash, source_id, address, region_code, city,
|
||||||
rooms, area_m2, floor, year_built, price_rub, price_per_m2,
|
rooms, area_m2, floor, year_built, price_rub, price_per_m2,
|
||||||
deal_date
|
deal_date, doc_type, raw_payload
|
||||||
)
|
)
|
||||||
VALUES (
|
VALUES (
|
||||||
'rosreestr',
|
'rosreestr',
|
||||||
|
|
@ -539,7 +631,9 @@ def import_rosreestr_dkp(
|
||||||
CAST(:year_built AS int),
|
CAST(:year_built AS int),
|
||||||
CAST(:price_rub AS bigint),
|
CAST(:price_rub AS bigint),
|
||||||
CAST(:price_per_m2 AS int),
|
CAST(:price_per_m2 AS int),
|
||||||
CAST(:deal_date AS date)
|
CAST(:deal_date AS date),
|
||||||
|
CAST(:doc_type AS text),
|
||||||
|
CAST(:raw_payload AS jsonb)
|
||||||
)
|
)
|
||||||
ON CONFLICT (dedup_hash) DO UPDATE SET
|
ON CONFLICT (dedup_hash) DO UPDATE SET
|
||||||
address = EXCLUDED.address,
|
address = EXCLUDED.address,
|
||||||
|
|
@ -551,7 +645,9 @@ def import_rosreestr_dkp(
|
||||||
year_built = EXCLUDED.year_built,
|
year_built = EXCLUDED.year_built,
|
||||||
price_rub = EXCLUDED.price_rub,
|
price_rub = EXCLUDED.price_rub,
|
||||||
price_per_m2 = EXCLUDED.price_per_m2,
|
price_per_m2 = EXCLUDED.price_per_m2,
|
||||||
deal_date = EXCLUDED.deal_date
|
deal_date = EXCLUDED.deal_date,
|
||||||
|
doc_type = EXCLUDED.doc_type,
|
||||||
|
raw_payload = EXCLUDED.raw_payload
|
||||||
WHERE deals.address IS DISTINCT FROM EXCLUDED.address
|
WHERE deals.address IS DISTINCT FROM EXCLUDED.address
|
||||||
OR deals.region_code IS DISTINCT FROM EXCLUDED.region_code
|
OR deals.region_code IS DISTINCT FROM EXCLUDED.region_code
|
||||||
OR deals.city IS DISTINCT FROM EXCLUDED.city
|
OR deals.city IS DISTINCT FROM EXCLUDED.city
|
||||||
|
|
@ -562,6 +658,8 @@ def import_rosreestr_dkp(
|
||||||
OR deals.price_rub IS DISTINCT FROM EXCLUDED.price_rub
|
OR deals.price_rub IS DISTINCT FROM EXCLUDED.price_rub
|
||||||
OR deals.price_per_m2 IS DISTINCT FROM EXCLUDED.price_per_m2
|
OR deals.price_per_m2 IS DISTINCT FROM EXCLUDED.price_per_m2
|
||||||
OR deals.deal_date IS DISTINCT FROM EXCLUDED.deal_date
|
OR deals.deal_date IS DISTINCT FROM EXCLUDED.deal_date
|
||||||
|
OR deals.doc_type IS DISTINCT FROM EXCLUDED.doc_type
|
||||||
|
OR deals.raw_payload IS DISTINCT FROM EXCLUDED.raw_payload
|
||||||
RETURNING (xmax = 0) AS was_inserted
|
RETURNING (xmax = 0) AS was_inserted
|
||||||
"""),
|
"""),
|
||||||
{
|
{
|
||||||
|
|
@ -577,6 +675,12 @@ def import_rosreestr_dkp(
|
||||||
"price_rub": row["price_rub"],
|
"price_rub": row["price_rub"],
|
||||||
"price_per_m2": row["price_per_m2"],
|
"price_per_m2": row["price_per_m2"],
|
||||||
"deal_date": row["deal_date"],
|
"deal_date": row["deal_date"],
|
||||||
|
"doc_type": row["doc_type"],
|
||||||
|
"raw_payload": (
|
||||||
|
json.dumps(row["raw_payload"], ensure_ascii=False)
|
||||||
|
if row["raw_payload"] is not None
|
||||||
|
else None
|
||||||
|
),
|
||||||
},
|
},
|
||||||
).fetchone()
|
).fetchone()
|
||||||
if result is None:
|
if result is None:
|
||||||
|
|
|
||||||
87
tradein-mvp/backend/data/sql/288_deals_doc_type.sql
Normal file
87
tradein-mvp/backend/data/sql/288_deals_doc_type.sql
Normal file
|
|
@ -0,0 +1,87 @@
|
||||||
|
-- 288_deals_doc_type.sql
|
||||||
|
-- deals.doc_type (#3051 п.3) + foreign table gendesign_rosreestr_deals: okato/
|
||||||
|
-- quarter_cad_number/district + disabled seed schedule rosreestr_dkp_import_77.
|
||||||
|
--
|
||||||
|
-- Dependencies: 002_core_tables.sql (deals), 072_scrape_schedules_seed_cian_rosreestr.sql
|
||||||
|
-- (gendesign_rosreestr_deals, scrape_schedules).
|
||||||
|
-- Apply after: 287_proxy_run_attribution.sql
|
||||||
|
--
|
||||||
|
-- WHY:
|
||||||
|
-- Трек 2 подготовки Mera к Москве — импорт сделок Росреестра параметризуется по
|
||||||
|
-- региону (66 Свердловская обл. / 77 Москва, code-часть в scheduler.py). deals
|
||||||
|
-- до сих пор не различал ДКП (вторичка) от ДДУ (застройщик) в самой строке —
|
||||||
|
-- различие жило только в WHERE-фильтре импортёра. Явная колонка нужна для
|
||||||
|
-- будущего ДДУ-импорта (#3051 п.3) и для аналитики, которая иначе не может
|
||||||
|
-- отличить типы сделок в одной таблице.
|
||||||
|
--
|
||||||
|
-- Foreign table расширена тремя колонками источника (okato, quarter_cad_number,
|
||||||
|
-- district) — они нужны import_rosreestr_dkp для raw_payload по регионам, где
|
||||||
|
-- city источника не используется как есть (Москва: city = муниципальный
|
||||||
|
-- округ/поселение, не город). Существование колонок в public.rosreestr_deals
|
||||||
|
-- на gendesign-стороне проверено live (2026-09-08, prod psql).
|
||||||
|
--
|
||||||
|
-- Seed-строка rosreestr_dkp_import_77 — ВЫКЛЮЧЕНА (enabled=false): миграция
|
||||||
|
-- только заводит расписание, включение и первый прогон по Москве — отдельное
|
||||||
|
-- решение main-сессии после ревью кода-части.
|
||||||
|
--
|
||||||
|
-- ИДЕМПОТЕНТНОСТЬ: ADD COLUMN IF NOT EXISTS × 4, COMMENT ON COLUMN (безусловны,
|
||||||
|
-- но перезаписывают тот же текст), ON CONFLICT (source) DO NOTHING для seed.
|
||||||
|
-- Backfill (doc_type='ДКП' WHERE source='rosreestr' AND doc_type IS NULL) —
|
||||||
|
-- повторный прогон no-op (второй раз IS NULL уже не матчит ни одну строку).
|
||||||
|
|
||||||
|
BEGIN;
|
||||||
|
|
||||||
|
SET LOCAL lock_timeout = '5s';
|
||||||
|
|
||||||
|
-- ── deals.doc_type ────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
ALTER TABLE deals
|
||||||
|
ADD COLUMN IF NOT EXISTS doc_type text;
|
||||||
|
|
||||||
|
-- Backfill: весь текущий rosreestr-импорт в deals уже отфильтрован по doc_type='ДКП'
|
||||||
|
-- на стороне import_rosreestr_dkp (WHERE doc_type = 'ДКП') — строки, попавшие в deals
|
||||||
|
-- ДО этой миграции, были все ДКП, ДДУ-импорта ещё нет ни одной строки.
|
||||||
|
UPDATE deals
|
||||||
|
SET doc_type = 'ДКП'
|
||||||
|
WHERE source = 'rosreestr'
|
||||||
|
AND doc_type IS NULL;
|
||||||
|
|
||||||
|
COMMENT ON COLUMN deals.doc_type IS
|
||||||
|
'Тип документа сделки Росреестра: ДКП (договор купли-продажи, вторичка) или '
|
||||||
|
'ДДУ (договор долевого участия, застройщик) — #3051 п.3. NULL для строк не из '
|
||||||
|
'rosreestr (etazhi/domklik_history) — у них своя типизация сделки, либо не '
|
||||||
|
'применимо. Заполняется import_rosreestr_dkp из значения источника (WHERE '
|
||||||
|
'уже фильтрует doc_type=''ДКП'', колонка носит его явно, а не только в фильтре).';
|
||||||
|
|
||||||
|
-- ── gendesign_rosreestr_deals: колонки под региональный raw_payload (#3051) ────
|
||||||
|
-- ALTER FOREIGN TABLE ADD COLUMN — только локальные метаданные (не трогает
|
||||||
|
-- реальную remote-таблицу), безопасно как ALTER TABLE ADD COLUMN без DEFAULT.
|
||||||
|
|
||||||
|
ALTER FOREIGN TABLE gendesign_rosreestr_deals
|
||||||
|
ADD COLUMN IF NOT EXISTS okato text;
|
||||||
|
|
||||||
|
ALTER FOREIGN TABLE gendesign_rosreestr_deals
|
||||||
|
ADD COLUMN IF NOT EXISTS quarter_cad_number text;
|
||||||
|
|
||||||
|
ALTER FOREIGN TABLE gendesign_rosreestr_deals
|
||||||
|
ADD COLUMN IF NOT EXISTS district text;
|
||||||
|
|
||||||
|
-- ── Seed: rosreestr_dkp_import_77 (выключено) ──────────────────────────────────
|
||||||
|
|
||||||
|
INSERT INTO scrape_schedules (
|
||||||
|
source,
|
||||||
|
enabled,
|
||||||
|
window_start_hour,
|
||||||
|
window_end_hour,
|
||||||
|
default_params
|
||||||
|
)
|
||||||
|
VALUES (
|
||||||
|
'rosreestr_dkp_import_77',
|
||||||
|
false,
|
||||||
|
4,
|
||||||
|
6,
|
||||||
|
'{"region_code": 77, "since": "2024-01-01", "batch_size": 2000}'::jsonb
|
||||||
|
)
|
||||||
|
ON CONFLICT (source) DO NOTHING;
|
||||||
|
|
||||||
|
COMMIT;
|
||||||
247
tradein-mvp/backend/tests/test_3051_rosreestr_region_param.py
Normal file
247
tradein-mvp/backend/tests/test_3051_rosreestr_region_param.py
Normal file
|
|
@ -0,0 +1,247 @@
|
||||||
|
"""#3051 п.3: параметризация import_rosreestr_dkp по региону + deals.doc_type.
|
||||||
|
|
||||||
|
Трек 2 подготовки Mera к Москве (region_code=77). Чисто-юнит: SQL-текст
|
||||||
|
(inspect.getsource), _FakeDb-двойники для чекпоинта, миграция 288 и реестр
|
||||||
|
регионов — без живого FDW/Postgres (тот же стиль, что test_rosreestr_dedup_key.py
|
||||||
|
и test_3168_backfill_cursor_resume.py).
|
||||||
|
|
||||||
|
Покрывает пункты задачи:
|
||||||
|
(a) region_code больше не литерал 66 в SQL — bind-параметр (см. также
|
||||||
|
test_rosreestr_dedup_key.py::test_live_import_region_code_is_bind_param).
|
||||||
|
(b) маппинг 77: city='Москва', address с префиксом, raw_payload с src_city/okato/
|
||||||
|
quarter_cad_number/district.
|
||||||
|
(c) маппинг 66 не изменился (canonical_city=None → старые SQL-выражения нетронуты).
|
||||||
|
(d) чекпоинт per-region: _resume_dkp_cursor(source=...) изолирует регионы.
|
||||||
|
(e) реестр хендлеров резолвит rosreestr_dkp_import_77 через wildcard.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import inspect
|
||||||
|
import os
|
||||||
|
from pathlib import Path
|
||||||
|
from typing import Any
|
||||||
|
from unittest.mock import MagicMock
|
||||||
|
|
||||||
|
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
|
||||||
|
|
||||||
|
from app.services import scheduler as sched
|
||||||
|
from app.services.regions import REGIONS
|
||||||
|
|
||||||
|
_IMPORT_SRC = inspect.getsource(sched.import_rosreestr_dkp)
|
||||||
|
_SQL_DIR = Path(__file__).resolve().parents[1] / "data" / "sql"
|
||||||
|
_MIGRATION_288 = _SQL_DIR / "288_deals_doc_type.sql"
|
||||||
|
|
||||||
|
|
||||||
|
# ── регион 66/77 в реестре (canonical_city) ──────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
def test_region_66_has_no_canonical_city_override() -> None:
|
||||||
|
"""66 — источник несёт свой city как есть, поведение byte-for-byte прежнее."""
|
||||||
|
assert REGIONS[66].canonical_city is None
|
||||||
|
|
||||||
|
|
||||||
|
def test_region_77_canonical_city_is_moskva() -> None:
|
||||||
|
"""77 — Росреестр отдаёт округ/поселение вместо города, нужен override."""
|
||||||
|
assert REGIONS[77].canonical_city == "Москва"
|
||||||
|
|
||||||
|
|
||||||
|
# ── _dkp_source_for_region — имя чекпоинта per-region ────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
def test_dkp_source_for_region_66_keeps_legacy_name() -> None:
|
||||||
|
"""66 — byte-for-byte прежнее имя, под которым годами писались scrape_runs."""
|
||||||
|
assert sched._dkp_source_for_region(66) == "rosreestr_dkp_import"
|
||||||
|
|
||||||
|
|
||||||
|
def test_dkp_source_for_region_77_gets_suffix() -> None:
|
||||||
|
"""77 — суффикс кода, тот же формат, что у строки scrape_schedules (миграция 288)."""
|
||||||
|
assert sched._dkp_source_for_region(77) == "rosreestr_dkp_import_77"
|
||||||
|
|
||||||
|
|
||||||
|
# ── валидация неизвестного региона — падает ДО любого обращения к БД ────────
|
||||||
|
|
||||||
|
|
||||||
|
def test_import_rejects_unknown_region_code_before_touching_db() -> None:
|
||||||
|
"""Неизвестный region_code — явный ValueError, БД не трогается вовсе."""
|
||||||
|
db = MagicMock()
|
||||||
|
try:
|
||||||
|
sched.import_rosreestr_dkp(db, run_id=1, params={"region_code": 404})
|
||||||
|
raised = False
|
||||||
|
except ValueError as exc:
|
||||||
|
raised = True
|
||||||
|
assert "404" in str(exc)
|
||||||
|
assert raised, "expected ValueError for unknown region_code"
|
||||||
|
db.execute.assert_not_called()
|
||||||
|
|
||||||
|
|
||||||
|
# ── SQL: region_code — bind-параметр, не литерал (см. также test_rosreestr_dedup_key) ─
|
||||||
|
|
||||||
|
|
||||||
|
def test_sql_selects_region_code_via_bind_param() -> None:
|
||||||
|
assert "WHERE region_code = CAST(:region_code AS int)" in _IMPORT_SRC
|
||||||
|
|
||||||
|
|
||||||
|
def test_no_python_branch_on_region_77() -> None:
|
||||||
|
"""Маппинг city/address/raw_payload идёт ОДНОЙ SQL-веткой на bind-параметре
|
||||||
|
:canonical_city (CASE WHEN), а не Python if/else на конкретный код региона."""
|
||||||
|
assert "== 77" not in _IMPORT_SRC
|
||||||
|
assert "region_code == 77" not in _IMPORT_SRC
|
||||||
|
|
||||||
|
|
||||||
|
# ── SQL: маппинг city/address через canonical_city (регион 77) ──────────────
|
||||||
|
|
||||||
|
|
||||||
|
def test_sql_address_uses_canonical_city_case() -> None:
|
||||||
|
assert "CAST(:canonical_city AS text) || ', ' || trim(street)" in _IMPORT_SRC, (
|
||||||
|
"address для canonical_city-региона обязан быть 'Москва, <street>'"
|
||||||
|
)
|
||||||
|
assert "trim(city) || ', ' || trim(street)" in _IMPORT_SRC, (
|
||||||
|
"ELSE-ветка (регион 66) обязана остаться прежней"
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def test_sql_city_uses_canonical_city_case() -> None:
|
||||||
|
# WHEN CAST(:canonical_city AS text) IS NOT NULL THEN CAST(:canonical_city AS text)
|
||||||
|
assert "THEN CAST(:canonical_city AS text)\n" in _IMPORT_SRC
|
||||||
|
assert "ELSE trim(city)\n" in _IMPORT_SRC
|
||||||
|
|
||||||
|
|
||||||
|
def test_sql_city_not_null_filter_skipped_when_canonical_city_set() -> None:
|
||||||
|
"""Фильтр city IS NOT NULL применяется ТОЛЬКО когда canonical_city не задан —
|
||||||
|
иначе на 77 теряется ~10% строк с пустым city источника."""
|
||||||
|
assert (
|
||||||
|
"CAST(:canonical_city AS text) IS NOT NULL\n"
|
||||||
|
" OR (city IS NOT NULL AND trim(city) <> '')" in _IMPORT_SRC
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def test_sql_raw_payload_carries_source_columns() -> None:
|
||||||
|
"""raw_payload (только когда canonical_city задан) несёт src_city/okato/
|
||||||
|
quarter_cad_number/district — исходные значения источника, не потерянные."""
|
||||||
|
for key in (
|
||||||
|
"'src_city', city",
|
||||||
|
"'okato', okato",
|
||||||
|
"'quarter_cad_number', quarter_cad_number",
|
||||||
|
"'district', district",
|
||||||
|
):
|
||||||
|
assert key in _IMPORT_SRC, f"raw_payload missing key expression: {key!r}"
|
||||||
|
|
||||||
|
|
||||||
|
def test_params_dict_binds_region_code_and_canonical_city() -> None:
|
||||||
|
assert '"region_code": region_code' in _IMPORT_SRC
|
||||||
|
assert '"canonical_city": region.canonical_city' in _IMPORT_SRC
|
||||||
|
|
||||||
|
|
||||||
|
def test_insert_carries_doc_type_and_raw_payload() -> None:
|
||||||
|
assert "doc_type, raw_payload" in _IMPORT_SRC
|
||||||
|
assert "doc_type = EXCLUDED.doc_type" in _IMPORT_SRC
|
||||||
|
assert "raw_payload = EXCLUDED.raw_payload" in _IMPORT_SRC
|
||||||
|
|
||||||
|
|
||||||
|
# ── чекпоинт per-region: 77 не подхватывает last_id региона 66 ──────────────
|
||||||
|
|
||||||
|
|
||||||
|
class _FakeRow:
|
||||||
|
def __init__(self, **kw: Any) -> None:
|
||||||
|
self.__dict__.update(kw)
|
||||||
|
|
||||||
|
|
||||||
|
class _SourceKeyedFakeDb:
|
||||||
|
"""Двойник сессии: отдаёт кандидата ТОЛЬКО для своего source, иначе None.
|
||||||
|
|
||||||
|
Эмулирует реальный `WHERE source = :source` — единственный способ честно
|
||||||
|
проверить изоляцию чекпоинтов между регионами без живого Postgres.
|
||||||
|
"""
|
||||||
|
|
||||||
|
def __init__(self, rows_by_source: dict[str, Any]) -> None:
|
||||||
|
self.rows_by_source = rows_by_source
|
||||||
|
self.seen_sources: list[str] = []
|
||||||
|
|
||||||
|
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> Any:
|
||||||
|
if params is not None and "counters" in params:
|
||||||
|
return MagicMock()
|
||||||
|
source = (params or {}).get("source")
|
||||||
|
self.seen_sources.append(source)
|
||||||
|
return MagicMock(fetchone=lambda: self.rows_by_source.get(source))
|
||||||
|
|
||||||
|
def commit(self) -> None:
|
||||||
|
pass
|
||||||
|
|
||||||
|
|
||||||
|
def test_resume_cursor_region_77_does_not_see_region_66_checkpoint() -> None:
|
||||||
|
"""Прогон source='rosreestr_dkp_import' (66) не резюмится под source региона 77."""
|
||||||
|
row_66 = _FakeRow(
|
||||||
|
prev_id=777,
|
||||||
|
prev_status="zombie",
|
||||||
|
prev_counters={"last_id": 999999},
|
||||||
|
age_h=1.0,
|
||||||
|
)
|
||||||
|
db = _SourceKeyedFakeDb({"rosreestr_dkp_import": row_66})
|
||||||
|
|
||||||
|
last_id_77, verdict_77 = sched._resume_dkp_cursor(
|
||||||
|
db, run_id=1, source="rosreestr_dkp_import_77"
|
||||||
|
)
|
||||||
|
assert last_id_77 == 0
|
||||||
|
assert verdict_77["resume_reason"] == "no_prev_run"
|
||||||
|
|
||||||
|
# Регион 66 по-прежнему видит СВОЙ чекпоинт.
|
||||||
|
last_id_66, verdict_66 = sched._resume_dkp_cursor(db, run_id=2, source="rosreestr_dkp_import")
|
||||||
|
assert last_id_66 == 999999
|
||||||
|
assert verdict_66["resume_reason"] == "ok"
|
||||||
|
|
||||||
|
|
||||||
|
def test_resume_cursor_default_source_is_backward_compatible() -> None:
|
||||||
|
"""Вызов без явного source (как в старых тестах/коде) — прежнее поведение (66)."""
|
||||||
|
row = _FakeRow(prev_id=1, prev_status="zombie", prev_counters={"last_id": 42}, age_h=1.0)
|
||||||
|
db = _SourceKeyedFakeDb({"rosreestr_dkp_import": row})
|
||||||
|
last_id, _verdict = sched._resume_dkp_cursor(db, run_id=1)
|
||||||
|
assert last_id == 42
|
||||||
|
assert db.seen_sources == ["rosreestr_dkp_import"]
|
||||||
|
|
||||||
|
|
||||||
|
# ── реестр хендлеров: rosreestr_dkp_import_77 резолвится через wildcard ──────
|
||||||
|
|
||||||
|
|
||||||
|
def test_handler_registry_region_77_uses_same_job_as_region_66() -> None:
|
||||||
|
from scraper_kit.orchestration.scheduler import build_registry, resolve_handler
|
||||||
|
|
||||||
|
from app.services.product_handlers import _job_rosreestr_dkp, build_product_handlers
|
||||||
|
|
||||||
|
registry = build_registry(build_product_handlers(ctx=None)) # type: ignore[arg-type]
|
||||||
|
h66 = resolve_handler("rosreestr_dkp_import", registry)
|
||||||
|
h77 = resolve_handler("rosreestr_dkp_import_77", registry)
|
||||||
|
assert h66 is not None and h77 is not None
|
||||||
|
assert h66.job is _job_rosreestr_dkp
|
||||||
|
assert h77.job is _job_rosreestr_dkp
|
||||||
|
|
||||||
|
|
||||||
|
# ── миграция 288 ──────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
def test_migration_288_exists() -> None:
|
||||||
|
assert _MIGRATION_288.is_file(), f"missing migration: {_MIGRATION_288}"
|
||||||
|
|
||||||
|
|
||||||
|
def test_migration_288_adds_doc_type_idempotently() -> None:
|
||||||
|
sql = _MIGRATION_288.read_text("utf-8")
|
||||||
|
assert "ADD COLUMN IF NOT EXISTS doc_type text" in sql
|
||||||
|
assert "SET doc_type = 'ДКП'" in sql
|
||||||
|
assert "WHERE source = 'rosreestr'" in sql
|
||||||
|
assert "AND doc_type IS NULL" in sql
|
||||||
|
|
||||||
|
|
||||||
|
def test_migration_288_seeds_disabled_region_77_schedule() -> None:
|
||||||
|
sql = _MIGRATION_288.read_text("utf-8")
|
||||||
|
assert "'rosreestr_dkp_import_77'" in sql
|
||||||
|
assert "false" in sql
|
||||||
|
assert '"region_code": 77' in sql
|
||||||
|
assert "ON CONFLICT (source) DO NOTHING" in sql
|
||||||
|
|
||||||
|
|
||||||
|
def test_migration_288_has_lock_timeout_before_alter_table() -> None:
|
||||||
|
sql = _MIGRATION_288.read_text("utf-8")
|
||||||
|
lt_pos = sql.find("SET LOCAL lock_timeout")
|
||||||
|
alter_pos = sql.find("ALTER TABLE deals")
|
||||||
|
assert lt_pos != -1 and alter_pos != -1
|
||||||
|
assert lt_pos < alter_pos
|
||||||
|
|
@ -151,15 +151,21 @@ def test_migration_077_converts_md5_to_plain_key() -> None:
|
||||||
def test_migration_077_shared_filter_matches_live_import() -> None:
|
def test_migration_077_shared_filter_matches_live_import() -> None:
|
||||||
"""Дедуп-релевантные фильтры src CTE миграции 077 совпадают с живым импортом.
|
"""Дедуп-релевантные фильтры src CTE миграции 077 совпадают с живым импортом.
|
||||||
|
|
||||||
Исключение — city ILIKE (см. test_live_import_dropped_ekb_city_filter ниже):
|
Исключения:
|
||||||
077 backfill'ил ЕКБ-строки под EKB-only scope того времени; живой импорт расширен
|
- city ILIKE (см. test_live_import_dropped_ekb_city_filter ниже): 077 backfill'ил
|
||||||
на всю Свердловскую область (region_code=66, все города), поэтому city-фильтр из
|
ЕКБ-строки под EKB-only scope того времени; живой импорт расширен на всю
|
||||||
живого импорта СНЯТ намеренно. Остальные клозы обязаны совпадать байт-в-байт,
|
Свердловскую область (region_code=66, все города), поэтому city-фильтр из
|
||||||
иначе 077 конвертировал бы не тот набор строк.
|
живого импорта СНЯТ намеренно.
|
||||||
|
- region_code (#3051 п.3): 077 — историческая миграция (уже применена на прод,
|
||||||
|
хардкод region_code=66 трогать нельзя и не нужно — она навсегда про 66). Живой
|
||||||
|
импорт с #3051 параметризован по региону (region_code = CAST(:region_code AS
|
||||||
|
int)) — литерала 66 в его SQL больше нет, см.
|
||||||
|
test_live_import_region_code_is_bind_param ниже. Остальные клозы обязаны
|
||||||
|
совпадать байт-в-байт, иначе 077 конвертировал бы не тот набор строк.
|
||||||
"""
|
"""
|
||||||
sql = _MIGRATION_077.read_text("utf-8")
|
sql = _MIGRATION_077.read_text("utf-8")
|
||||||
|
assert "region_code = 66" in sql, "migration 077 must keep its historical literal"
|
||||||
for clause in (
|
for clause in (
|
||||||
"region_code = 66",
|
|
||||||
"realestate_type_code = '002001003000'",
|
"realestate_type_code = '002001003000'",
|
||||||
"area BETWEEN 18 AND 200",
|
"area BETWEEN 18 AND 200",
|
||||||
"deal_price BETWEEN 1000000 AND 100000000",
|
"deal_price BETWEEN 1000000 AND 100000000",
|
||||||
|
|
@ -170,6 +176,18 @@ def test_migration_077_shared_filter_matches_live_import() -> None:
|
||||||
assert clause in _IMPORT_SRC, f"missing filter clause in import: {clause!r}"
|
assert clause in _IMPORT_SRC, f"missing filter clause in import: {clause!r}"
|
||||||
|
|
||||||
|
|
||||||
|
def test_live_import_region_code_is_bind_param() -> None:
|
||||||
|
"""#3051 п.3: живой импорт параметризован по региону, литерала 66 в SQL нет.
|
||||||
|
|
||||||
|
До #3051 живой импорт хардкодил `WHERE region_code = 66` — единственный регион
|
||||||
|
покрытия. С параметризацией region_code приходит из params (default 66 — обратная
|
||||||
|
совместимость), в SQL идёт bind-параметром через CAST, не литералом.
|
||||||
|
"""
|
||||||
|
assert "region_code = CAST(:region_code AS int)" in _IMPORT_SRC
|
||||||
|
assert "region_code = 66" not in _IMPORT_BODY
|
||||||
|
assert 'params.get("region_code", 66)' in _IMPORT_SRC
|
||||||
|
|
||||||
|
|
||||||
def test_live_import_dropped_ekb_city_filter() -> None:
|
def test_live_import_dropped_ekb_city_filter() -> None:
|
||||||
"""Живой импорт БОЛЬШЕ не фильтрует по городу — Mera расширена на всю обл. 66.
|
"""Живой импорт БОЛЬШЕ не фильтрует по городу — Mera расширена на всю обл. 66.
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -57,6 +57,10 @@ from scraper_kit.orchestration.scheduler import (
|
||||||
_PRODUCT_SOURCES: set[str] = {
|
_PRODUCT_SOURCES: set[str] = {
|
||||||
"cian_history_backfill",
|
"cian_history_backfill",
|
||||||
"rosreestr_dkp_import",
|
"rosreestr_dkp_import",
|
||||||
|
# #3051 п.3: member-source семейства "rosreestr_dkp_import_*" (Москва, миграция 288
|
||||||
|
# — то же раскрытие wildcard в конкретный member, что deactivate_stale_avito/yandex/
|
||||||
|
# cian ниже, а не сам wildcard "rosreestr_dkp_import_*").
|
||||||
|
"rosreestr_dkp_import_77",
|
||||||
"listing_source_snapshot",
|
"listing_source_snapshot",
|
||||||
"asking_to_sold_ratio_refresh",
|
"asking_to_sold_ratio_refresh",
|
||||||
"deal_city_price_bands_refresh",
|
"deal_city_price_bands_refresh",
|
||||||
|
|
|
||||||
|
|
@ -2,12 +2,22 @@
|
||||||
# Импорт реальных сделок Росреестра из gendesign-БД в tradein deals.
|
# Импорт реальных сделок Росреестра из gendesign-БД в tradein deals.
|
||||||
#
|
#
|
||||||
# Источник: gendesign-postgres-1 / rosreestr_deals (6.8М строк, партиц.).
|
# Источник: gendesign-postgres-1 / rosreestr_deals (6.8М строк, партиц.).
|
||||||
# Берём ЕКБ квартиры (realestate_type_code=002001003000) за последние ~18 мес.
|
# Берём квартиры региона REGION_CODE (realestate_type_code=002001003000) за
|
||||||
# rooms выводим из площади (Росреестр не отдаёт кол-во комнат).
|
# последние ~18 мес. rooms выводим из площади (Росреестр не отдаёт кол-во комнат).
|
||||||
# Координаты — NULL, проставляются отдельно (geocode по street).
|
# Координаты — NULL, проставляются отдельно (geocode по street).
|
||||||
#
|
#
|
||||||
# Только ДКП (вторичка) — ДДУ застройщиков скёюят median, отделены сознательно. См. PR-A.
|
# Только ДКП (вторичка) — ДДУ застройщиков скёюят median, отделены сознательно. См. PR-A.
|
||||||
#
|
#
|
||||||
|
# #3051 п.3: REGION_CODE параметризован (default 66 — Свердловская обл., byte-for-byte
|
||||||
|
# прежнее поведение). Этот bash-путь — НЕ region-generic: для region_code=77 (Москва) он
|
||||||
|
# НЕ подставляет canonical_city вместо city источника (Росреестр по Москве отдаёт
|
||||||
|
# муниципальный округ/поселение, не сам город) и не пишет raw_payload с
|
||||||
|
# okato/quarter_cad_number/district — эту логику несёт только Python-путь
|
||||||
|
# (app/services/scheduler.py::import_rosreestr_dkp, боевой планировщик). Если этот
|
||||||
|
# скрипт когда-нибудь запустят вручную с REGION_CODE=77 — city/address будут
|
||||||
|
# "муниципальный округ Раменки, ..." как есть из источника, НЕ "Москва, ...". Держать
|
||||||
|
# паритет фильтров (area/price/doc_type) обязательно, паритет city-override — нет.
|
||||||
|
#
|
||||||
# Запуск на прод-хосте: ./import-rosreestr.sh
|
# Запуск на прод-хосте: ./import-rosreestr.sh
|
||||||
# Повторяемо: дедуп по dedup_hash, новый запуск подтянет свежие кварталы.
|
# Повторяемо: дедуп по dedup_hash, новый запуск подтянет свежие кварталы.
|
||||||
|
|
||||||
|
|
@ -18,8 +28,9 @@ DST_PG="${DST_PG:-tradein-postgres}"
|
||||||
SRC_DB="${SRC_DB:-gendesign}"
|
SRC_DB="${SRC_DB:-gendesign}"
|
||||||
SRC_USER="${SRC_USER:-gendesign}"
|
SRC_USER="${SRC_USER:-gendesign}"
|
||||||
SINCE="${SINCE:-2024-01-01}"
|
SINCE="${SINCE:-2024-01-01}"
|
||||||
|
REGION_CODE="${REGION_CODE:-66}"
|
||||||
|
|
||||||
echo "[$(date -u +%H:%M:%S)] import-rosreestr: ЕКБ квартиры с $SINCE"
|
echo "[$(date -u +%H:%M:%S)] import-rosreestr: регион $REGION_CODE, квартиры с $SINCE"
|
||||||
|
|
||||||
# 1. Staging-таблица в tradein.
|
# 1. Staging-таблица в tradein.
|
||||||
docker exec "$DST_PG" psql -U tradein -d tradein -v ON_ERROR_STOP=on -c "
|
docker exec "$DST_PG" psql -U tradein -d tradein -v ON_ERROR_STOP=on -c "
|
||||||
|
|
@ -53,7 +64,7 @@ docker exec "$SRC_PG" psql -U "$SRC_USER" -d "$SRC_DB" -v ON_ERROR_STOP=on -c "
|
||||||
round(price_per_sqm)::int AS price_per_m2,
|
round(price_per_sqm)::int AS price_per_m2,
|
||||||
period_start_date AS deal_date
|
period_start_date AS deal_date
|
||||||
FROM rosreestr_deals
|
FROM rosreestr_deals
|
||||||
WHERE region_code = 66
|
WHERE region_code = $REGION_CODE
|
||||||
AND city IS NOT NULL AND trim(city) <> ''
|
AND city IS NOT NULL AND trim(city) <> ''
|
||||||
AND realestate_type_code = '002001003000'
|
AND realestate_type_code = '002001003000'
|
||||||
AND area BETWEEN 18 AND 200
|
AND area BETWEEN 18 AND 200
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue