gendesign/tradein-mvp/backend/app/services/house_imv_backfill.py
bot-backend 43aaf91b97
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
fix(tradein): писатели наконец пишут то, что обещает схема — фото подсказок, статус «снято», события объявлений (#2674)
Три находки одного класса из эпика: колонка есть, писатель есть, тест на писателя
зелёный — а данные не появляются. Тестами это не ловится по построению, только
сверкой с продом.

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.
2026-08-06 02:29:58 +05:00

802 lines
36 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

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