Merge pull request 'feat(tradein/scraper): чекпоинт по якорям для cian_city_sweep (#3074)' (#3117) from feat/3074-cian-anchor-checkpoint into main
All checks were successful
Deploy Trade-In / changes (push) Successful in 11s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-backend (push) Successful in 1m39s
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 3m48s
Deploy Trade-In / deploy (push) Successful in 1m36s
Deploy Trade-In / deploy-status (push) Successful in 1s
Deploy Trade-In / perimeter-smoke (push) Successful in 10s

This commit is contained in:
bot-backend 2026-08-26 14:28:15 +00:00
commit 54438039e5
3 changed files with 217 additions and 1 deletions

View 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, "исправный якорь не зафиксирован"

View file

@ -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):

View file

@ -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),
)