feat(tradein/estimate): внешние оценки не ждутся в запросе — 9 секунд превращаются в 1 #3055

Merged
lekss361 merged 1 commit from feat/estimate-external-sources-background into main 2026-08-22 12:58:04 +00:00
5 changed files with 344 additions and 15 deletions

View file

@ -723,6 +723,26 @@ class Settings(BaseSettings):
# ESTIMATE_HOUSE_META_TIMEOUT_S.
estimate_yandex_valuation_timeout_s: float = 8.0
estimate_cian_valuation_timeout_s: float = 8.0
# Внешние оценки (Yandex/Cian) не ждать в запросе, а догружать в фоне.
#
# Замер на проде 2026-08-22: расчёт по НОВОМУ адресу занимает 8-12 с, из них
# ~6 с ждёт Yandex и ~1.5 с Cian. Собственные запросы к базе и сам расчёт
# укладываются в секунду. По уже виденному адресу (кэш 24 ч) — 0.4-0.8 с.
#
# True: в запросе делается только чтение кэша; при промахе источник
# деградирует в None, а свежая загрузка уходит в фон и наполняет кэш к
# следующему обращению по тому же адресу. Ответ отдаётся за ~1 с.
#
# Деградация в None — НЕ новое состояние ответа: ровно так же ведёт себя
# таймаут `estimate_*_valuation_timeout_s`, этот путь работает в проде
# сегодня. Поэтому переключение не меняет контракт API.
#
# False (дефолт) — прежнее поведение: ждать источники в запросе.
# ENV: ESTIMATE_EXTERNAL_SOURCES_BACKGROUND.
estimate_external_sources_background: bool = Field(
default=False, validation_alias="ESTIMATE_EXTERNAL_SOURCES_BACKGROUND"
)
estimate_geocode_budget_s: float = 12.0
estimate_house_meta_timeout_s: float = 8.0

View file

@ -27,7 +27,7 @@ import math
import re
import statistics
import time
from collections.abc import Callable, Iterable
from collections.abc import Awaitable, Callable, Iterable
from dataclasses import dataclass
from datetime import UTC, date, datetime, timedelta
from typing import Any, Literal
@ -55,6 +55,7 @@ from sqlalchemy import text
from sqlalchemy.orm import Session
from app.core.config import Settings, settings
from app.core.db import SessionLocal
from app.schemas.trade_in import (
AggregatedEstimate,
AnalogLot,
@ -862,6 +863,63 @@ def _yandex_valuation_cache_key(address: str, offer_category: str, offer_type: s
return hashlib.sha256(payload.encode("utf-8")).hexdigest()
# ── отложенная догрузка внешних источников (фоновый режим) ───────────────────
#
# Держим ссылки на задачи: asyncio не хранит их сам, и без этого сборщик мусора
# может убить задачу на полпути (документированное поведение create_task).
_DEFERRED_REFRESH_TASKS: set[asyncio.Task[None]] = set()
# Подобран под текущий прод: 1 воркер uvicorn, mem_limit 768m, max_connections
# 100 у Postgres. Догрузка почти всё время ждёт чужой HTTP, поэтому небольшой
# потолок не занижает пропускную способность прогрева.
_MAX_DEFERRED_REFRESH_TASKS = 8
def _defer_external_refresh(label: str, work: Callable[[Session], Awaitable[object]]) -> None:
"""Догрузить внешний источник ПОСЛЕ ответа, наполнив кэш к следующему разу.
Своя сессия намеренно: сессия запроса закрывается вместе с ответом, а обе
функции источников делают внутри себя `db.commit()`. Переиспользование чужой
сессии зафиксировало бы её незавершённую работу.
Все ошибки гасятся: это прогрев кэша, а не часть ответа. Провал означает лишь
то, что следующий запрос по адресу снова промахнётся мимо кэша.
"""
async def _run() -> None:
db = SessionLocal()
try:
await work(db)
except Exception:
logger.exception("deferred %s: догрузка не удалась (кэш не прогрет)", label)
finally:
db.close()
# Потолок одновременных догрузок. Без него всплеск по НОВЫМ адресам (ровно
# тот случай, ради которого режим и сделан) породил бы сотни параллельных
# задач, каждая с сессией к базе и HTTP-клиентом — при max_connections=100
# и mem_limit 768m у backend это отказ вместо ускорения.
# Переполнение не ошибка: прогрев кэша необязателен, следующий запрос по
# адресу попробует снова.
if len(_DEFERRED_REFRESH_TASKS) >= _MAX_DEFERRED_REFRESH_TASKS:
logger.info(
"deferred %s: очередь догрузки заполнена (%d) — пропуск",
label,
_MAX_DEFERRED_REFRESH_TASKS,
)
return
try:
task = asyncio.create_task(_run())
except RuntimeError:
# Нет активного loop (синхронный вызов из теста/скрипта) — молча пропускаем:
# прогрев кэша необязателен по определению.
logger.debug("deferred %s: нет event loop — пропуск", label)
return
_DEFERRED_REFRESH_TASKS.add(task)
task.add_done_callback(_DEFERRED_REFRESH_TASKS.discard)
async def _get_or_fetch_yandex_valuation_cached(
db: Session,
*,
@ -869,6 +927,7 @@ async def _get_or_fetch_yandex_valuation_cached(
offer_category: str = YANDEX_VALUATION_DEFAULT_CATEGORY,
offer_type: str = YANDEX_VALUATION_DEFAULT_TYPE,
house_id: int | None = None,
fetch_on_miss: bool = True,
) -> YandexValuationResult | None:
"""Cached Yandex Valuation lookup. TTL 24h via external_valuations table.
@ -922,6 +981,12 @@ async def _get_or_fetch_yandex_valuation_cached(
except Exception as e:
logger.warning("yandex_valuation: cache deserialize failed — refetching: %s", e)
# fetch_on_miss=False — режим «только кэш» (фоновый режим внешних источников): вызывающий код
# не хочет ждать внешний HTTP в запросе и сам поставит догрузку в фон.
if not fetch_on_miss:
logger.info("yandex_valuation: cache MISS key=%s — fetch отложен", cache_key[:8])
return None
# Fresh fetch
try:
async with YandexValuationScraper(
@ -4285,13 +4350,24 @@ async def estimate_quality(
# ── Stage 8: Yandex Valuation as on-demand source (anonymous, cached 24h) ──
yandex_val: YandexValuationResult | None = None
if geo is not None and geo.full_address:
_y_addr, _y_house = geo.full_address, target_house_id
_bg = settings.estimate_external_sources_background
yandex_val = await _with_budget(
_get_or_fetch_yandex_valuation_cached(
db, address=geo.full_address, house_id=target_house_id
db, address=_y_addr, house_id=_y_house, fetch_on_miss=not _bg
),
settings.estimate_yandex_valuation_timeout_s,
label="yandex_valuation",
)
if _bg and yandex_val is None:
# Промах кэша: ответ не ждёт ~6 с внешнего HTTP, догрузка уходит в фон
# и наполняет кэш к следующему обращению по этому адресу.
_defer_external_refresh(
"yandex_valuation",
lambda bg_db: _get_or_fetch_yandex_valuation_cached(
bg_db, address=_y_addr, house_id=_y_house
),
)
if yandex_val is not None:
saved_hist = await asyncio.to_thread(_save_yandex_history_items, db, yandex_val)
logger.info(
@ -4311,24 +4387,30 @@ async def estimate_quality(
and payload.floor is not None
and payload.total_floors is not None
):
_c_bg = settings.estimate_external_sources_background
_c_kwargs: dict[str, Any] = {
"config": RealScraperConfig(),
"address": geo.full_address,
"total_area": payload.area_m2,
"rooms_count": payload.rooms,
"floor": payload.floor,
"total_floors": payload.total_floors,
"repair_type": "cosmetic",
"deal_type": "sale",
"use_cache": True,
"house_id": target_house_id,
}
try:
cian_val = await _with_budget(
estimate_via_cian_valuation(
db,
config=RealScraperConfig(),
address=geo.full_address,
total_area=payload.area_m2,
rooms_count=payload.rooms,
floor=payload.floor,
total_floors=payload.total_floors,
repair_type="cosmetic",
deal_type="sale",
use_cache=True,
house_id=target_house_id,
),
estimate_via_cian_valuation(db, fetch_on_miss=not _c_bg, **_c_kwargs),
settings.estimate_cian_valuation_timeout_s,
label="cian_valuation",
)
if _c_bg and cian_val is None:
_defer_external_refresh(
"cian_valuation",
lambda bg_db: estimate_via_cian_valuation(bg_db, **_c_kwargs),
)
if cian_val is not None and cian_val.sale_price_rub:
logger.info(
"cian_valuation: price=%s accuracy=%s house_id=%s",

View file

@ -0,0 +1,212 @@
"""Внешние источники оценки не ждутся в запросе, а догружаются в фоне.
Замер на проде 2026-08-22: расчёт по новому адресу 8-12 с, из них ~6 с ждёт
Yandex и ~1.5 с Cian; собственный расчёт укладывается в секунду. По уже
виденному адресу (кэш 24 ч) 0.4-0.8 с.
Ключевое свойство, которое делает переключение безопасным: деградация источника
в None НЕ новое состояние ответа. Ровно так же ведёт себя таймаут
`estimate_*_valuation_timeout_s`, и этот путь работает в проде сегодня.
Тесты проверяют три вещи, каждая из которых при поломке молча вернула бы
секунды ожидания в пользовательский путь:
1. `fetch_on_miss=False` действительно НЕ ходит в сеть при промахе кэша;
2. попадание в кэш работает одинаково в обоих режимах;
3. отложенная задача берёт СВОЮ сессию и гасит свои ошибки.
"""
from __future__ import annotations
import asyncio
from typing import Any
from unittest.mock import MagicMock
import pytest
class _NoRowsDB:
"""Сессия, у которой кэш всегда пуст."""
def __init__(self) -> None:
self.commits = 0
self.closed = False
def execute(self, *_a: Any, **_kw: Any) -> Any:
result = MagicMock()
result.mappings.return_value.first.return_value = None
return result
def commit(self) -> None:
self.commits += 1
def rollback(self) -> None:
return None
def close(self) -> None:
self.closed = True
# ── Yandex ───────────────────────────────────────────────────────────────────
def test_yandex_cache_only_does_not_touch_network(monkeypatch: pytest.MonkeyPatch) -> None:
"""fetch_on_miss=False → None и НИ ОДНОГО обращения к скраперу."""
from app.services import estimator
def _boom(*_a: Any, **_kw: Any) -> Any:
raise AssertionError("сеть не должна дёргаться в режиме «только кэш»")
monkeypatch.setattr(estimator, "YandexValuationScraper", _boom)
out = asyncio.run(
estimator._get_or_fetch_yandex_valuation_cached(
_NoRowsDB(), # type: ignore[arg-type]
address="Екатеринбург, улица Сурикова, 4",
fetch_on_miss=False,
)
)
assert out is None
def test_yandex_default_still_fetches(monkeypatch: pytest.MonkeyPatch) -> None:
"""Дефолт не изменился: без флага источник по-прежнему грузится в запросе."""
from app.services import estimator
called: dict[str, bool] = {}
class _Scraper:
def __init__(self, *_a: Any, **_kw: Any) -> None:
called["constructed"] = True
async def __aenter__(self) -> _Scraper:
return self
async def __aexit__(self, *_a: Any) -> None:
return None
async def fetch_house_history(self, **_kw: Any) -> None:
return None
monkeypatch.setattr(estimator, "YandexValuationScraper", _Scraper)
asyncio.run(
estimator._get_or_fetch_yandex_valuation_cached(
_NoRowsDB(), # type: ignore[arg-type]
address="Екатеринбург, улица Сурикова, 4",
)
)
assert called.get("constructed") is True
# ── отложенная догрузка ──────────────────────────────────────────────────────
def test_deferred_refresh_uses_its_own_session(monkeypatch: pytest.MonkeyPatch) -> None:
"""Сессия запроса закрывается вместе с ответом — фон обязан открыть свою.
Переиспользование чужой сессии зафиксировало бы её незавершённую работу:
обе функции источников делают внутри себя `db.commit()`.
"""
from app.services import estimator
own = _NoRowsDB()
request_db = _NoRowsDB()
monkeypatch.setattr(estimator, "SessionLocal", lambda: own)
seen: dict[str, Any] = {}
async def _work(db: Any) -> None:
seen["db"] = db
async def _drive() -> None:
estimator._defer_external_refresh("test", _work)
await asyncio.sleep(0)
await asyncio.sleep(0)
asyncio.run(_drive())
assert seen["db"] is own, "фоновая задача взяла не свою сессию"
assert seen["db"] is not request_db
assert own.closed is True, "фоновая сессия осталась незакрытой"
def test_deferred_refresh_swallows_errors(monkeypatch: pytest.MonkeyPatch) -> None:
"""Провал прогрева кэша не должен всплывать: это не часть ответа."""
from app.services import estimator
own = _NoRowsDB()
monkeypatch.setattr(estimator, "SessionLocal", lambda: own)
async def _work(_db: Any) -> None:
raise RuntimeError("внешний сервис лёг")
async def _drive() -> None:
estimator._defer_external_refresh("test", _work)
await asyncio.sleep(0)
await asyncio.sleep(0)
asyncio.run(_drive()) # не должно поднять
assert own.closed is True
def test_deferred_refresh_without_event_loop_is_noop() -> None:
"""Синхронный контекст (скрипт/тест) — прогрев необязателен, не падаем."""
from app.services import estimator
async def _work(_db: Any) -> None:
raise AssertionError("не должно вызваться")
estimator._defer_external_refresh("test", _work) # без запущенного loop
def test_task_reference_is_retained(monkeypatch: pytest.MonkeyPatch) -> None:
"""Ссылка на задачу удерживается: иначе GC может убить её на полпути."""
from app.services import estimator
monkeypatch.setattr(estimator, "SessionLocal", _NoRowsDB)
started = asyncio.Event()
release = asyncio.Event()
async def _work(_db: Any) -> None:
started.set()
await release.wait()
async def _drive() -> None:
estimator._defer_external_refresh("test", _work)
await started.wait()
assert len(estimator._DEFERRED_REFRESH_TASKS) == 1
release.set()
await asyncio.sleep(0)
await asyncio.sleep(0)
assert len(estimator._DEFERRED_REFRESH_TASKS) == 0, "ссылка не убрана после завершения"
asyncio.run(_drive())
def test_deferred_queue_is_bounded(monkeypatch: pytest.MonkeyPatch) -> None:
"""Всплеск по новым адресам не должен породить сотни задач.
Ровно тот случай, ради которого режим и сделан: у backend один воркер,
mem_limit 768m, у Postgres max_connections 100. Неограниченный create_task
превратил бы ускорение в отказ.
"""
from app.services import estimator
monkeypatch.setattr(estimator, "SessionLocal", _NoRowsDB)
monkeypatch.setattr(estimator, "_MAX_DEFERRED_REFRESH_TASKS", 3)
release = asyncio.Event()
async def _work(_db: Any) -> None:
await release.wait()
async def _drive() -> None:
for _ in range(50):
estimator._defer_external_refresh("test", _work)
await asyncio.sleep(0)
assert len(estimator._DEFERRED_REFRESH_TASKS) <= 3, (
f"очередь догрузки не ограничена: {len(estimator._DEFERRED_REFRESH_TASKS)} задач"
)
release.set()
for _ in range(5):
await asyncio.sleep(0)
asyncio.run(_drive())

View file

@ -246,6 +246,13 @@ services:
- path: ./backend/.env.runtime
required: false
environment:
# Внешние оценки (Yandex/Cian) не ждать в запросе — догружать в фоне.
# Замер 2026-08-22: расчёт по новому адресу 8-12 с, из них ~6 с Yandex и
# ~1.5 с Cian, собственный расчёт < 1 с. Публикация в РБК 30.08 приведёт
# аудиторию на НОВЫЕ адреса, то есть мимо суточного кэша.
# Деградация источника в None — существующее состояние ответа (так же
# ведёт себя таймаут), контракт API не меняется.
ESTIMATE_EXTERNAL_SOURCES_BACKGROUND: "true"
DATABASE_URL: "postgresql+psycopg://${TRADEIN_POSTGRES_USER:-tradein}:${TRADEIN_POSTGRES_PASSWORD}@postgres:5432/tradein"
PUBLIC_URL: "https://gendsgn.ru/trade-in"
CORS_ORIGINS: '["https://gendsgn.ru"]'

View file

@ -143,6 +143,7 @@ async def estimate_via_cian_valuation(
repair_type: str = "cosmetic",
deal_type: str = "sale",
use_cache: bool = True,
fetch_on_miss: bool = True,
house_id: int | None = None,
listing_id: int | None = None,
proxy_provider: ProxyProvider | None = None,
@ -164,6 +165,13 @@ async def estimate_via_cian_valuation(
logger.info("Cian valuation cache HIT for cache_key=%s...", cache_key[:12])
return cached
# fetch_on_miss=False — режим «только кэш»: вызывающий код не хочет ждать
# внешний HTTP в запросе и сам поставит догрузку в фон. Возврат None здесь
# идёт по тому же graceful-пути, что таймаут и сетевая ошибка.
if not fetch_on_miss:
logger.info("Cian valuation: cache MISS key=%s... — fetch отложен", cache_key[:12])
return None
# 2. Load auth cookies из cian_session_cookies
cookies = load_session(db, config=config)
if cookies is None: