Три куска, каждый нужен, чтобы поиск и оценка по Москве заработали end-to-end. 1. Импортёр msk_raw -> listings (app/tasks/msk_raw_import.py). Переиспользует штатный save_listings из кита: писатель уже параметризован регионом, свой не нужен. payload в msk_raw — сериализованный ScrapedLot один в один, так что импорт сводится к сборке модели и вызову писателя. Москва отбирается по префиксу административного округа в адресе, а не по bbox. Причина: адрес Циан не содержит города, а границы региона 77 захватывают ближний пояс области. Замер по проду: с округом 35 552, все внутри bbox 77; без округа внутри bbox 17 576 — это область. Отдельно отсекаются 212 карточек с адресом «Екатеринбург (Cian)», артефакт парсера. listing_segment ПЕРЕСЧИТЫВАЕТСЯ перед записью, а не копируется из payload. Кит ставит novostroyki по одному наличию offer.newbuilding.id. Замер по всем 60 464 карточкам: is_from_developer=true у НУЛЯ, false у 29 000, отсутствует у 31 464. Застройщик не продаёт ни одной карточки корпуса. В проде есть гвард (estimator.py): в аналоги идут строки только с listing_segment IS NULL или 'vtorichka' — копирование метки как есть выбросило бы 29 000 строк из подбора. Авито импортируется только с явным --allow-unfiltered: в его адресе нет ни города, ни округа, координат нет ни у одной из 50 335 карточек, отличить область от Москвы нечем. Прогон на проде: прочитано 60 464, записано 35 552, все с геометрией, 35 182 привязаны к дому, создано 10 960 домов. Повторный проход строки не дублирует — idempotency на dedup_hash, проверено. 2. Подсказки адреса стали региональными (geocoder.suggest, api/v1/geocode). Раньше suggest вообще не принимал регион: DaData звалась с жёстким region='Свердловская', Nominatim — с viewbox 66-го и bounded=1. Московский адрес давал ПУСТОЙ список молча, без ошибки; в коде это уже было описано как известный баг. Механику по регионам переиспользовали из geocode(), вторую не писали. Кадастровый тир для не-66 не зовётся: он на ЕКБ-данных. 3. Оценка перестала геокодировать Москву свердловским скоупом (estimator). geocode() звалась без региона, то есть с дефолтом 66, и московский адрес возвращал бы пустую оценку с причиной address_not_geocoded даже с рабочими подсказками. Регион запроса определяется по координатам через реестр, затем по city_hint, затем дефолт. Fast-path клиентских координат стал региононезависимым: OBLAST66_BBOX и REGIONS[66].bbox_region совпадают байт-в-байт, поэтому для 66 поведение прежнее, добавились координаты Москвы. Регресс-нейтральность по Свердловской области — главный критерий всех трёх кусков. Тесты: 1160 passed по затронутым областям. Известные ограничения. Границы 77 захватывают ближний пояс области, Химки резолвятся в Москву. У региона 77 нет ни одного тира обогащения, оценка поедет на аналогах и сделках. Ценовая полоса по Москве одна на весь город. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01VQ8jqr4SFirX5tFLwdSrXh
322 lines
14 KiB
Python
322 lines
14 KiB
Python
"""Импорт московского сырья (`msk_raw.*_latest`) в `listings`.
|
||
|
||
Сырьё собрано отдельным коллектором и лежит в прод-схеме `msk_raw`: каждая строка
|
||
несёт `payload` — сериализованный `ScrapedLot` один в один (те же 54 ключа, что и
|
||
поля модели, см. `scraper_kit/base.py`). Свой писатель поэтому не нужен: собираем
|
||
`ScrapedLot(**payload)` и отдаём в штатный `save_listings(..., region_code=77)`.
|
||
|
||
Отбор Москвы (source=cian). Адрес карточки Циана города НЕ содержит, зато
|
||
начинается с округа: «ЦАО, ...», «СВАО, ...». По этому префиксу Москва и
|
||
опознаётся. Замер по проду (60 464 карточки): с округом — 35 551, ВСЕ внутри
|
||
bbox региона 77; без округа внутри bbox — 17 576 (это Московская область, регион
|
||
50, которого в реестре ещё нет, в этот импорт не берём); без округа вне bbox —
|
||
7 337. Отдельно 212 карточек с адресом вида «Екатеринбург (Cian)» — артефакт
|
||
парсера, считаются своим счётчиком, чтобы не растворяться в «не Москва».
|
||
|
||
Отбор Москвы (source=avito) НЕВОЗМОЖЕН по адресу: у Авито адрес — голая улица с
|
||
домом («Варшавское ш.,62к1»), ни города, ни округа, и координат нет НИ У ОДНОЙ
|
||
карточки. Поэтому:
|
||
* префиксный фильтр к Авито не применяется — он отбросил бы 100% строк;
|
||
* запись Авито требует явного `--allow-unfiltered`: молча залить в регион 77
|
||
вперемешку Москву и область — хуже, чем не залить ничего;
|
||
* строки Авито лягут БЕЗ geom (lat/lon пусты) — они не попадут в radius-подбор
|
||
аналогов estimator'а, пока их не догеокодит `geocode_missing`.
|
||
|
||
Пересчёт `listing_segment` (пункт, ради которого нельзя копировать payload как
|
||
есть). Кит ставит 'novostroyki' по одному лишь наличию `offer.newbuilding.id`,
|
||
то есть по ссылке на ЖК, а не по продаже застройщиком. Замер по всем 60 464:
|
||
`raw_payload.is_from_developer` = true у НУЛЯ карточек, false у 29 000,
|
||
отсутствует у 31 464 — застройщик в этом корпусе не продаёт ничего, это вся
|
||
вторичка. В estimator'е стоит гвард (`estimator.py:5992-5995`): в аналоги идут
|
||
только строки с `listing_segment IS NULL` или 'vtorichka'. Скопируй мы метку
|
||
кита — 29 000 карточек выпали бы из подбора. Поэтому метка считается заново:
|
||
is_from_developer is True → 'novostroyki', иначе → 'vtorichka'.
|
||
|
||
Идемпотентность — на стороне `save_listings`: он делает upsert
|
||
`ON CONFLICT (dedup_hash) DO UPDATE` плюс reconcile-UPDATE по
|
||
`(source, source_id)` на случай дрейфа хеша. `dedup_hash` = sha256(source +
|
||
source_id) считает сам кит (`ScrapedLot.compute_dedup_hash`), цена в ключ не
|
||
входит. Повторный прогон поэтому обновляет те же строки, а не плодит дубли;
|
||
курсор идёт по `id` вью, так что порядок и полнота обхода от прогона к прогону
|
||
одинаковы.
|
||
|
||
Запуск:
|
||
python -m app.tasks.msk_raw_import --dry-run
|
||
python -m app.tasks.msk_raw_import --limit 500
|
||
python -m app.tasks.msk_raw_import --source avito --allow-unfiltered
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import argparse
|
||
import logging
|
||
import re
|
||
from dataclasses import dataclass
|
||
|
||
from pydantic import ValidationError
|
||
from scraper_kit.base import ScrapedLot, save_listings
|
||
from sqlalchemy import text
|
||
from sqlalchemy.orm import Session
|
||
|
||
from app.core.db import SessionLocal
|
||
from app.services.scraper_adapters import RealMatcherAdapter
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
MOSCOW_REGION_CODE = 77
|
||
MOSCOW_CITY = "Москва"
|
||
DEFAULT_BATCH_SIZE = 500
|
||
|
||
# Префиксы административных округов Москвы — единственный признак города в адресе
|
||
# карточки Циана (сам город в адрес не попадает).
|
||
MOSCOW_OKRUGS = (
|
||
"ЦАО",
|
||
"САО",
|
||
"СВАО",
|
||
"ВАО",
|
||
"ЮВАО",
|
||
"ЮАО",
|
||
"ЮЗАО",
|
||
"ЗАО",
|
||
"СЗАО",
|
||
"ЗелАО",
|
||
"НАО",
|
||
"ТАО",
|
||
)
|
||
# Lookahead вместо \b: следом за округом идёт запятая/пробел, но НЕ буква — иначе
|
||
# «ЗАО» матчило бы начало гипотетического «ЗАОзёрная».
|
||
_MOSCOW_OKRUG_RE = re.compile(
|
||
r"^(?:" + "|".join(MOSCOW_OKRUGS) + r")(?![А-Яа-яЁёA-Za-z])",
|
||
)
|
||
# Артефакт парсера: адрес вида «Екатеринбург (Cian)» в московском корпусе.
|
||
_ARTIFACT_RE = re.compile(r"Екатеринбург", re.IGNORECASE)
|
||
|
||
# Вью-источники. Только whitelist: имя подставляется в SQL текстом, параметром
|
||
# идентификатор не передать.
|
||
SOURCE_VIEWS = {
|
||
"cian": "msk_raw.cian_latest",
|
||
"avito": "msk_raw.avito_latest",
|
||
}
|
||
|
||
_PAGE_SQL = """
|
||
SELECT id, payload
|
||
FROM {view}
|
||
WHERE id > :after
|
||
ORDER BY id
|
||
LIMIT :limit
|
||
"""
|
||
|
||
|
||
@dataclass
|
||
class ImportCounters:
|
||
"""Разбор прогона. Числа обязаны сходиться, см. `check()`."""
|
||
|
||
read: int = 0
|
||
skipped_artifact: int = 0
|
||
skipped_not_moscow: int = 0
|
||
skipped_invalid: int = 0
|
||
selected: int = 0
|
||
inserted: int = 0
|
||
updated: int = 0
|
||
|
||
@property
|
||
def written(self) -> int:
|
||
return self.inserted + self.updated
|
||
|
||
@property
|
||
def writer_skipped(self) -> int:
|
||
"""Отобрано, но писатель строку не тронул.
|
||
|
||
`save_listings` возвращает только (inserted, updated); неизменные строки,
|
||
уже виденные сегодня, он пропускает своим гейтом (#2992). Остаток честно
|
||
показываем отдельно, а не растворяем в «записано».
|
||
"""
|
||
return self.selected - self.written
|
||
|
||
def check(self) -> bool:
|
||
return (
|
||
self.read
|
||
== self.selected
|
||
+ self.skipped_artifact
|
||
+ self.skipped_not_moscow
|
||
+ self.skipped_invalid
|
||
)
|
||
|
||
|
||
def is_artifact_address(address: str | None) -> bool:
|
||
"""Адрес чужого города в московском корпусе (артефакт парсера)."""
|
||
return bool(address) and _ARTIFACT_RE.search(address) is not None
|
||
|
||
|
||
def is_moscow_address(address: str | None) -> bool:
|
||
"""Москва опознаётся префиксом административного округа."""
|
||
if not address:
|
||
return False
|
||
return _MOSCOW_OKRUG_RE.match(address.strip()) is not None
|
||
|
||
|
||
def recompute_listing_segment(payload: dict) -> str:
|
||
"""Заново считаем сегмент: 'novostroyki' только при продаже застройщиком.
|
||
|
||
Обоснование — в докстринге модуля: метка кита означает лишь ссылку на ЖК.
|
||
"""
|
||
raw = payload.get("raw_payload") or {}
|
||
if not isinstance(raw, dict):
|
||
return "vtorichka"
|
||
return "novostroyki" if raw.get("is_from_developer") is True else "vtorichka"
|
||
|
||
|
||
def build_lot(payload: dict) -> ScrapedLot:
|
||
"""`payload` → `ScrapedLot` с пересчитанным сегментом.
|
||
|
||
Ключи, которых в модели нет, отбрасываем явно (по `model_fields`), а не
|
||
полагаемся на настройку extra у pydantic-модели.
|
||
"""
|
||
known = {k: v for k, v in payload.items() if k in ScrapedLot.model_fields}
|
||
known["listing_segment"] = recompute_listing_segment(payload)
|
||
return ScrapedLot(**known)
|
||
|
||
|
||
def _iter_pages(db: Session, view: str, *, batch_size: int, limit: int | None):
|
||
"""Keyset-пагинация по `id` — весь корпус в память не тянем."""
|
||
after = 0
|
||
taken = 0
|
||
sql = text(_PAGE_SQL.format(view=view))
|
||
while True:
|
||
page_size = batch_size
|
||
if limit is not None:
|
||
page_size = min(batch_size, limit - taken)
|
||
if page_size <= 0:
|
||
return
|
||
rows = db.execute(sql, {"after": after, "limit": page_size}).mappings().all()
|
||
if not rows:
|
||
return
|
||
after = rows[-1]["id"]
|
||
taken += len(rows)
|
||
yield rows
|
||
|
||
|
||
def import_msk_raw(
|
||
db: Session,
|
||
*,
|
||
source: str = "cian",
|
||
batch_size: int = DEFAULT_BATCH_SIZE,
|
||
limit: int | None = None,
|
||
dry_run: bool = False,
|
||
allow_unfiltered: bool = False,
|
||
) -> ImportCounters:
|
||
"""Переливает сырьё `msk_raw` в `listings`. Коммит — на каждом батче."""
|
||
view = SOURCE_VIEWS[source]
|
||
counters = ImportCounters()
|
||
matcher = RealMatcherAdapter()
|
||
|
||
filter_by_okrug = source == "cian"
|
||
if not filter_by_okrug:
|
||
# У Авито в адресе нет ни города, ни округа, и нет координат — отсечь
|
||
# область нечем. Пишем только по явному разрешению.
|
||
if not (dry_run or allow_unfiltered):
|
||
raise SystemExit(
|
||
f"source={source}: адрес не содержит признака города, Москву от "
|
||
"области не отличить. Нужен --allow-unfiltered (или --dry-run)."
|
||
)
|
||
logger.warning(
|
||
"source=%s: фильтр по округу НЕ применяется (в адресе нет города); "
|
||
"строки лягут без geom — координат нет ни у одной карточки",
|
||
source,
|
||
)
|
||
|
||
for rows in _iter_pages(db, view, batch_size=batch_size, limit=limit):
|
||
lots: list[ScrapedLot] = []
|
||
for row in rows:
|
||
counters.read += 1
|
||
payload = row["payload"] or {}
|
||
address = payload.get("address")
|
||
if is_artifact_address(address):
|
||
counters.skipped_artifact += 1
|
||
continue
|
||
if filter_by_okrug and not is_moscow_address(address):
|
||
counters.skipped_not_moscow += 1
|
||
continue
|
||
try:
|
||
lots.append(build_lot(payload))
|
||
except ValidationError as exc:
|
||
counters.skipped_invalid += 1
|
||
logger.warning("msk_raw id=%s не собрался в ScrapedLot: %s", row["id"], exc)
|
||
|
||
counters.selected += len(lots)
|
||
if dry_run or not lots:
|
||
continue
|
||
|
||
inserted, updated = save_listings(
|
||
db,
|
||
lots,
|
||
matcher=matcher,
|
||
region_code=MOSCOW_REGION_CODE,
|
||
city=MOSCOW_CITY,
|
||
)
|
||
counters.inserted += inserted
|
||
counters.updated += updated
|
||
db.commit() # батч зафиксирован — обрыв не отматывает всю работу
|
||
logger.info(
|
||
"msk_raw %s: прочитано=%d отобрано=%d записано=%d (new=%d upd=%d)",
|
||
source,
|
||
counters.read,
|
||
counters.selected,
|
||
counters.written,
|
||
counters.inserted,
|
||
counters.updated,
|
||
)
|
||
|
||
logger.info(
|
||
"msk_raw %s ИТОГ%s: прочитано=%d отобрано=%d записано=%d "
|
||
"(new=%d upd=%d, писатель пропустил=%d) | пропущено: не Москва=%d "
|
||
"артефакт=%d невалидный payload=%d | сходится=%s",
|
||
source,
|
||
" (dry-run)" if dry_run else "",
|
||
counters.read,
|
||
counters.selected,
|
||
counters.written,
|
||
counters.inserted,
|
||
counters.updated,
|
||
counters.writer_skipped if not dry_run else 0,
|
||
counters.skipped_not_moscow,
|
||
counters.skipped_artifact,
|
||
counters.skipped_invalid,
|
||
counters.check(),
|
||
)
|
||
return counters
|
||
|
||
|
||
def main() -> None:
|
||
logging.basicConfig(
|
||
level=logging.INFO,
|
||
format="%(asctime)s %(levelname)s %(name)s: %(message)s",
|
||
)
|
||
parser = argparse.ArgumentParser(description="Импорт сырья msk_raw в listings (регион 77)")
|
||
parser.add_argument("--source", choices=sorted(SOURCE_VIEWS), default="cian")
|
||
parser.add_argument("--batch-size", type=int, default=DEFAULT_BATCH_SIZE)
|
||
parser.add_argument("--limit", type=int, default=None, help="обработать не больше N карточек")
|
||
parser.add_argument("--dry-run", action="store_true", help="ничего не пишет, только счётчики")
|
||
parser.add_argument(
|
||
"--allow-unfiltered",
|
||
action="store_true",
|
||
help="разрешить запись источника без признака города в адресе (avito)",
|
||
)
|
||
args = parser.parse_args()
|
||
|
||
db = SessionLocal()
|
||
try:
|
||
import_msk_raw(
|
||
db,
|
||
source=args.source,
|
||
batch_size=args.batch_size,
|
||
limit=args.limit,
|
||
dry_run=args.dry_run,
|
||
allow_unfiltered=args.allow_unfiltered,
|
||
)
|
||
finally:
|
||
db.close()
|
||
|
||
|
||
if __name__ == "__main__":
|
||
main()
|