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
981 lines
46 KiB
Python
981 lines
46 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,
|
||
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"
|