feat(tradein/scraper): чекпоинт по якорям для cian_city_sweep (#3074) #3117
3 changed files with 217 additions and 1 deletions
159
tradein-mvp/backend/tests/test_3074_cian_anchor_checkpoint.py
Normal file
159
tradein-mvp/backend/tests/test_3074_cian_anchor_checkpoint.py
Normal file
|
|
@ -0,0 +1,159 @@
|
|||
"""Чекпоинт по якорям для cian_city_sweep (#3074).
|
||||
|
||||
Замер за 60 дней, по которому выбран источник: 65 прогонов, среднее 35 минут,
|
||||
максимум 72, две отмены деплоем. Пятиминутного дренажа (#3029) на такие прогоны
|
||||
не хватает — убитый на 35-й минуте сбор начинался заново с первого якоря.
|
||||
|
||||
Ключ чекпоинта — ИМЯ якоря, а не индекс: состав списка зависит от `city_slug`
|
||||
(областные свипы идут по своим наборам), позиция между городами не устойчива.
|
||||
|
||||
ИНВАРИАНТ. В чекпоинт попадает только якорь, пройденный до конца. У циана эта
|
||||
граница уже проведена потоком управления: все ветки отказа делают `return` или
|
||||
`continue`, и до записи не доходят. Этим он отличается от avito-свипа, где успех
|
||||
и неудача сходились в одной строке и потребовался отдельный флаг. Тест ниже
|
||||
проверяет, что граница не нарушена: упавший якорь не должен попасть в чекпоинт,
|
||||
иначе следующий прогон пропустит его навсегда — молча, потому что прогон
|
||||
завершится штатно, просто часть города не соберётся.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
|
||||
# Settings собирается автофикстурой conftest'а и требует database_url.
|
||||
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
|
||||
|
||||
import json
|
||||
import types
|
||||
from typing import Any
|
||||
from unittest.mock import MagicMock, patch
|
||||
|
||||
import pytest
|
||||
|
||||
ANCHOR_A = (56.83, 60.60, "ekb-center")
|
||||
ANCHOR_B = (56.79, 60.63, "ekb-south")
|
||||
|
||||
|
||||
class _FakeDb:
|
||||
def __init__(self, prev_counters: dict[str, Any] | None = None) -> None:
|
||||
self.prev_counters = prev_counters or {}
|
||||
self.heartbeats: list[dict[str, Any]] = []
|
||||
|
||||
def execute(self, _stmt: Any, params: dict[str, Any] | None = None) -> Any:
|
||||
if params and "counters" in params:
|
||||
self.heartbeats.append(json.loads(params["counters"]))
|
||||
return MagicMock()
|
||||
if params and "rid" in params:
|
||||
return MagicMock(fetchone=lambda: types.SimpleNamespace(counters=self.prev_counters))
|
||||
return MagicMock()
|
||||
|
||||
def commit(self) -> None: ...
|
||||
def rollback(self) -> None: ...
|
||||
|
||||
|
||||
class _FakeScraper:
|
||||
"""Двойник CianScraper: помнит, за какими якорями реально ходили."""
|
||||
|
||||
visited: list[tuple[float, float]] = [] # noqa: RUF012 — тестовый сборник
|
||||
raise_on: tuple[float, float] | None = None
|
||||
|
||||
def __init__(self, *_a: Any, **_kw: Any) -> None:
|
||||
# Счётчики, которые конвейер читает у скрапера после каждого якоря
|
||||
# (#2625 — диагностика «все якоря упали»). Без них падает не проверяемая
|
||||
# логика, а сам двойник.
|
||||
self.state_extraction_attempts = 1
|
||||
self.state_extraction_failures = 0
|
||||
self.request_delay_sec = 0.0
|
||||
|
||||
async def __aenter__(self) -> _FakeScraper:
|
||||
return self
|
||||
|
||||
async def __aexit__(self, *_e: Any) -> None:
|
||||
return None
|
||||
|
||||
async def fetch_around_multi_room(self, lat: float, lon: float, *_a: Any, **_kw: Any) -> list:
|
||||
_FakeScraper.visited.append((lat, lon))
|
||||
if _FakeScraper.raise_on == (lat, lon):
|
||||
raise RuntimeError("якорь упал")
|
||||
return []
|
||||
|
||||
|
||||
def _config() -> types.SimpleNamespace:
|
||||
return types.SimpleNamespace(
|
||||
scraper_proxy_url=None,
|
||||
scraper_fetch_mode="cffi",
|
||||
use_proxy_pool_browser=False,
|
||||
browser_http_endpoint=None,
|
||||
environment="test",
|
||||
)
|
||||
|
||||
|
||||
async def _run(prev: dict[str, Any] | None, raise_on: tuple[float, float] | None = None) -> _FakeDb:
|
||||
from scraper_kit.orchestration import pipeline as pl
|
||||
|
||||
_FakeScraper.visited = []
|
||||
_FakeScraper.raise_on = raise_on
|
||||
db = _FakeDb(prev)
|
||||
|
||||
with (
|
||||
patch.object(pl, "CianScraper", _FakeScraper),
|
||||
patch.object(pl, "save_listings", lambda *_a, **_kw: (0, 0)),
|
||||
patch.object(pl.runs, "is_cancelled", lambda *_a: False),
|
||||
):
|
||||
await pl.run_cian_city_sweep(
|
||||
db, # type: ignore[arg-type]
|
||||
run_id=8001,
|
||||
config=_config(),
|
||||
matcher=MagicMock(),
|
||||
anchors=[ANCHOR_A, ANCHOR_B],
|
||||
enrich_houses=False,
|
||||
detail_top_n=0,
|
||||
request_delay_sec=0.0,
|
||||
resume_run_id=7999 if prev is not None else None,
|
||||
)
|
||||
return db
|
||||
|
||||
|
||||
def _last_checkpoint(db: _FakeDb) -> list[str]:
|
||||
with_ckpt = [hb for hb in db.heartbeats if "done_buckets" in hb]
|
||||
assert with_ckpt, "ни один heartbeat не унёс done_buckets — чекпоинт не персистится"
|
||||
return with_ckpt[-1]["done_buckets"]
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_checkpointed_anchor_is_skipped_without_a_single_request() -> None:
|
||||
"""Якорь из чекпоинта не опрашивается вовсе."""
|
||||
await _run({"done_buckets": ["ekb-center"]})
|
||||
assert (ANCHOR_A[0], ANCHOR_A[1]) not in _FakeScraper.visited, (
|
||||
"якорь из чекпоинта всё-таки опрашивали"
|
||||
)
|
||||
assert (ANCHOR_B[0], ANCHOR_B[1]) in _FakeScraper.visited, "второй якорь не обошли"
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_checkpoint_accumulates_over_inherited() -> None:
|
||||
"""Пройденный якорь дописывается поверх унаследованных, а не затирает их."""
|
||||
db = await _run({"done_buckets": ["ekb-center"]})
|
||||
assert _last_checkpoint(db) == ["ekb-center", "ekb-south"]
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_without_resume_all_anchors_are_visited() -> None:
|
||||
"""Без чекпоинта поведение прежнее — обходятся все якоря."""
|
||||
db = await _run(None)
|
||||
assert len(_FakeScraper.visited) == 2
|
||||
assert _last_checkpoint(db) == ["ekb-center", "ekb-south"]
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_failed_anchor_does_not_enter_checkpoint() -> None:
|
||||
"""Упавший якорь НЕ считается пройденным.
|
||||
|
||||
У циана это обеспечено потоком управления, а не флагом: ветка отказа делает
|
||||
`continue` до записи. Тест сторожит именно это — рефакторинг, сливающий ветки
|
||||
в одну, сломает инвариант незаметно.
|
||||
"""
|
||||
db = await _run(None, raise_on=(ANCHOR_A[0], ANCHOR_A[1]))
|
||||
ckpt = _last_checkpoint(db)
|
||||
assert "ekb-center" not in ckpt, "упавший якорь попал в чекпоинт"
|
||||
assert "ekb-south" in ckpt, "исправный якорь не зафиксирован"
|
||||
|
|
@ -2766,6 +2766,7 @@ async def run_cian_city_sweep(
|
|||
enrich_houses: bool = True,
|
||||
newbuilding_only: bool = True,
|
||||
region_code: int = DEFAULT_REGION_CODE,
|
||||
resume_run_id: int | None = None,
|
||||
) -> CianCitySweepCounters:
|
||||
"""Cian newbuilding city sweep: SERP → detail(+price-history) → newbuilding/houses.
|
||||
|
||||
|
|
@ -2782,6 +2783,33 @@ async def run_cian_city_sweep(
|
|||
mark_done вызывается ВСЕГДА (finally outer).
|
||||
"""
|
||||
_anchors = anchors if anchors is not None else EKB_ANCHORS
|
||||
|
||||
# ── Checkpoint/resume (#3074): единица обхода — ЯКОРЬ ──────────────────────
|
||||
# Замер за 60 дней: 65 прогонов, среднее 35 минут, максимум 72, две отмены
|
||||
# деплоем. Пятиминутного дренажа (#3029) на такие прогоны не хватает, и
|
||||
# убитый на 35-й минуте сбор начинается заново с первого якоря.
|
||||
#
|
||||
# Ключ — ИМЯ якоря, а не индекс: состав списка зависит от `city_slug`
|
||||
# (областные свипы идут по своим наборам), позиция между городами не
|
||||
# устойчива. По той же причине гарда по числу якорей не нужна — в отличие от
|
||||
# combo-чекпоинта яндекса, где ключ якоря не содержал.
|
||||
_skip_anchors: set[str] = set()
|
||||
if resume_run_id is not None:
|
||||
_prev = db.execute(
|
||||
text("SELECT counters FROM scrape_runs WHERE id = CAST(:rid AS bigint)"),
|
||||
{"rid": resume_run_id},
|
||||
).fetchone()
|
||||
if _prev is not None and _prev.counters:
|
||||
_pc: dict = _prev.counters if isinstance(_prev.counters, dict) else {}
|
||||
_skip_anchors = set(_pc.get("done_buckets", []))
|
||||
logger.info(
|
||||
"cian-sweep run_id=%d: resuming from run %s — %d якорей уже пройдено",
|
||||
run_id,
|
||||
resume_run_id,
|
||||
len(_skip_anchors),
|
||||
)
|
||||
_done_anchors: set[str] = set(_skip_anchors)
|
||||
|
||||
# city_slug (#12): region_id города-цели → CianScraper.city_region_id скоупит SERP
|
||||
# на город вместо дефолтного ЕКБ. None/неизвестный slug → ЕКБ-дефолт в конструкторе.
|
||||
_loc = get_city_location(city_slug)
|
||||
|
|
@ -2827,6 +2855,19 @@ async def run_cian_city_sweep(
|
|||
|
||||
try:
|
||||
for idx, (lat, lon, name) in enumerate(_anchors, start=1):
|
||||
if name in _skip_anchors:
|
||||
# #3074: якорь собран предыдущим оборванным прогоном — ни одного
|
||||
# HTTP-запроса. Счётчик двигаем, чтобы `anchors_done` продолжал
|
||||
# означать «докуда дошли по списку», а не «сколько собрал этот run».
|
||||
counters.anchors_done = idx
|
||||
logger.info(
|
||||
"cian-sweep run_id=%d: anchor #%d/%d (%s) пропущен — есть в чекпоинте",
|
||||
run_id,
|
||||
idx,
|
||||
len(_anchors),
|
||||
name,
|
||||
)
|
||||
continue
|
||||
# Cooperative cancel перед каждым anchor
|
||||
if runs.is_cancelled(db, run_id):
|
||||
logger.info(
|
||||
|
|
@ -3178,7 +3219,19 @@ async def run_cian_city_sweep(
|
|||
continue
|
||||
|
||||
counters.anchors_done = idx
|
||||
runs.update_heartbeat(db, run_id, counters.to_dict())
|
||||
# #3074: сюда попадают ТОЛЬКО якоря, пройденные до конца. Все ветки
|
||||
# отказа выше либо `return`, либо `continue` — в чекпоинт они не
|
||||
# заходят. Разница с avito-свипом, где успех и неудача сходились в
|
||||
# одной строке и потребовался отдельный флаг: здесь граница уже
|
||||
# проведена самим потоком управления, и её достаточно не нарушать.
|
||||
#
|
||||
# Записать упавший якорь пройденным значило бы, что следующий прогон
|
||||
# пропустит его навсегда — молча, потому что прогон завершится
|
||||
# штатно, просто часть города не соберётся.
|
||||
_done_anchors.add(name)
|
||||
runs.update_heartbeat(
|
||||
db, run_id, {**counters.to_dict(), "done_buckets": sorted(_done_anchors)}
|
||||
)
|
||||
|
||||
# Пауза между anchor'ами (поверх per-request sleep внутри scraper'а)
|
||||
if idx < len(_anchors):
|
||||
|
|
|
|||
|
|
@ -920,6 +920,10 @@ async def _job_cian_city_sweep(
|
|||
detail_top_n=int(params.get("detail_top_n", 10)),
|
||||
enrich_houses=bool(params.get("enrich_houses", True)),
|
||||
newbuilding_only=bool(params.get("newbuilding_only", True)),
|
||||
# #3074: подхват якорей у оборванного предшественника. Ключ чекпоинта —
|
||||
# имя якоря, от их количества не зависит, поэтому гарда как у combo-
|
||||
# чекпоинта яндекса здесь не требуется.
|
||||
resume_run_id=_pick_resume(db, run_id),
|
||||
)
|
||||
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue