All checks were successful
CI Trade-In / changes (pull_request) Successful in 7s
CI / changes (pull_request) Successful in 7s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 2m56s
Три находки одного класса из эпика: колонка есть, писатель есть, тест на писателя зелёный — а данные не появляются. Тестами это не ловится по построению, только сверкой с продом. 1. house_suggestions: парсер выбрасывал imageLink, а INSERT не перечислял image_link + area_m2/rooms/floor/total_floors. 25 055 строк с NULL во всех пяти колонках, ~74 дня с миграции 064. Метрики парсятся из title тем же _parse_title, что и у placementHistory. 2. listings_snapshots.status: 'active' у всех 394 299 строк при 55 448 реально неактивных объявлений. Оба места вызова с литералом 'active' честны — там объявление действительно видели; не писал никто ветку «снято». Теперь оба места деактивации пишут снимок 'closed' в ТОЙ ЖЕ транзакции: TTL-задача (data-modifying CTE, все 4 источника через один deactivate_stale_listings) и 404 из avito_detail_backfill. Дата снятия перестаёт быть догадкой. 3. listing_source_events: схема знает 5 типов, писался 1 (price_change, 8288 строк). Дописаны ветки delisted/relisted/edited/first_seen в тот же set-based statement — данные для них уже лежат в снимке. JOIN → LEFT JOIN LATERAL, иначе first_seen недостижим по построению; план #2607 (per-row index point-lookup по idx_lss_source_date) сохранён, проверено EXPLAIN на проде. Счётчики прогона теперь по типам, все пять всегда присутствуют — ровно они показали бы четыре нуля из пяти. Миграция не нужна: все колонки и CHECK уже существуют. Тесты: tests/test_2674_writers_honor_schema.py. Гейты сверяют писателя со СХЕМОЙ (колонки INSERT против CREATE TABLE 064, типы событий против CHECK 079), поэтому ловят и следующую забытую колонку. Фальсификация патч-методом: без фикса 1 — 6 красных, без фикса 2 — 6, без фикса 3 — 4.
802 lines
36 KiB
Python
802 lines
36 KiB
Python
"""T7: House-level IMV backfill service.
|
||
|
||
Iterates houses with imv_status='pending' (or 'transient_error') that have
|
||
lat/lon coordinates and at least one linked listing with rooms+area data.
|
||
|
||
For each house:
|
||
1. Pick median lot-params from linked listings (rooms, area_m2, floor,
|
||
total_floors, house_type).
|
||
2. Enrich address with region prefix (prevents Avito geocoder ambiguity).
|
||
3. Call avito_imv.evaluate_via_imv().
|
||
4. Persist house_imv_evaluations + house_placement_history + house_suggestions.
|
||
5. Mark houses.imv_status accordingly (ok/not_found/no_params/error/transient_error).
|
||
|
||
Resumable: re-running picks up from where the previous run stopped.
|
||
process_houses_imv_batch() IS wired into the scheduler: run_avito_city_sweep()
|
||
(scraper_kit.orchestration.pipeline, via the injected EnrichmentJobs protocol)
|
||
calls it automatically as the final IMV-phase of every avito city sweep.
|
||
backfill_house_imv() (batch/pending-status entrypoint) remains manual-only,
|
||
triggered via admin API.
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import asyncio
|
||
import json
|
||
import logging
|
||
import time
|
||
from collections.abc import Callable
|
||
from dataclasses import dataclass, field
|
||
from typing import Literal
|
||
|
||
from scraper_kit.browser_fetcher import BrowserFetcher
|
||
from scraper_kit.house_type_normalizer import normalize_house_type
|
||
|
||
# #2337 (Group E4, эпик #2277): переключено на scraper_kit — тот же периметр риска,
|
||
# что и estimator.py (обе точки читают/пишут house_imv_evaluations, #651 IMV/Yandex
|
||
# blend). config=RealScraperConfig() обязателен для evaluate_via_imv — без него kit
|
||
# молча уходит без прокси (прямое подключение) вместо настроенного (см. #2334).
|
||
from scraper_kit.providers.avito.imv import (
|
||
IMVAddressNotFoundError,
|
||
IMVAuthError,
|
||
IMVEvaluation,
|
||
IMVTransientError,
|
||
evaluate_via_imv,
|
||
)
|
||
from sqlalchemy import text
|
||
from sqlalchemy.orm import Session
|
||
|
||
from app.core.config import settings
|
||
|
||
# RealScraperConfig НЕ импортируется на уровне модуля: app.services.scraper_adapters
|
||
# сам импортирует backfill_house_imv/process_houses_imv_batch ИЗ этого модуля
|
||
# (RealEnrichmentJobs-адаптер) — top-level import здесь создал бы circular import
|
||
# (проверено: ImportError "cannot import name 'RealScraperConfig' from partially
|
||
# initialized module" при импорте scraper_adapters первым). Ленивый import внутри
|
||
# _process_one_house() ломает цикл — к моменту вызова оба модуля уже полностью
|
||
# инициализированы.
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
# Каждые N домов дёргаем heartbeat-колбэк (если передан) внутри длинного per-house
|
||
# цикла. Без этого долгая IMV-фаза (десятки/сотни домов × request_delay_sec) не
|
||
# обновляет scrape_runs.heartbeat_at дольше ZOMBIE_THRESHOLD_HOURS, и scheduler
|
||
# ложно помечает живой run 'zombie' (после чего терминальный mark_done — no-op). См. #1363.
|
||
_HEARTBEAT_EVERY_N_HOUSES = 5
|
||
|
||
# ── house_type normalisation ─────────────────────────────────────────────────
|
||
|
||
# Ключи — КАНОНИЧНЫЕ значения listings.house_type (после normalize_house_type),
|
||
# значения — вокабуляр Avito IMV.
|
||
_HOUSE_TYPE_TO_IMV: dict[str, str] = {
|
||
"panel": "panel",
|
||
"brick": "brick",
|
||
"monolith": "monolithic",
|
||
"monolith_brick": "monolithic", # Avito API не принимает гибриды
|
||
"block": "block",
|
||
"wood": "wood",
|
||
}
|
||
|
||
|
||
def _map_house_type(raw: str | None) -> str | None:
|
||
"""Наш house_type → вокабуляр Avito IMV. None = тип неизвестен, запрос не шлём.
|
||
|
||
Сырое значение сначала прогоняем через общий normalize_house_type (#2007): он
|
||
знает camelCase-вокабуляр Циана (monolithBrick / gasSilicateBlock /
|
||
aerocreteBlock / stalin / ...) и SCREAMING-вокабуляр Яндекса, а нераспознанное
|
||
('other', 'wireframe', пустое) схлопывает в None. Приведения к нижнему регистру
|
||
тут мало: ключ канона пишется через подчёркивание (monolith_brick), поэтому
|
||
'monolithbrick' в словарь не попадал.
|
||
|
||
#2674: раньше здесь стоял дефолт 'panel' — и когда типа нет вовсе, и когда он
|
||
есть, но не распознан. Панель — почти самый дешёвый класс (медиана по нашим же
|
||
2685 оценкам: block 122.6k < panel 128.8k < brick 131.1k < monolithic 145.9k
|
||
₽/м²), то есть дефолт систематически ЗАНИЖАЛ оценку: на проде 363 дома совсем
|
||
без типа + 75 домов с camelCase-типом (56 из них monolithBrick, −11.7% к
|
||
monolithic) уехали как панель. Теперь неизвестный тип → None → дом помечается
|
||
и запрос к площадке не тратится (см. _process_one_house).
|
||
"""
|
||
canon = normalize_house_type(raw)
|
||
if canon is None:
|
||
return None
|
||
return _HOUSE_TYPE_TO_IMV.get(canon)
|
||
|
||
|
||
def _map_renovation_type(repair_state: str | None) -> str:
|
||
"""listings.repair_state → renovation_type вокабуляра Avito IMV.
|
||
|
||
Переиспользуем _IMV_REPAIR_MAP эстиматора — единственный источник правды для
|
||
этого соответствия (needs_repair→required / standard→cosmetic / good→euro /
|
||
excellent→designer). Импорт ленивый: estimator тянет scraper_adapters, а тот
|
||
импортирует этот модуль (circular — см. блок импортов выше).
|
||
|
||
#2674: раньше здесь стоял литерал 'cosmetic' — все 2685 запросов ушли как
|
||
«косметический ремонт», хотя мода по объявлениям этих же домов совсем другая
|
||
(standard 4564 / good 4118 / needs_repair 2279 / excellent 1631 — косметика
|
||
лишь 36%).
|
||
|
||
Неизвестный ремонт (498 домов из 2685 — ни одного объявления с repair_state)
|
||
ОСТАЁТСЯ 'cosmetic', в отличие от неизвестного типа дома: 'cosmetic'
|
||
(=standard) — это одновременно МОДА и МЕДИАННАЯ категория популяции
|
||
(standard 7984 / good 7116 / needs_repair 4738 / excellent 2562; кумулятивно
|
||
needs_repair 21.2%, +standard 56.8%), то есть наилучшая одиночная догадка.
|
||
У типа дома такой догадки нет: 'panel' — почти край шкалы, а не её середина.
|
||
|
||
Асимметрия осознанная, а не недосмотр: поштучный путь эстиматора при
|
||
неизвестном ремонте IMV вообще не зовёт (estimator.py, `imv_renovation is not
|
||
None`), а домовой дефолтит — иначе теряем ещё ~32% домов очереди поверх тех,
|
||
что уже отсекает неизвестный тип дома.
|
||
"""
|
||
from app.services.estimator import _IMV_REPAIR_MAP # lazy — см. import-блок
|
||
|
||
mapped = _IMV_REPAIR_MAP.get(repair_state)
|
||
if mapped is None and repair_state:
|
||
# Непустое, но незнакомое значение — признак дрейфа вокабуляра на ингесте
|
||
# (сырых repair-значений в listings больше, чем нормализованных). Паритет
|
||
# с house_type_normalizer, который такой случай уже логирует.
|
||
logger.debug("house_imv: unmapped repair_state %r — падаем в 'cosmetic'", repair_state)
|
||
return mapped or "cosmetic"
|
||
|
||
|
||
# ── Region bbox prefix для Avito geocoder ────────────────────────────────────
|
||
|
||
_REGION_BBOX: list[tuple[str, float, float, float, float]] = [
|
||
# (region_name, lat_min, lat_max, lon_min, lon_max)
|
||
("Свердловская область, Екатеринбург", 56.5, 57.1, 59.9, 61.0),
|
||
("Свердловская область", 56.0, 61.0, 58.0, 65.0),
|
||
("Кировская область, Киров", 58.4, 58.8, 49.4, 49.9),
|
||
("Ставропольский край, Ставрополь", 44.9, 45.2, 41.7, 42.2),
|
||
("Тюменская область, Тюмень", 57.0, 57.3, 65.3, 65.8),
|
||
("Челябинская область, Челябинск", 55.0, 55.4, 61.2, 61.7),
|
||
]
|
||
_REGION_DEFAULT = "Свердловская область, Екатеринбург"
|
||
|
||
|
||
def _detect_region_prefix(lat: float | None, lon: float | None) -> str:
|
||
if lat is None or lon is None:
|
||
return _REGION_DEFAULT
|
||
for name, lat_min, lat_max, lon_min, lon_max in _REGION_BBOX:
|
||
if lat_min <= lat <= lat_max and lon_min <= lon <= lon_max:
|
||
return name
|
||
return _REGION_DEFAULT
|
||
|
||
|
||
def _enrich_address_for_imv(raw: str, lat: float | None, lon: float | None) -> str:
|
||
"""Prepend region prefix if address lacks a region marker."""
|
||
addr = (raw or "").strip()
|
||
addr_lower = addr.lower()
|
||
if any(
|
||
m in addr_lower for m in ("область,", " обл.,", " обл,", "край,", "респ.,", "республика,")
|
||
):
|
||
return addr
|
||
prefix = _detect_region_prefix(lat, lon)
|
||
return f"{prefix}, {addr}"
|
||
|
||
|
||
# ── DB helpers ────────────────────────────────────────────────────────────────
|
||
|
||
|
||
def pick_lot_params(db: Session, house_id: int) -> dict:
|
||
"""Pick representative lot-params — median values from linked listings."""
|
||
row = (
|
||
db.execute(
|
||
text("""
|
||
SELECT
|
||
CAST(percentile_cont(0.5) WITHIN GROUP (ORDER BY rooms)
|
||
AS integer) AS rooms,
|
||
CAST(percentile_cont(0.5) WITHIN GROUP (ORDER BY area_m2)
|
||
AS numeric(8,2)) AS area_m2,
|
||
CAST(percentile_cont(0.5) WITHIN GROUP (ORDER BY floor)
|
||
AS integer) AS floor,
|
||
CAST(percentile_cont(0.5) WITHIN GROUP (ORDER BY total_floors)
|
||
AS integer) AS total_floors,
|
||
mode() WITHIN GROUP (ORDER BY house_type) AS house_type,
|
||
mode() WITHIN GROUP (ORDER BY repair_state) AS repair_state
|
||
FROM listings
|
||
WHERE house_id_fk = :hid
|
||
AND rooms IS NOT NULL
|
||
AND area_m2 IS NOT NULL
|
||
"""),
|
||
{"hid": house_id},
|
||
)
|
||
.mappings()
|
||
.first()
|
||
)
|
||
|
||
if row is None or row["rooms"] is None:
|
||
return {}
|
||
|
||
house = (
|
||
db.execute(
|
||
text("SELECT house_type, total_floors FROM houses WHERE id = :hid"),
|
||
{"hid": house_id},
|
||
)
|
||
.mappings()
|
||
.first()
|
||
)
|
||
|
||
rooms_raw = int(row["rooms"])
|
||
floor = int(row["floor"] or 1)
|
||
floor_at_home = int(row["total_floors"] or (house and house["total_floors"]) or 9)
|
||
# floor и floor_at_home считаются независимо (медиана vs default-9). Если у
|
||
# всех listings и у дома total_floors=NULL, floor_at_home молча падает на 9,
|
||
# а медианный floor может быть выше → Avito IMV получает кв. на этаже выше дома.
|
||
# Консервативный clamp: floor не может превышать этажность дома.
|
||
floor = min(floor, floor_at_home)
|
||
return {
|
||
"rooms": max(rooms_raw, 1), # Avito IMV отклоняет rooms=0 (студия)
|
||
"area_m2": float(row["area_m2"]),
|
||
"floor": floor,
|
||
"floor_at_home": floor_at_home,
|
||
"house_type": _map_house_type(row["house_type"] or (house and house["house_type"])),
|
||
"renovation_type": _map_renovation_type(row["repair_state"]),
|
||
# has_balcony/has_loggia остаются константами намеренно (#2674): покрытие
|
||
# listings.has_balcony 13.8%, listings.balcony_loggia 9.4%, и две колонки
|
||
# противоречат друг другу (по has_balcony «есть» у 62%, а по
|
||
# balcony_loggia самый частый случай — loggia 5650 против balcony 2794).
|
||
# Мода по одному-двум объявлениям на таком покрытии — шум, а не данные.
|
||
"has_balcony": True,
|
||
"has_loggia": False,
|
||
}
|
||
|
||
|
||
def save_imv_result(db: Session, house_id: int, params: dict, result: IMVEvaluation) -> None:
|
||
"""Persist evaluation to house_imv_evaluations, house_placement_history, house_suggestions."""
|
||
# 1. House-level evaluation (UPSERT on house_id)
|
||
db.execute(
|
||
text("""
|
||
INSERT INTO house_imv_evaluations (
|
||
house_id, cache_key, recommended_price, lower_price, higher_price,
|
||
market_count, raw_response, fetched_at,
|
||
rooms, area_m2, floor, floor_at_home, house_type,
|
||
renovation_type, has_balcony, has_loggia
|
||
) VALUES (
|
||
:hid, :ck, :rec, :lo, :hi,
|
||
:mc, CAST(:raw AS jsonb), NOW(),
|
||
:rooms, CAST(:area AS numeric), :floor, :fah, :htype,
|
||
:rtype, :balcony, :loggia
|
||
)
|
||
ON CONFLICT (house_id) DO UPDATE SET
|
||
cache_key = EXCLUDED.cache_key,
|
||
recommended_price = EXCLUDED.recommended_price,
|
||
lower_price = EXCLUDED.lower_price,
|
||
higher_price = EXCLUDED.higher_price,
|
||
market_count = EXCLUDED.market_count,
|
||
raw_response = EXCLUDED.raw_response,
|
||
fetched_at = NOW(),
|
||
rooms = EXCLUDED.rooms,
|
||
area_m2 = EXCLUDED.area_m2,
|
||
floor = EXCLUDED.floor,
|
||
floor_at_home = EXCLUDED.floor_at_home,
|
||
house_type = EXCLUDED.house_type,
|
||
renovation_type = EXCLUDED.renovation_type,
|
||
has_balcony = EXCLUDED.has_balcony,
|
||
has_loggia = EXCLUDED.has_loggia
|
||
"""),
|
||
{
|
||
"hid": house_id,
|
||
"ck": result.cache_key,
|
||
"rec": result.recommended_price,
|
||
"lo": result.lower_price,
|
||
"hi": result.higher_price,
|
||
"mc": result.market_count,
|
||
"raw": (
|
||
json.dumps(result.raw_response, ensure_ascii=False) if result.raw_response else None
|
||
),
|
||
"rooms": params["rooms"],
|
||
"area": params["area_m2"],
|
||
"floor": params["floor"],
|
||
"fah": params["floor_at_home"],
|
||
"htype": params["house_type"],
|
||
"rtype": params["renovation_type"],
|
||
"balcony": params["has_balcony"],
|
||
"loggia": params["has_loggia"],
|
||
},
|
||
)
|
||
|
||
# 2. Placement history items
|
||
for item in result.placement_history:
|
||
db.execute(
|
||
text("""
|
||
INSERT INTO house_placement_history (
|
||
source, house_id, ext_item_id, title, rooms, area_m2,
|
||
floor, total_floors, start_price, start_price_date,
|
||
last_price, last_price_date, removed_date, exposure_days,
|
||
raw_payload, scraped_at
|
||
) VALUES (
|
||
'avito_imv', :hid, :ext, :title, :rooms, CAST(:area AS numeric),
|
||
:floor, :total_floors, :sp, :spd, :lp, :lpd, :rmd, :exp,
|
||
CAST(:raw AS jsonb), NOW()
|
||
)
|
||
ON CONFLICT (source, ext_item_id) DO UPDATE SET
|
||
house_id = EXCLUDED.house_id,
|
||
title = EXCLUDED.title,
|
||
rooms = EXCLUDED.rooms,
|
||
area_m2 = EXCLUDED.area_m2,
|
||
floor = EXCLUDED.floor,
|
||
total_floors = EXCLUDED.total_floors,
|
||
start_price = EXCLUDED.start_price,
|
||
start_price_date = EXCLUDED.start_price_date,
|
||
last_price = EXCLUDED.last_price,
|
||
last_price_date = EXCLUDED.last_price_date,
|
||
removed_date = EXCLUDED.removed_date,
|
||
exposure_days = EXCLUDED.exposure_days,
|
||
raw_payload = EXCLUDED.raw_payload,
|
||
scraped_at = NOW()
|
||
"""),
|
||
{
|
||
"hid": house_id,
|
||
"ext": item.ext_item_id,
|
||
"title": item.title,
|
||
"rooms": item.rooms,
|
||
"area": item.area_m2,
|
||
"floor": item.floor,
|
||
"total_floors": item.total_floors,
|
||
"sp": item.start_price,
|
||
"spd": item.start_price_date,
|
||
"lp": item.last_price,
|
||
"lpd": item.last_price_date,
|
||
"rmd": item.removed_date,
|
||
"exp": item.exposure_days,
|
||
"raw": (
|
||
json.dumps(item.raw_payload, ensure_ascii=False) if item.raw_payload else None
|
||
),
|
||
},
|
||
)
|
||
|
||
# 3. Suggestions
|
||
# #2674: до этого фикса в INSERT не входили image_link + area_m2/rooms/floor/
|
||
# total_floors — колонки есть с миграции 064, но писатель их не заполнял
|
||
# (25 055 строк на проде с NULL во всех пяти). Ссылка на фото приходит в
|
||
# suggestions.items[].imageLink, метрики квартиры парсятся из title.
|
||
for sug in result.suggestions:
|
||
db.execute(
|
||
text("""
|
||
INSERT INTO house_suggestions (
|
||
house_id, ext_item_id, title, address, price_rub,
|
||
area_m2, rooms, floor, total_floors,
|
||
exposure_days, publish_date,
|
||
item_link, image_link, metro_name, metro_distance,
|
||
has_good_price_badge, raw_payload, fetched_at
|
||
) VALUES (
|
||
:hid, :ext, :title, :addr, :price,
|
||
CAST(:area AS numeric), :rooms, :floor, :total_floors,
|
||
:exp, :pdate,
|
||
:link, :img, :mname, :mdist,
|
||
:gpb, CAST(:raw AS jsonb), NOW()
|
||
)
|
||
ON CONFLICT (house_id, ext_item_id) DO UPDATE SET
|
||
title = EXCLUDED.title,
|
||
price_rub = EXCLUDED.price_rub,
|
||
area_m2 = EXCLUDED.area_m2,
|
||
rooms = EXCLUDED.rooms,
|
||
floor = EXCLUDED.floor,
|
||
total_floors = EXCLUDED.total_floors,
|
||
exposure_days = EXCLUDED.exposure_days,
|
||
publish_date = EXCLUDED.publish_date,
|
||
item_link = EXCLUDED.item_link,
|
||
image_link = EXCLUDED.image_link,
|
||
metro_name = EXCLUDED.metro_name,
|
||
metro_distance = EXCLUDED.metro_distance,
|
||
has_good_price_badge = EXCLUDED.has_good_price_badge,
|
||
raw_payload = EXCLUDED.raw_payload,
|
||
fetched_at = NOW()
|
||
"""),
|
||
{
|
||
"hid": house_id,
|
||
"ext": sug.ext_item_id,
|
||
"title": sug.title,
|
||
"addr": sug.address,
|
||
"price": sug.price_rub,
|
||
"area": sug.area_m2,
|
||
"rooms": sug.rooms,
|
||
"floor": sug.floor,
|
||
"total_floors": sug.total_floors,
|
||
"exp": sug.exposure_days,
|
||
"pdate": sug.publish_date,
|
||
"link": sug.item_url,
|
||
"img": sug.image_link,
|
||
"mname": sug.metro_name,
|
||
"mdist": sug.metro_distance,
|
||
"gpb": sug.has_good_price_badge,
|
||
"raw": (
|
||
json.dumps(sug.raw_payload, ensure_ascii=False) if sug.raw_payload else None
|
||
),
|
||
},
|
||
)
|
||
|
||
# 4. Mark house success
|
||
db.execute(
|
||
text("""
|
||
UPDATE houses
|
||
SET imv_status = 'ok',
|
||
last_imv_attempt_at = NOW(),
|
||
imv_error_reason = NULL
|
||
WHERE id = :hid
|
||
"""),
|
||
{"hid": house_id},
|
||
)
|
||
|
||
|
||
def _mark_status(
|
||
db: Session,
|
||
house_id: int,
|
||
status: str,
|
||
reason: str | None = None,
|
||
) -> None:
|
||
db.execute(
|
||
text("""
|
||
UPDATE houses
|
||
SET imv_status = :s,
|
||
last_imv_attempt_at = NOW(),
|
||
imv_error_reason = :r
|
||
WHERE id = :hid
|
||
"""),
|
||
{"hid": house_id, "s": status, "r": reason},
|
||
)
|
||
db.commit()
|
||
|
||
|
||
# ── Top-level orchestrator ────────────────────────────────────────────────────
|
||
|
||
_IMVStatus = Literal[
|
||
"ok", "no_params", "no_address", "not_found", "auth_error", "transient", "error"
|
||
]
|
||
|
||
|
||
@dataclass
|
||
class HouseIMVBackfillResult:
|
||
checked: int = 0
|
||
saved: int = 0
|
||
skipped: int = 0
|
||
errors: int = 0
|
||
duration_sec: float = field(default=0.0)
|
||
status_counts: dict[str, int] = field(default_factory=dict)
|
||
|
||
|
||
def _beat(heartbeat: Callable[[], None] | None) -> None:
|
||
"""Best-effort вызов heartbeat-колбэка. Сбой heartbeat не должен ронять backfill."""
|
||
if heartbeat is None:
|
||
return
|
||
try:
|
||
heartbeat()
|
||
except Exception:
|
||
logger.warning("house_imv_backfill: heartbeat callback failed (ignored)", exc_info=True)
|
||
|
||
|
||
async def backfill_house_imv(
|
||
db: Session,
|
||
*,
|
||
batch_size: int = 50,
|
||
request_delay_sec: float = 5.0,
|
||
only_status: str = "pending",
|
||
house_id: int | None = None,
|
||
heartbeat: Callable[[], None] | None = None,
|
||
) -> HouseIMVBackfillResult:
|
||
"""Run Avito IMV evaluation for each house in scope, save results.
|
||
|
||
Args:
|
||
db: SQLAlchemy session.
|
||
batch_size: max houses to process (ignored when house_id given).
|
||
request_delay_sec: sleep between Avito API calls (default 5s — anti-bot).
|
||
only_status: process houses with this imv_status (default 'pending').
|
||
Use 'transient_error' to retry failures.
|
||
house_id: process a single specific house (debug).
|
||
heartbeat: optional callback дёргается каждые _HEARTBEAT_EVERY_N_HOUSES
|
||
домов — caller обновляет scrape_runs.heartbeat_at, чтобы reap_zombies
|
||
не пометил живой долгий run 'zombie' (#1363). Best-effort: исключения
|
||
из колбэка логируются и не прерывают backfill.
|
||
|
||
Returns:
|
||
HouseIMVBackfillResult with aggregate counters.
|
||
"""
|
||
result = HouseIMVBackfillResult()
|
||
t0 = time.time()
|
||
|
||
if house_id is not None:
|
||
rows = (
|
||
db.execute(
|
||
text("""
|
||
SELECT id, address, full_address, lat, lon
|
||
FROM houses
|
||
WHERE id = CAST(:hid AS bigint)
|
||
AND lat IS NOT NULL
|
||
AND lon IS NOT NULL
|
||
"""),
|
||
{"hid": house_id},
|
||
)
|
||
.mappings()
|
||
.all()
|
||
)
|
||
else:
|
||
rows = (
|
||
db.execute(
|
||
text("""
|
||
SELECT id, address, full_address, lat, lon
|
||
FROM houses
|
||
WHERE imv_status = :status
|
||
AND lat IS NOT NULL
|
||
AND lon IS NOT NULL
|
||
AND address IS NOT NULL
|
||
ORDER BY last_imv_attempt_at NULLS FIRST, id
|
||
LIMIT :batch
|
||
"""),
|
||
{"status": only_status, "batch": batch_size},
|
||
)
|
||
.mappings()
|
||
.all()
|
||
)
|
||
|
||
result.checked = len(rows)
|
||
if not rows:
|
||
logger.info("house_imv_backfill: nothing to process (status=%r)", only_status)
|
||
result.duration_sec = time.time() - t0
|
||
return result
|
||
|
||
logger.info(
|
||
"house_imv_backfill: %d houses (status=%r delay=%.1fs)",
|
||
result.checked,
|
||
only_status,
|
||
request_delay_sec,
|
||
)
|
||
|
||
async def _run_loop(_bf: BrowserFetcher | None) -> None:
|
||
for i, house in enumerate(rows):
|
||
hid: int = house["id"]
|
||
status_str = await _process_one_house(db, dict(house), browser_fetcher=_bf)
|
||
|
||
result.status_counts[status_str] = result.status_counts.get(status_str, 0) + 1
|
||
if status_str == "ok":
|
||
result.saved += 1
|
||
elif status_str in ("no_params", "no_address"):
|
||
result.skipped += 1
|
||
else:
|
||
result.errors += 1
|
||
|
||
logger.info(
|
||
"[%d/%d] house_id=%d status=%s",
|
||
i + 1,
|
||
result.checked,
|
||
hid,
|
||
status_str,
|
||
)
|
||
|
||
# Периодический heartbeat внутри длинного цикла — иначе reap_zombies
|
||
# может пометить живой run 'zombie' за время IMV-фазы (#1363).
|
||
if (i + 1) % _HEARTBEAT_EVERY_N_HOUSES == 0:
|
||
_beat(heartbeat)
|
||
|
||
if i < len(rows) - 1:
|
||
if status_str == "auth_error":
|
||
logger.warning("IMV auth error — extra delay 60s before next house")
|
||
await asyncio.sleep(60)
|
||
else:
|
||
await asyncio.sleep(request_delay_sec)
|
||
|
||
# #915 Stage 3: при флаге открываем ОДИН browser на весь батч (warmed-сессия
|
||
# + прокси переиспользуются всеми домами; обходит datacenter-403, #562/#853).
|
||
# Флаг OFF → _bf=None → evaluate_via_imv делает свою curl-сессию как раньше
|
||
# (поведение байт-в-байт идентично доспринтовому).
|
||
if settings.avito_imv_use_browser_fetcher:
|
||
async with BrowserFetcher(source="avito", endpoint=settings.browser_http_endpoint) as _bf:
|
||
await _run_loop(_bf)
|
||
else:
|
||
await _run_loop(None)
|
||
|
||
result.duration_sec = time.time() - t0
|
||
logger.info(
|
||
"house_imv_backfill done: checked=%d saved=%d skipped=%d errors=%d %.1fs %s",
|
||
result.checked,
|
||
result.saved,
|
||
result.skipped,
|
||
result.errors,
|
||
result.duration_sec,
|
||
result.status_counts,
|
||
)
|
||
return result
|
||
|
||
|
||
async def process_houses_imv_batch(
|
||
db: Session,
|
||
house_ids: set[int],
|
||
*,
|
||
request_delay_sec: float = 5.0,
|
||
heartbeat: Callable[[], None] | None = None,
|
||
) -> HouseIMVBackfillResult:
|
||
"""Обработать конкретный набор house_id (для sweep-интеграции).
|
||
|
||
Зеркало backfill_house_imv, но работает по явному списку house_id
|
||
(не по imv_status-фильтру). Предназначен для финальной IMV-фазы
|
||
run_avito_city_sweep — обрабатываем только дома, тронутые за этот sweep.
|
||
|
||
Дома без lat/lon или address пропускаются (статус no_params/no_address).
|
||
Дома уже с imv_status='ok' тоже пропускаются — UPSERT всё равно обновит,
|
||
поэтому лучше не гонять лишних API-запросов, если данные свежие.
|
||
Для форс-обновления — оставь это решение caller'у (sweep по умолчанию
|
||
обходит дома с ok, обновляя только pending/not_found/error/transient_error).
|
||
|
||
heartbeat: optional callback дёргается каждые _HEARTBEAT_EVERY_N_HOUSES домов,
|
||
чтобы scrape_runs.heartbeat_at обновлялся в течение долгой IMV-фазы и
|
||
reap_zombies не пометил живой run 'zombie' (#1363). Best-effort.
|
||
"""
|
||
result = HouseIMVBackfillResult()
|
||
t0 = time.time()
|
||
|
||
if not house_ids:
|
||
result.duration_sec = time.time() - t0
|
||
return result
|
||
|
||
ids_list = list(house_ids)
|
||
|
||
# Загружаем только дома, у которых есть координаты + адрес + imv_status != 'ok'
|
||
# (уже оцененные дома пропускаем — не тратим API-квоту повторно)
|
||
placeholders = ", ".join(f":hid_{i}" for i in range(len(ids_list)))
|
||
params: dict[str, object] = {f"hid_{i}": hid for i, hid in enumerate(ids_list)}
|
||
rows = (
|
||
db.execute(
|
||
text(f"""
|
||
SELECT id, address, full_address, lat, lon
|
||
FROM houses
|
||
WHERE id IN ({placeholders})
|
||
AND lat IS NOT NULL
|
||
AND lon IS NOT NULL
|
||
AND address IS NOT NULL
|
||
AND COALESCE(imv_status, 'pending') != 'ok'
|
||
ORDER BY id
|
||
"""),
|
||
params,
|
||
)
|
||
.mappings()
|
||
.all()
|
||
)
|
||
|
||
result.checked = len(rows)
|
||
if not rows:
|
||
logger.info(
|
||
"process_houses_imv_batch: nothing to process (%d house_ids, all ok/no-coords)",
|
||
len(house_ids),
|
||
)
|
||
result.duration_sec = time.time() - t0
|
||
return result
|
||
|
||
logger.info(
|
||
"process_houses_imv_batch: %d/%d houses eligible (delay=%.1fs)",
|
||
result.checked,
|
||
len(house_ids),
|
||
request_delay_sec,
|
||
)
|
||
|
||
for i, house in enumerate(rows):
|
||
hid: int = house["id"]
|
||
status_str = await _process_one_house(db, dict(house))
|
||
|
||
result.status_counts[status_str] = result.status_counts.get(status_str, 0) + 1
|
||
if status_str == "ok":
|
||
result.saved += 1
|
||
elif status_str in ("no_params", "no_address"):
|
||
result.skipped += 1
|
||
else:
|
||
result.errors += 1
|
||
|
||
logger.info(
|
||
"[%d/%d] imv_batch house_id=%d status=%s",
|
||
i + 1,
|
||
result.checked,
|
||
hid,
|
||
status_str,
|
||
)
|
||
|
||
# Периодический heartbeat внутри длинного цикла — иначе reap_zombies
|
||
# может пометить живой run 'zombie' за время IMV-фазы (#1363).
|
||
if (i + 1) % _HEARTBEAT_EVERY_N_HOUSES == 0:
|
||
_beat(heartbeat)
|
||
|
||
if i < len(rows) - 1:
|
||
if status_str == "auth_error":
|
||
logger.warning("IMV auth error — extra delay 60s before next house")
|
||
await asyncio.sleep(60)
|
||
else:
|
||
await asyncio.sleep(request_delay_sec)
|
||
|
||
result.duration_sec = time.time() - t0
|
||
logger.info(
|
||
"process_houses_imv_batch done: checked=%d saved=%d skipped=%d errors=%d %.1fs %s",
|
||
result.checked,
|
||
result.saved,
|
||
result.skipped,
|
||
result.errors,
|
||
result.duration_sec,
|
||
result.status_counts,
|
||
)
|
||
return result
|
||
|
||
|
||
async def _process_one_house(
|
||
db: Session, house: dict, *, browser_fetcher: BrowserFetcher | None = None
|
||
) -> str:
|
||
"""Process a single house. Returns final status string.
|
||
|
||
browser_fetcher: при #915 Stage 3 backfill пробрасывает один общий
|
||
BrowserFetcher на весь батч → evaluate_via_imv ходит через /fetch-json
|
||
sidecar (обходит datacenter-403). None → curl_cffi-путь (как раньше).
|
||
"""
|
||
from app.services.scraper_adapters import RealScraperConfig # lazy — см. import-блок
|
||
|
||
hid: int = house["id"]
|
||
|
||
params = pick_lot_params(db, hid)
|
||
if not params:
|
||
_mark_status(db, hid, "no_params", "no listings with rooms+area")
|
||
return "no_params"
|
||
|
||
# #2674: тип дома неизвестен (нет ни в объявлениях, ни в houses — либо
|
||
# вокабуляр не распознан). Раньше такой дом молча уезжал как 'panel' и
|
||
# занижал оценку. Лучше не тратить запрос и честно пометить дом — тот же
|
||
# путь, что и при отсутствии комнат/площади.
|
||
if params["house_type"] is None:
|
||
_mark_status(db, hid, "no_params", "unknown house_type")
|
||
return "no_params"
|
||
|
||
address = house.get("address") or house.get("full_address")
|
||
if not address:
|
||
_mark_status(db, hid, "no_address", "house.address is NULL")
|
||
return "no_address"
|
||
|
||
enriched = _enrich_address_for_imv(address, house.get("lat"), house.get("lon"))
|
||
if enriched != address:
|
||
logger.debug("house_imv: enriched address house=%d %r -> %r", hid, address, enriched)
|
||
|
||
try:
|
||
eval_result = await evaluate_via_imv(
|
||
address=enriched,
|
||
browser_fetcher=browser_fetcher,
|
||
config=RealScraperConfig(),
|
||
**params,
|
||
)
|
||
except IMVAddressNotFoundError as exc:
|
||
_mark_status(db, hid, "not_found", str(exc)[:200])
|
||
return "not_found"
|
||
except IMVAuthError as exc:
|
||
_mark_status(db, hid, "transient_error", f"auth: {exc!s}"[:200])
|
||
return "auth_error"
|
||
except IMVTransientError as exc:
|
||
_mark_status(db, hid, "transient_error", str(exc)[:200])
|
||
return "transient"
|
||
except Exception as exc:
|
||
_mark_status(db, hid, "error", repr(exc)[:200])
|
||
logger.error("house_imv: unexpected error house=%d: %r", hid, exc)
|
||
return "error"
|
||
|
||
# Адрес найден, но рыночной цены нет: avito_imv возвращает recommended_price=0
|
||
# ("нет оценки"). Не считаем это успехом — иначе дом залипает в imv_status='ok'
|
||
# и больше не переоценивается (backfill_house_imv фильтрует pending,
|
||
# process_houses_imv_batch — != 'ok', estimator — recommended_price>0).
|
||
# Помечаем not_found (адрес есть, но usable-оценки нет) — дом остаётся
|
||
# доступным для повторной обработки sweep'ом (imv_status != 'ok').
|
||
if eval_result.recommended_price <= 0:
|
||
logger.info(
|
||
"house_imv: empty IMV (no market price) house=%d market_count=%r",
|
||
hid,
|
||
eval_result.market_count,
|
||
)
|
||
_mark_status(
|
||
db,
|
||
hid,
|
||
"not_found",
|
||
f"imv_empty: recommended_price=0 market_count={eval_result.market_count}"[:200],
|
||
)
|
||
return "not_found"
|
||
|
||
try:
|
||
save_imv_result(db, hid, params, eval_result)
|
||
db.commit()
|
||
except Exception as exc:
|
||
logger.error("house_imv: save failed house=%d: %r", hid, exc)
|
||
try:
|
||
db.rollback()
|
||
except Exception:
|
||
pass
|
||
_mark_status(db, hid, "error", f"save_failed: {exc!s}"[:200])
|
||
return "error"
|
||
|
||
return "ok"
|