feat(tradein/cian): у Циана не было добора карточек — только побочный эффект задачи про историю (#3284) #3285

Merged
lekss361 merged 1 commit from feat/3284-cian-detail-backfill into main 2026-08-30 12:41:37 +00:00
5 changed files with 374 additions and 18 deletions

View file

@ -706,6 +706,15 @@ def build_product_handlers(ctx: SchedulerContext) -> dict[str, Handler]:
"cian_history_backfill",
pre_claim=_cian_pre_claim,
),
# #3284: то же тело, другая очередь. Разные source нужны именно как РАЗНЫЕ
# строки расписания — у них свои окна, свой next_run_at и свой счётчик
# прогонов; гонять оба режима под одним source нельзя, планировщик держит
# на source ровно один активный прогон.
"cian_detail_backfill": Handler(
_job_cian_history_backfill,
"cian_detail_backfill",
pre_claim=_cian_pre_claim,
),
"rosreestr_dkp_import": Handler(_job_rosreestr_dkp, "rosreestr_dkp_import"),
"listing_source_snapshot": Handler(_job_listing_source_snapshot, "listing_source_snapshot"),
"asking_to_sold_ratio_refresh": Handler(

View file

@ -99,10 +99,18 @@ async def _execute_cian_backfill(
Params (from default_params jsonb):
batch_size: int rows per run (listings + houses counted separately).
listings_pending: str "history" (дефолт) | "detail", см. #3284.
do_houses: bool дефолт true; у cian_detail_backfill выключен.
"""
from app.tasks.cian_history_backfill import CianBackfillResult, backfill_cian_history
batch_size = int(params.get("batch_size", 100))
# #3284: одно тело обслуживает ДВА расписания. cian_history_backfill идёт с
# дефолтами (история + дома), cian_detail_backfill — с listings_pending="detail"
# и do_houses=false: дома у него уже разбирает суточный сосед, а гонять их
# круглосуточно незачем.
listings_pending = str(params.get("listings_pending", "history"))
do_houses = bool(params.get("do_houses", True))
def _counters(result: CianBackfillResult) -> dict[str, int]:
return {
@ -143,9 +151,10 @@ async def _execute_cian_backfill(
db,
batch_size=batch_size,
do_listings=True,
do_houses=True,
do_houses=do_houses,
do_valuations=False,
on_progress=_heartbeat,
listings_pending=listings_pending,
)
counters = {**_counters(result), "duration_sec": int(result.duration_sec)}

View file

@ -110,6 +110,41 @@ def _note_refusal(result: CianBackfillResult, status: int | None) -> str | None:
return kind
# ── Выборки «что ещё не добрано» ─────────────────────────────────────────────
# Две РАЗНЫЕ цели, которые до #3284 были склеены в одну. Историческая выборка
# ключуется по offer_price_history, поэтому объявление, у которого история уже
# есть, а карточки нет, не вернётся ей НИКОГДА (на 30.08 таких 1697). Вторая
# выборка закрывает ровно этот пробел и заодно даёт Циану то, что у avito /
# domclick / yandex есть давно, — добор по признаку «нет карточки».
_LISTINGS_PENDING_SQL: dict[str, str] = {
# Прежнее поведение, побайтово. Менять его правкой про добор карточек нельзя:
# по этой выборке живёт суточный cian_history_backfill.
"history": """
SELECT l.id, l.source_url
FROM listings l
LEFT JOIN offer_price_history oph ON oph.listing_id = l.id
WHERE l.source = 'cian'
AND l.source_url IS NOT NULL
AND oph.listing_id IS NULL
LIMIT :lim
""",
# #3284. ORDER BY здесь есть, а в "history" нет, и это намеренно: очередь
# карточек (на 30.08 — 20 728 объявлений) заведомо длиннее любого батча, и
# порядок решает, что мы успеем добрать. Свежие важнее: по ним считается
# оценка. У "history" очередь того же порядка, но её сортировку трогать —
# отдельное решение с отдельной проверкой, не побочный эффект этой правки.
"detail": """
SELECT l.id, l.source_url
FROM listings l
WHERE l.source = 'cian'
AND l.source_url IS NOT NULL
AND l.detail_enriched_at IS NULL
ORDER BY l.last_seen_at DESC NULLS LAST
LIMIT :lim
""",
}
async def backfill_cian_history(
db: Session,
*,
@ -119,6 +154,7 @@ async def backfill_cian_history(
do_valuations: bool = False,
dry_run: bool = False,
on_progress: Callable[[CianBackfillResult], None] | None = None,
listings_pending: str = "history",
) -> CianBackfillResult:
"""Iterate Cian listings + houses with missing history, fetch+save.
@ -133,6 +169,11 @@ async def backfill_cian_history(
do_valuations: process Cian Valuation Calculator batch (external_valuations backfill).
Default False opt-in because each call hits Cian auth-gated API.
dry_run: skip all fetch+save; only count and log pending rows.
listings_pending: какую очередь разбирает блок listings (#3284) --
"history" (дефолт, прежнее поведение): нет строки в offer_price_history;
"detail": нет карточки (detail_enriched_at IS NULL), свежие первыми.
Неизвестное значение -- ValueError, а не молчаливый дефолт: пустой
батч из-за опечатки в default_params выглядел бы как «всё добрано».
on_progress: колбэк живости (#2725) — вызывается на каждой сущности ЛЮБОГО из
трёх блоков, до её обработки, с текущим (мутируемым) result. Caller пишет
heartbeat; исключения колбэка на его совести (планировщик глушит их сам),
@ -145,24 +186,16 @@ async def backfill_cian_history(
t0 = time.time()
delay = get_scraper_delay("cian") # seconds; default 5.0
# ── 1. Listings: missing offer_price_history ──────────────────────────────
# ── 1. Listings: pending-выборка (см. _LISTINGS_PENDING_SQL) ──────────────
if do_listings:
rows = (
db.execute(
text("""
SELECT l.id, l.source_url
FROM listings l
LEFT JOIN offer_price_history oph ON oph.listing_id = l.id
WHERE l.source = 'cian'
AND l.source_url IS NOT NULL
AND oph.listing_id IS NULL
LIMIT :lim
"""),
{"lim": batch_size},
)
.mappings()
.all()
)
try:
pending_sql = _LISTINGS_PENDING_SQL[listings_pending]
except KeyError:
raise ValueError(
f"listings_pending={listings_pending!r} неизвестен; "
f"допустимы {sorted(_LISTINGS_PENDING_SQL)}"
) from None
rows = db.execute(text(pending_sql), {"lim": batch_size}).mappings().all()
result.listings_total = len(rows)
if dry_run:

View file

@ -0,0 +1,70 @@
-- 282_scrape_schedules_seed_cian_detail_backfill.sql
-- Отдельный добор карточек Циана (issue #3284).
--
-- ЧТО БЫЛО НЕ ТАК
-- У avito, domclick и yandex есть выделенный *_detail_backfill. У Циана его не было
-- вообще: карточки доставались побочным эффектом cian_history_backfill — суточной
-- задачи про историю цен, — и её выборка ключуется по offer_price_history, а не по
-- наличию карточки. Следствия на 30.08.2026:
-- * 25 222 объявления Циана, карточка есть у 4 494 (17.8%) — последнее место среди
-- четырёх источников при втором месте по объёму (domklik 76%, yandex 50%,
-- avito 24%);
-- * 19 046 объявлений без истории цен при квоте batch_size=100 в сутки — это 190
-- дней на текущий остаток, притом что Циан приносит ~4 900 объявлений за двое
-- суток, то есть очередь растёт примерно в двадцать раз быстрее, чем разбирается;
-- * 1 697 объявлений имеют историю цен и НЕ имеют карточки — исторической выборке
-- они уже «обработаны» и не вернутся к ней никогда.
-- При этом сам фетчер исправен: прогоны 5154/5240/5328 дали 100/100, 99/100, 100/100.
-- Чинить нечего — не выдана мощность.
--
-- ЧТО ДЕЛАЕТ ЭТА СТРОКА
-- Второе расписание поверх ТОГО ЖЕ тела (app.services.scheduler._execute_cian_backfill),
-- отличающееся двумя параметрами:
-- listings_pending = "detail" — выборка по detail_enriched_at IS NULL, свежие первыми
-- (ORDER BY last_seen_at DESC): очередь заведомо длиннее батча, и
-- порядок решает, что успеем добрать; по свежим считается оценка.
-- do_houses = false — дома разбирает суточный cian_history_backfill, дублировать их
-- круглосуточно незачем.
--
-- ОКНО 023
-- Ровно как у avito_detail_backfill и yandex_detail_backfill. Трёхчасовое окно, как у
-- истории, здесь бессмысленно: прогон 5328 уложился в 54 минуты из трёх часов, то есть
-- две трети окна простаивали, а очередь при этом росла.
--
-- BATCH_SIZE 400 — арифметика, а не вкус
-- Замеренный темп ~28 с на объявление (задержка 5 с, остальное сеть). 400 × 28 с ≈ 3.1 ч.
-- Порог reap_zombies — 6 ч по heartbeat_at, и прогон обязан заканчиваться заметно раньше:
-- у backfill_cian_history НЕТ бюджетного сторожа, единственный ограничитель — batch_size.
-- Отсюда же не 800, как у yandex: 800 × 28 с ≈ 6.2 ч — прогон убивало бы зомби-жнецом.
--
-- ENABLED=false
-- Как и у domclick_detail_backfill (миграция 175): включение — отдельный осознанный шаг
-- после деплоя и дымовой пробы. Причина не в куках (гейт _cian_pre_claim общий с историей
-- и проходит), а в нагрузке: это ТРЕТИЙ круглосуточный добор на общий пул из четырёх
-- узлов, и его влияние на соседей надо посмотреть, а не предположить.
--
-- ЗАВИСИМОСТИ: 052_scrape_schedules.sql (таблица + UNIQUE(source)).
-- Идемпотентно: ON CONFLICT (source) DO NOTHING.
BEGIN;
INSERT INTO scrape_schedules (
source,
enabled,
window_start_hour,
window_end_hour,
next_run_at,
default_params
)
VALUES
(
'cian_detail_backfill',
false,
0,
23,
((CURRENT_DATE + INTERVAL '1 day')) AT TIME ZONE 'UTC',
'{"batch_size": 400, "listings_pending": "detail", "do_houses": false}'::jsonb
)
ON CONFLICT (source) DO NOTHING;
COMMIT;

View file

@ -0,0 +1,235 @@
"""#3284: у Циана не было добора карточек — только побочный эффект задачи про историю цен.
Историческая выборка ключуется по offer_price_history. Объявление, у которого история
уже есть, а карточки нет, для неё «обработано» и не вернётся никогда на 30.08.2026
таких 1697. Плюс квота 100/сутки против очереди в 19 046 это 190 дней, при том что
Циан приносит ~4900 объявлений за двое суток.
Правка добавляет второй режим выборки (`listings_pending="detail"`) и второе расписание
поверх того же тела. Тесты ниже закрепляют ровно то, что делает режим полезным:
выборку по detail_enriched_at, свежие первыми, неизменность прежнего режима и отказ
на опечатке в default_params (иначе пустой батч читался бы как «всё добрано»).
"""
from __future__ import annotations
import os
import re
from types import SimpleNamespace
from typing import Any
from unittest.mock import AsyncMock, MagicMock, patch
import pytest
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db")
from app.tasks import cian_history_backfill
from app.tasks.cian_history_backfill import _LISTINGS_PENDING_SQL
def _norm(sql: str) -> str:
"""Схлопнуть пробелы — сравниваем смысл запроса, а не его отступы."""
return re.sub(r"\s+", " ", sql).strip()
class _FakeBrowserFetcher:
def __init__(self, **kwargs: Any) -> None:
self.last_response_status = None
async def __aenter__(self) -> _FakeBrowserFetcher:
return self
async def __aexit__(self, *_: object) -> None:
return None
def _enrichment() -> SimpleNamespace:
return SimpleNamespace(price_changes=[])
def _db_returning(rows: list[dict[str, Any]]) -> MagicMock:
db = MagicMock()
db.execute.return_value.mappings.return_value.all.return_value = rows
return db
async def _run(db: MagicMock, **kwargs: Any) -> Any:
with (
patch.object(cian_history_backfill, "BrowserFetcher", _FakeBrowserFetcher),
patch.object(cian_history_backfill, "fetch_detail", AsyncMock(return_value=_enrichment())),
patch.object(cian_history_backfill, "save_detail_enrichment", MagicMock()),
patch.object(cian_history_backfill, "RealMatcherAdapter", MagicMock()),
patch("asyncio.sleep", new_callable=AsyncMock),
):
return await cian_history_backfill.backfill_cian_history(
db, do_houses=False, do_valuations=False, **kwargs
)
# ── выборка "detail" ─────────────────────────────────────────────────────────
def test_detail_queue_selects_by_missing_card_not_by_history() -> None:
"""Ключ выборки — detail_enriched_at, а не offer_price_history.
Это вся суть #3284: пока признаком служит история, 1697 объявлений с историей
и без карточки недостижимы.
"""
sql = _norm(_LISTINGS_PENDING_SQL["detail"])
assert "l.detail_enriched_at IS NULL" in sql
assert "offer_price_history" not in sql
def test_detail_queue_takes_freshest_first() -> None:
"""Очередь длиннее батча, поэтому порядок решает, что мы успеем добрать."""
sql = _norm(_LISTINGS_PENDING_SQL["detail"])
assert "ORDER BY l.last_seen_at DESC NULLS LAST" in sql
assert sql.index("ORDER BY") < sql.index("LIMIT")
def test_detail_queue_stays_scoped_to_cian() -> None:
"""Источник обязан быть в WHERE: иначе добор Циана заберёт чужие объявления."""
assert "l.source = 'cian'" in _norm(_LISTINGS_PENDING_SQL["detail"])
assert "l.source_url IS NOT NULL" in _norm(_LISTINGS_PENDING_SQL["detail"])
# ── прежний режим не тронут ──────────────────────────────────────────────────
def test_history_queue_unchanged() -> None:
"""По этой выборке живёт суточный cian_history_backfill — она обязана остаться прежней."""
sql = _norm(_LISTINGS_PENDING_SQL["history"])
assert "LEFT JOIN offer_price_history oph ON oph.listing_id = l.id" in sql
assert "oph.listing_id IS NULL" in sql
assert "detail_enriched_at" not in sql
# Сортировки в исторической выборке не было и не появилось: её добавление —
# отдельное решение с отдельной проверкой, а не побочный эффект #3284.
assert "ORDER BY" not in sql
def test_default_mode_is_history() -> None:
"""Дефолт обязан оставаться прежним: расписание истории параметра не передаёт."""
db = _db_returning([])
import asyncio
asyncio.run(_run(db))
used = _norm(str(db.execute.call_args_list[0].args[0]))
assert "offer_price_history" in used
# ── режим доезжает до запроса ────────────────────────────────────────────────
async def test_detail_mode_reaches_the_query() -> None:
"""Параметр не должен потеряться по дороге — проверяем сам исполненный SQL."""
db = _db_returning([])
await _run(db, listings_pending="detail")
used = _norm(str(db.execute.call_args_list[0].args[0]))
assert "detail_enriched_at IS NULL" in used
assert "offer_price_history" not in used
async def test_detail_mode_processes_rows_normally() -> None:
"""Смена выборки не меняет обработку: строки те же, счётчики те же."""
db = _db_returning([{"id": 7, "source_url": "https://cian.ru/7"}])
result = await _run(db, listings_pending="detail", batch_size=10)
assert result.listings_processed == 1
assert result.listings_succeeded == 1
async def test_batch_size_reaches_the_query() -> None:
"""batch_size обязан доезжать как :lim — иначе квота из расписания ничего не значит."""
db = _db_returning([])
await _run(db, listings_pending="detail", batch_size=400)
assert db.execute.call_args_list[0].args[1] == {"lim": 400}
# ── опечатка не должна читаться как «всё добрано» ────────────────────────────
async def test_unknown_mode_raises_instead_of_silently_defaulting() -> None:
"""Опечатка в default_params обязана падать громко.
Молчаливый откат на "history" дал бы прогон с нулём добранных карточек, который
выглядит как штатный: очередь якобы пуста. Такую ошибку ищут днями.
"""
db = _db_returning([])
with pytest.raises(ValueError, match="listings_pending"):
await _run(db, listings_pending="detali")
async def test_unknown_mode_names_the_allowed_values() -> None:
"""Сообщение должно называть допустимые значения — иначе оно не помогает."""
db = _db_returning([])
with pytest.raises(ValueError) as e:
await _run(db, listings_pending="")
assert "detail" in str(e.value) and "history" in str(e.value)
# ── параметры доезжают от расписания до задачи ───────────────────────────────
# Три предыдущих теста проверяют саму задачу. Эти — путь от строки расписания:
# default_params → _execute_cian_backfill → backfill_cian_history. Без них правка
# в теле-посреднике молча вернула бы cian_detail_backfill к разбору истории.
async def test_scheduler_passes_detail_mode_from_params() -> None:
"""listings_pending из default_params обязан доехать до задачи."""
from app.services import scheduler as scheduler_mod
captured: dict[str, Any] = {}
async def _fake_backfill(db: Any, **kwargs: Any) -> Any:
captured.update(kwargs)
# Настоящий результат, а не SimpleNamespace: посредник читает у него
# поля, которых в самодельной заглушке легко недосчитаться.
return cian_history_backfill.CianBackfillResult()
with (
patch("app.tasks.cian_history_backfill.backfill_cian_history", _fake_backfill),
patch.object(scheduler_mod.runs_mod, "update_heartbeat", MagicMock()),
patch.object(scheduler_mod.runs_mod, "mark_done", MagicMock()),
):
await scheduler_mod._execute_cian_backfill(
MagicMock(),
run_id=1,
params={"batch_size": 400, "listings_pending": "detail", "do_houses": False},
)
assert captured["listings_pending"] == "detail"
assert captured["do_houses"] is False
assert captured["batch_size"] == 400
async def test_scheduler_defaults_stay_history_and_houses() -> None:
"""Без параметров поведение прежнее — по нему живёт суточный cian_history_backfill."""
from app.services import scheduler as scheduler_mod
captured: dict[str, Any] = {}
async def _fake_backfill(db: Any, **kwargs: Any) -> Any:
captured.update(kwargs)
# Настоящий результат, а не SimpleNamespace: посредник читает у него
# поля, которых в самодельной заглушке легко недосчитаться.
return cian_history_backfill.CianBackfillResult()
with (
patch("app.tasks.cian_history_backfill.backfill_cian_history", _fake_backfill),
patch.object(scheduler_mod.runs_mod, "update_heartbeat", MagicMock()),
patch.object(scheduler_mod.runs_mod, "mark_done", MagicMock()),
):
await scheduler_mod._execute_cian_backfill(MagicMock(), run_id=1, params={})
assert captured["listings_pending"] == "history"
assert captured["do_houses"] is True
def test_both_cian_sources_are_registered() -> None:
"""Оба source обязаны быть в реестре: планировщик держит на source один прогон,
поэтому два режима не могут делить одно имя."""
from app.services.product_handlers import build_product_handlers
handlers = build_product_handlers(MagicMock())
assert "cian_history_backfill" in handlers
assert "cian_detail_backfill" in handlers
# Гейт кук общий: он читает schedule_row["source"], а не хардкодит имя.
assert handlers["cian_detail_backfill"].pre_claim is handlers["cian_history_backfill"].pre_claim