gendesign/tradein-mvp/backend/app/services/house_imv_backfill.py
bot-backend c927b77777
Some checks failed
Deploy Trade-In / changes (push) Successful in 12s
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / test (push) Failing after 3m29s
Deploy Trade-In / build-backend (push) Has been skipped
Deploy Trade-In / deploy (push) Has been skipped
fix(tradein/imv): «временная» ошибка снова временная — 1390 домов возвращаются в очередь (#2843)
2026-08-12 16:06:24 +00:00

981 lines
46 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,
imv_transient_attempts = 0
WHERE id = :hid
"""),
{"hid": house_id},
)
def _mark_status(
db: Session,
house_id: int,
status: str,
reason: str | None = None,
) -> None:
# #2674: счётчик растёт ТОЛЬКО на transient_error — это «сколько раз подряд дом
# падал по временной причине», а не «сколько раз его трогали». no_params /
# no_address / not_found счётчик не двигают: они не занимают retry-слот.
db.execute(
text("""
UPDATE houses
SET imv_status = :s,
last_imv_attempt_at = NOW(),
imv_error_reason = :r,
imv_transient_attempts = CASE
WHEN :s = 'transient_error' THEN imv_transient_attempts + 1
ELSE imv_transient_attempts
END
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"
]
# #2674: сколько раз подряд дом может упасть в transient_error, прежде чем
# перестанет занимать retry-слот. Число из замера: после починки сайдкара (04.08)
# доля отказов на попытку — 2/27 и 3/25 (прогоны 3708/3467), т.е. ~10%. На 1039
# застрявших это ~104 повторных отказа на первом проходе, ~10 на втором, ~1 на
# третьем. Порог 3 стоит максимум ~115 слотов ВСЕГО (≈2 прогона) и гарантирует,
# что дом со СВОЕЙ (не инфраструктурной) причиной не крутится в пакете вечно.
# Исчерпавшие лимит не исчезают из наблюдаемости: они остаются imv_status=
# 'transient_error' и считаются как
# WHERE imv_status='transient_error' AND imv_transient_attempts >= 3.
_MAX_TRANSIENT_ATTEMPTS = 3
# Доля пакета под повтор transient_error. Половина — потому что остальные слоты
# после #2674 достаются ТОЛЬКО домам, по которым реально будет запрос к площадке
# (см. _premark_unusable): раньше из 50 слотов до площадки доходили 17 (замер
# головы очереди на 12.08), так что pending на половине пакета всё равно идёт
# быстрее, чем на целом до правки.
_RETRY_SLOTS_SHARE = 0.5
# Дом без пригодных параметров backfill всё равно пометит no_params — но только
# заплатив слотом пакета и паузой request_delay_sec. Тот же вердикт берётся одним
# запросом: нет ни одного объявления с rooms+area (pick_lot_params вернёт {}) ИЛИ
# не из чего взять house_type (_map_house_type вернёт None → «unknown house_type»).
# Причины пишем ТЕМИ ЖЕ строками, что и поштучный путь, — старые разрезы по
# imv_error_reason продолжают работать.
# Условие сознательно УЖЕ питоновского: normalize_house_type схлопывает в None ещё
# и нераспознанный вокабуляр ('other', 'wireframe'), который тут остаётся текстом.
# Промахнуться можно только в безопасную сторону — пометить меньше, чем пометил бы
# поштучный путь.
_PREMARK_UNUSABLE_SQL = text("""
WITH unusable AS (
SELECT h.id,
CASE WHEN NOT EXISTS (
SELECT 1 FROM listings l
WHERE l.house_id_fk = h.id
AND l.rooms IS NOT NULL
AND l.area_m2 IS NOT NULL)
THEN 'no listings with rooms+area'
ELSE 'unknown house_type'
END AS reason
FROM houses h
WHERE h.imv_status = ANY(CAST(:statuses AS text[]))
AND h.lat IS NOT NULL
AND h.lon IS NOT NULL
AND h.address IS NOT NULL
AND (
NOT EXISTS (
SELECT 1 FROM listings l
WHERE l.house_id_fk = h.id
AND l.rooms IS NOT NULL
AND l.area_m2 IS NOT NULL)
OR COALESCE(
NULLIF(TRIM((
SELECT mode() WITHIN GROUP (ORDER BY l.house_type)
FROM listings l
WHERE l.house_id_fk = h.id
AND l.rooms IS NOT NULL
AND l.area_m2 IS NOT NULL)), ''),
NULLIF(TRIM(h.house_type), '')
) IS NULL
)
)
UPDATE houses
SET imv_status = 'no_params',
last_imv_attempt_at = NOW(),
imv_error_reason = unusable.reason
FROM unusable
WHERE houses.id = unusable.id
""")
# Основная очередь: один статус, как и было (only_status — публичный параметр
# admin-API, семантику не трогаем).
_QUEUE_SQL = 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
""")
# Retry-очередь (#2674). Отдельный запрос, а не OR к основной: у pending
# last_imv_attempt_at всегда NULL, поэтому при общем ORDER BY ... NULLS FIRST
# transient_error не попал бы в пакет, пока не кончится pending (по замеру
# 12.08 — 5747 домов ≈ год). Отдельная квота = отдельный проход.
_RETRY_QUEUE_SQL = text("""
SELECT id, address, full_address, lat, lon
FROM houses
WHERE imv_status = 'transient_error'
AND imv_transient_attempts < :max_attempts
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
""")
@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)
# #2674: сколько домов пакета пришло из retry-очереди transient_error и
# сколько помечено no_params до пакета (без запроса к площадке).
retried: int = 0
premarked: int = 0
def _premark_unusable(db: Session, statuses: list[str]) -> int:
"""Пометить no_params дома, по которым запрос к площадке невозможен. → сколько.
Не новое поведение, а тот же вердикт _process_one_house одним запросом: на
12.08 в очереди 1925 таких домов из 5143 (113 без объявлений с rooms+area,
1812 без house_type) — каждый занимал слот пакета и паузу, чтобы получить
ответ, который виден в SQL.
"""
res = db.execute(_PREMARK_UNUSABLE_SQL, {"statuses": statuses})
db.commit()
return int(res.rowcount or 0)
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. При значении по умолчанию
часть пакета (_RETRY_SLOTS_SHARE) автоматически уходит на повтор
transient_error с непотраченным лимитом попыток (#2674) — явно
переданный only_status этот проход отключает, оператор получает
ровно то, что попросил, включая исчерпавшие лимит дома.
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:
# Повторный проход только на расписании (only_status по умолчанию): явный
# only_status от оператора — это ручной запрос ровно одного статуса.
retry_lane = only_status == "pending"
statuses = [only_status] + (["transient_error"] if retry_lane else [])
result.premarked = _premark_unusable(db, statuses)
if result.premarked:
logger.info(
"house_imv_backfill: %d домов помечены no_params до пакета (нет rooms+area "
"или house_type) — слоты пакета не потрачены",
result.premarked,
)
retry_rows: list = []
if retry_lane:
retry_rows = (
db.execute(
_RETRY_QUEUE_SQL,
{
"max_attempts": _MAX_TRANSIENT_ATTEMPTS,
"batch": int(batch_size * _RETRY_SLOTS_SHARE),
},
)
.mappings()
.all()
)
result.retried = len(retry_rows)
# Недобор retry-очереди (она кончится раньше pending: 1039 против 3218 на
# 12.08) возвращается pending — пакет не простаивает.
fresh_rows = (
db.execute(
_QUEUE_SQL,
{"status": only_status, "batch": max(batch_size - result.retried, 0)},
)
.mappings()
.all()
)
rows = list(fresh_rows) + list(retry_rows)
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 retry=%d premarked=%d delay=%.1fs)",
result.checked,
only_status,
result.retried,
result.premarked,
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-сессию как раньше
# (поведение байт-в-байт идентично доспринтовому).
#
# #2698: proxy_provider/use_pool/environment — обязательная часть проводки, а не
# опция. Без них BrowserFetcher не кладёт "proxy" в тело POST /fetch-json, и сайдкар
# берёт свой env-прокси SCRAPER_PROXY_URL — на проде это узел пула id=1
# (asocks-residential-1, provider_affinity='domclick'), который proxy_pool.acquire
# («affinity IN (provider,'any')» + защита последнего узла выделенной affinity от
# fallback) для avito не выдал бы НИКОГДА. Результат: 03.07-05.08 все 35 из 35 попыток
# каждого прогона падали на геокодере A (1240 домов — 503 «browser unavailable», затем
# 500 «Page.goto: NS_ERROR_PROXY_BAD_GATEWAY» и 403 от самого Авито), пока
# avito_city_sweep/avito_newbuilding_sweep в те же дни тянули сотни объявлений через
# ТОТ ЖЕ сайдкар и тот же инстанс камуфокса — они пул подключают (pipeline.py). Хуже:
# запрос без "proxy" в теле ещё и роняет сайдкару желаемый прокси на env → relaunch
# камуфокса на каждый дом (server.py::_ensure_browser).
if settings.avito_imv_use_browser_fetcher:
# lazy import — тот же цикл scraper_adapters↔этот модуль, что и у RealScraperConfig.
from app.services.scraper_adapters import RealProxyProvider, RealScraperConfig
_cfg = RealScraperConfig()
async with BrowserFetcher(
source="avito",
endpoint=settings.browser_http_endpoint,
proxy_provider=RealProxyProvider(),
use_pool=_cfg.use_proxy_pool_browser,
# #2616 шаг 1: без environment прод-отказ «пул пуст» мёртв на этом пути —
# фетчер молча ушёл бы на тот самый env-прокси (см. _acquire_lease).
environment=_cfg.environment,
) 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 "
"retried=%d premarked=%d %.1fs %s",
result.checked,
result.saved,
result.skipped,
result.errors,
result.retried,
result.premarked,
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"