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,
|
enrich_houses: bool = True,
|
||||||
newbuilding_only: bool = True,
|
newbuilding_only: bool = True,
|
||||||
region_code: int = DEFAULT_REGION_CODE,
|
region_code: int = DEFAULT_REGION_CODE,
|
||||||
|
resume_run_id: int | None = None,
|
||||||
) -> CianCitySweepCounters:
|
) -> CianCitySweepCounters:
|
||||||
"""Cian newbuilding city sweep: SERP → detail(+price-history) → newbuilding/houses.
|
"""Cian newbuilding city sweep: SERP → detail(+price-history) → newbuilding/houses.
|
||||||
|
|
||||||
|
|
@ -2782,6 +2783,33 @@ async def run_cian_city_sweep(
|
||||||
mark_done вызывается ВСЕГДА (finally outer).
|
mark_done вызывается ВСЕГДА (finally outer).
|
||||||
"""
|
"""
|
||||||
_anchors = anchors if anchors is not None else EKB_ANCHORS
|
_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
|
# city_slug (#12): region_id города-цели → CianScraper.city_region_id скоупит SERP
|
||||||
# на город вместо дефолтного ЕКБ. None/неизвестный slug → ЕКБ-дефолт в конструкторе.
|
# на город вместо дефолтного ЕКБ. None/неизвестный slug → ЕКБ-дефолт в конструкторе.
|
||||||
_loc = get_city_location(city_slug)
|
_loc = get_city_location(city_slug)
|
||||||
|
|
@ -2827,6 +2855,19 @@ async def run_cian_city_sweep(
|
||||||
|
|
||||||
try:
|
try:
|
||||||
for idx, (lat, lon, name) in enumerate(_anchors, start=1):
|
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
|
# Cooperative cancel перед каждым anchor
|
||||||
if runs.is_cancelled(db, run_id):
|
if runs.is_cancelled(db, run_id):
|
||||||
logger.info(
|
logger.info(
|
||||||
|
|
@ -3178,7 +3219,19 @@ async def run_cian_city_sweep(
|
||||||
continue
|
continue
|
||||||
|
|
||||||
counters.anchors_done = idx
|
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'а)
|
# Пауза между anchor'ами (поверх per-request sleep внутри scraper'а)
|
||||||
if idx < len(_anchors):
|
if idx < len(_anchors):
|
||||||
|
|
|
||||||
|
|
@ -920,6 +920,10 @@ async def _job_cian_city_sweep(
|
||||||
detail_top_n=int(params.get("detail_top_n", 10)),
|
detail_top_n=int(params.get("detail_top_n", 10)),
|
||||||
enrich_houses=bool(params.get("enrich_houses", True)),
|
enrich_houses=bool(params.get("enrich_houses", True)),
|
||||||
newbuilding_only=bool(params.get("newbuilding_only", 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