feat(tradein/estimate): внешние оценки не ждутся в запросе — 9 секунд превращаются в 1 #3055
5 changed files with 344 additions and 15 deletions
|
|
@ -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
|
||||
|
||||
|
|
|
|||
|
|
@ -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",
|
||||
|
|
|
|||
|
|
@ -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())
|
||||
|
|
@ -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"]'
|
||||
|
|
|
|||
|
|
@ -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:
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue