feat(tradein/domclick): чекпоинты для city_sweep — корзина как единица возобновления (#3118) #3140
4 changed files with 236 additions and 42 deletions
144
tradein-mvp/backend/tests/test_3118_domclick_checkpoint.py
Normal file
144
tradein-mvp/backend/tests/test_3118_domclick_checkpoint.py
Normal file
|
|
@ -0,0 +1,144 @@
|
|||
"""#3118/#3074: чекпоинты для domclick_city_sweep — корзина ROOM_BUCKETS как единица.
|
||||
|
||||
Прод-факты (#3118): QRATOR обрывает свип внутри 1-2-й корзины при ЛЮБОМ старте
|
||||
(`buckets_completed ≤ 1` из 6 во всех прогонах), banned-прогоны при этом собирают
|
||||
59–1815 лотов. Сдвиг #2854 лишь распределяет потери; чекпоинт превращает
|
||||
случайную ротацию в систематический обход — шесть прогонов закрывают шесть
|
||||
корзин. Моё раннее «чекпоинтить нечего» (комментарий в #3074) опровергнуто
|
||||
данными #3118 — этот файл и есть исправление того вывода кодом.
|
||||
|
||||
Три слоя (зеркально yandex #3074): провайдер skip_buckets + имена завершённых;
|
||||
пайплайн — resume-чтение и done_buckets в counters (мерж jsonb, финализаторы
|
||||
не затирают); планировщик — generic _pick_resume (#2845) бесплатно.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
|
||||
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
|
||||
|
||||
from types import SimpleNamespace
|
||||
from typing import Any
|
||||
from unittest.mock import MagicMock, patch
|
||||
|
||||
from scraper_kit.orchestration import scheduler as sched
|
||||
from scraper_kit.providers.domclick.serp import ROOM_BUCKETS, DomClickBlockedError, DomClickScraper
|
||||
|
||||
# ── 1. Провайдер: скип, имена, порядок с ротацией #2854 ──────────────────────
|
||||
|
||||
|
||||
async def _no_sleep() -> None:
|
||||
pass
|
||||
|
||||
|
||||
def _scraper() -> DomClickScraper:
|
||||
s = DomClickScraper(SimpleNamespace(scraper_proxy_url=None))
|
||||
s.sleep_between_requests = _no_sleep # type: ignore[method-assign]
|
||||
return s
|
||||
|
||||
|
||||
async def _run_fetch_city(
|
||||
scraper: DomClickScraper,
|
||||
*,
|
||||
start: int = 0,
|
||||
skip: set[str] | None = None,
|
||||
block_on: str | None = None,
|
||||
) -> list[str]:
|
||||
"""fetch_city с замоканными сеткой и обходом корзины; возвращает порядок
|
||||
корзин, по которым РЕАЛЬНО пошёл обход."""
|
||||
visited: list[str] = []
|
||||
|
||||
async def _fake_sweep_bucket(*, rooms: str, **_kw: Any) -> None:
|
||||
visited.append(rooms)
|
||||
if block_on is not None and rooms == block_on:
|
||||
raise DomClickBlockedError("QRATOR")
|
||||
|
||||
class _FakeFetcherCtx:
|
||||
async def __aenter__(self) -> SimpleNamespace:
|
||||
return SimpleNamespace(report_ban=lambda *_a, **_k: None)
|
||||
|
||||
async def __aexit__(self, *_exc: Any) -> None:
|
||||
return None
|
||||
|
||||
with (
|
||||
patch.object(scraper, "_sweep_bucket", _fake_sweep_bucket),
|
||||
patch(
|
||||
"scraper_kit.providers._base.build_browser_fetcher",
|
||||
lambda *_a, **_k: _FakeFetcherCtx(),
|
||||
),
|
||||
):
|
||||
await scraper.fetch_city(city_id=1, pages=1, start_bucket_index=start, skip_buckets=skip)
|
||||
return visited
|
||||
|
||||
|
||||
async def test_skip_buckets_are_never_fetched() -> None:
|
||||
"""Корзины из чекпоинта не получают ни одного обхода; buckets_total честно
|
||||
сжимается до объёма ЭТОГО прогона (иначе honest-status читал бы
|
||||
возобновлённый прогон как вечно-частичный)."""
|
||||
s = _scraper()
|
||||
visited = await _run_fetch_city(s, start=0, skip={"st", "1", "2"})
|
||||
assert visited == ["3", "4", "5+"], visited
|
||||
assert s.buckets_total == 3
|
||||
assert s.completed_buckets == ["3", "4", "5+"]
|
||||
|
||||
|
||||
async def test_block_midway_records_completed_names() -> None:
|
||||
"""Блок на 2-й корзине: имена завершённых до блока — источник чекпоинта."""
|
||||
s = _scraper()
|
||||
visited = await _run_fetch_city(s, start=1, block_on="2")
|
||||
# старт со сдвигом 1: порядок 1,2,... — блок на '2' после завершения '1'
|
||||
assert visited[:2] == ["1", "2"]
|
||||
assert s.completed_buckets == ["1"], s.completed_buckets
|
||||
assert s.blocked is True
|
||||
|
||||
|
||||
async def test_all_skipped_is_honest_noop() -> None:
|
||||
"""Цепочка накопила все 6 корзин → пустой обход без падения (гард)."""
|
||||
s = _scraper()
|
||||
visited = await _run_fetch_city(s, skip=set(ROOM_BUCKETS))
|
||||
assert visited == []
|
||||
assert s.completed_buckets == []
|
||||
|
||||
|
||||
# ── 2. Планировщик: диспатч отдаёт точку (зеркало test_930/test_3074) ────────
|
||||
|
||||
|
||||
def _candidate() -> SimpleNamespace:
|
||||
return SimpleNamespace(
|
||||
prev_id=6001,
|
||||
prev_status="banned",
|
||||
prev_counters={"done_buckets": ["st", "1"]},
|
||||
same_params=True,
|
||||
age_h=20.0,
|
||||
interval_days="1",
|
||||
)
|
||||
|
||||
|
||||
class _FakeDb:
|
||||
def __init__(self, row: Any) -> None:
|
||||
self.row = row
|
||||
|
||||
def execute(self, _stmt: Any, params: dict[str, Any] | None = None) -> Any:
|
||||
if params and "counters" in params:
|
||||
return MagicMock()
|
||||
return MagicMock(fetchone=lambda: self.row)
|
||||
|
||||
def commit(self) -> None:
|
||||
pass
|
||||
|
||||
|
||||
async def test_scheduler_hands_checkpoint_to_domclick_sweep() -> None:
|
||||
"""Красный на main по значению: kwargs без resume_run_id → None != 6001."""
|
||||
db = _FakeDb(_candidate())
|
||||
captured: dict[str, Any] = {}
|
||||
|
||||
async def _spy(*_a: Any, **kw: Any) -> None:
|
||||
captured.update(kw)
|
||||
|
||||
with patch.object(sched, "run_domclick_city_sweep", _spy):
|
||||
await sched._job_domclick_city_sweep(db, 7000, {}, MagicMock())
|
||||
|
||||
assert captured.get("resume_run_id") == 6001, (
|
||||
"планировщик не отдал чекпоинт домклик-свипу — прогон пойдёт с нуля"
|
||||
)
|
||||
|
|
@ -221,6 +221,7 @@ def ban_kind_of_exception(exc: BaseException) -> str:
|
|||
return BAN_KIND_PLATFORM
|
||||
return BAN_KIND_UNKNOWN
|
||||
|
||||
|
||||
# #2160: константы для расчёта watchdog-таймаута Cian city sweep. При
|
||||
# USE_PROXY_POOL_BROWSER=true каждый SERP-фетч идёт через camoufox с relaunch при смене
|
||||
# прокси (page.goto timeout 60s + overhead) = 13-45s/страница, а якорь cian = 4 room-buckets
|
||||
|
|
@ -259,6 +260,7 @@ def _cian_anchor_timeout_s(
|
|||
+ (_CIAN_HOUSES_BUDGET_S if enrich_houses else 0.0),
|
||||
)
|
||||
|
||||
|
||||
# Default anchors ЕКБ — 5 точек покрытия города
|
||||
EKB_ANCHORS: list[tuple[float, float, str]] = [
|
||||
(56.8400, 60.6050, "Центр"),
|
||||
|
|
@ -1461,15 +1463,13 @@ async def run_avito_city_sweep(
|
|||
>= _AVITO_PIPELINE_CONSECUTIVE_BLOCK_ABORT
|
||||
and sweep_rotations_done < config.avito_proxy_max_rotations
|
||||
):
|
||||
rotated, sweep_rotations_done = (
|
||||
await _try_rotate_within_budget(
|
||||
config,
|
||||
reason=(
|
||||
f"city-sweep house "
|
||||
f"consecutive={sweep_consecutive_house_blocks}"
|
||||
),
|
||||
rotations_done=sweep_rotations_done,
|
||||
)
|
||||
rotated, sweep_rotations_done = await _try_rotate_within_budget(
|
||||
config,
|
||||
reason=(
|
||||
f"city-sweep house "
|
||||
f"consecutive={sweep_consecutive_house_blocks}"
|
||||
),
|
||||
rotations_done=sweep_rotations_done,
|
||||
)
|
||||
if rotated:
|
||||
sweep_consecutive_house_blocks = 0
|
||||
|
|
@ -1585,15 +1585,12 @@ async def run_avito_city_sweep(
|
|||
)
|
||||
if consecutive_blocks >= 3:
|
||||
# #1790: попытка ротации перед abort'ом.
|
||||
rotated, sweep_rotations_done = (
|
||||
await _try_rotate_within_budget(
|
||||
config,
|
||||
reason=(
|
||||
f"city-sweep detail "
|
||||
f"consecutive={consecutive_blocks}"
|
||||
),
|
||||
rotations_done=sweep_rotations_done,
|
||||
)
|
||||
rotated, sweep_rotations_done = await _try_rotate_within_budget(
|
||||
config,
|
||||
reason=(
|
||||
f"city-sweep detail consecutive={consecutive_blocks}"
|
||||
),
|
||||
rotations_done=sweep_rotations_done,
|
||||
)
|
||||
if rotated:
|
||||
consecutive_blocks = 0
|
||||
|
|
@ -1686,19 +1683,19 @@ async def run_avito_city_sweep(
|
|||
)
|
||||
if (
|
||||
sweep_consecutive_detail_house_blocks >= 3
|
||||
and sweep_rotations_done
|
||||
< config.avito_proxy_max_rotations
|
||||
and sweep_rotations_done < config.avito_proxy_max_rotations
|
||||
):
|
||||
rotated, sweep_rotations_done = (
|
||||
await _try_rotate_within_budget(
|
||||
config,
|
||||
reason=(
|
||||
"city-sweep house(detail) "
|
||||
f"consecutive="
|
||||
f"{sweep_consecutive_detail_house_blocks}"
|
||||
),
|
||||
rotations_done=sweep_rotations_done,
|
||||
)
|
||||
(
|
||||
rotated,
|
||||
sweep_rotations_done,
|
||||
) = await _try_rotate_within_budget(
|
||||
config,
|
||||
reason=(
|
||||
"city-sweep house(detail) "
|
||||
f"consecutive="
|
||||
f"{sweep_consecutive_detail_house_blocks}"
|
||||
),
|
||||
rotations_done=sweep_rotations_done,
|
||||
)
|
||||
if rotated:
|
||||
sweep_consecutive_detail_house_blocks = 0
|
||||
|
|
@ -3411,9 +3408,7 @@ async def run_cian_full_load(
|
|||
raise RuntimeError("shutdown")
|
||||
if not lots:
|
||||
_mark_bucket(bucket_key, complete)
|
||||
runs.update_heartbeat(
|
||||
db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)}
|
||||
)
|
||||
runs.update_heartbeat(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)})
|
||||
return
|
||||
inserted, updated = save_listings(
|
||||
db,
|
||||
|
|
@ -3502,8 +3497,7 @@ async def run_cian_full_load(
|
|||
|
||||
counters.dropped_novostroyki = getattr(scraper, "last_dropped_nb", 0)
|
||||
logger.info(
|
||||
"cian-full-load run_id=%d: fetch done — unique=%d ins=%d upd=%d "
|
||||
"dropped_novostroyki=%d",
|
||||
"cian-full-load run_id=%d: fetch done — unique=%d ins=%d upd=%d dropped_novostroyki=%d",
|
||||
run_id,
|
||||
counters.unique_fetched,
|
||||
counters.saved_inserted,
|
||||
|
|
@ -3752,9 +3746,7 @@ async def run_yandex_full_load(
|
|||
raise RuntimeError("shutdown")
|
||||
if not lots:
|
||||
done.add(bucket_key)
|
||||
runs.update_heartbeat(
|
||||
db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)}
|
||||
)
|
||||
runs.update_heartbeat(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)})
|
||||
return
|
||||
inserted, updated = save_listings(
|
||||
db,
|
||||
|
|
@ -4036,9 +4028,7 @@ async def run_avito_full_load(
|
|||
raise RuntimeError("shutdown")
|
||||
if not lots:
|
||||
_mark_bucket(bucket_key, complete)
|
||||
runs.update_heartbeat(
|
||||
db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)}
|
||||
)
|
||||
runs.update_heartbeat(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)})
|
||||
return
|
||||
inserted, updated = save_listings(
|
||||
db,
|
||||
|
|
@ -4228,6 +4218,7 @@ async def run_domclick_city_sweep(
|
|||
pages: int = 100,
|
||||
request_delay_sec: float | None = None,
|
||||
region_code: int = DEFAULT_REGION_CODE,
|
||||
resume_run_id: int | None = None,
|
||||
) -> DomClickCitySweepCounters:
|
||||
"""DomClick citywide sweep через BFF JSON API.
|
||||
|
||||
|
|
@ -4290,6 +4281,28 @@ async def run_domclick_city_sweep(
|
|||
|
||||
lots: list[ScrapedLot] = []
|
||||
|
||||
# ── Checkpoint/resume (#3118): корзина ROOM_BUCKETS — единица обхода ──
|
||||
# Вместе со сдвигом #2854 подхват превращает случайную ротацию в
|
||||
# систематический обход: банимый на 1-й корзине источник закрывает все
|
||||
# шесть за несколько прогонов вместо повторов случайных.
|
||||
skip_buckets: set[str] = set()
|
||||
if resume_run_id is not None:
|
||||
_prev_row = db.execute(
|
||||
text("SELECT counters FROM scrape_runs WHERE id = CAST(:rid AS bigint)"),
|
||||
{"rid": resume_run_id},
|
||||
).fetchone()
|
||||
if _prev_row is not None and _prev_row.counters:
|
||||
_prev_counters: dict = (
|
||||
_prev_row.counters if isinstance(_prev_row.counters, dict) else {}
|
||||
)
|
||||
skip_buckets = set(_prev_counters.get("done_buckets", []))
|
||||
logger.info(
|
||||
"domclick-sweep run_id=%d: resuming from run %s — %d корзин уже собрано",
|
||||
run_id,
|
||||
resume_run_id,
|
||||
len(skip_buckets),
|
||||
)
|
||||
|
||||
async def _domclick_phase() -> None:
|
||||
"""Единственная citywide-фаза: fetch_city + save."""
|
||||
nonlocal lots
|
||||
|
|
@ -4315,6 +4328,7 @@ async def run_domclick_city_sweep(
|
|||
rooms=rooms,
|
||||
pages=pages,
|
||||
start_bucket_index=run_id % len(ROOM_BUCKETS),
|
||||
skip_buckets=skip_buckets or None,
|
||||
)
|
||||
counters.lots_fetched += len(lots)
|
||||
if lots:
|
||||
|
|
@ -4358,6 +4372,11 @@ async def run_domclick_city_sweep(
|
|||
# скрейпер живая, а его счётчик показывает, докуда прогон дошёл.
|
||||
counters.buckets_completed = _s.buckets_completed
|
||||
counters.buckets_total = _s.buckets_total
|
||||
# #3118: чекпоинт = унаследованное ∪ завершённое в этом прогоне.
|
||||
# Пишем heartbeat'ом СЕЙЧАС (мерж jsonb) — финализаторы ключ не
|
||||
# затирают, и оборванный болезнью финализации прогон его не теряет.
|
||||
_done_now = sorted(skip_buckets | set(_s.completed_buckets))
|
||||
runs.update_heartbeat(db, run_id, {**counters.to_dict(), "done_buckets": _done_now})
|
||||
counters.bucket_start_index = _s.bucket_start_index
|
||||
|
||||
# pages_fetched: worst-case число страниц (buckets × pages cap).
|
||||
|
|
|
|||
|
|
@ -1015,6 +1015,7 @@ async def _job_domclick_city_sweep(
|
|||
rooms=params.get("rooms"),
|
||||
pages=int(params.get("pages_per_anchor", 5)),
|
||||
request_delay_sec=float(params.get("request_delay_sec", 6.0)),
|
||||
resume_run_id=_pick_resume(db, run_id),
|
||||
)
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -280,6 +280,9 @@ class DomClickScraper(BaseScraper):
|
|||
# Поэтому прогон, прошедший 3 бакета из 6, был неотличим от полного.
|
||||
self.buckets_total: int = len(ROOM_BUCKETS)
|
||||
self.buckets_completed: int = 0
|
||||
# #3118: ИМЕНА фактически завершённых корзин этого прогона — источник
|
||||
# чекпоинта done_buckets (счётчик выше даёт число, resume нужен состав).
|
||||
self.completed_buckets: list[str] = []
|
||||
# #2854: с какой корзины начался обход. Без этого прогон неатрибутируем —
|
||||
# по данным нельзя отличить «корзина не собралась» от «до неё не дошли».
|
||||
self.bucket_start_index: int = 0
|
||||
|
|
@ -304,6 +307,7 @@ class DomClickScraper(BaseScraper):
|
|||
rooms: list[int] | None = None,
|
||||
pages: int = 100,
|
||||
start_bucket_index: int = 0,
|
||||
skip_buckets: set[str] | None = None,
|
||||
) -> list[ScrapedLot]:
|
||||
"""Citywide sweep через BFF JSON API.
|
||||
|
||||
|
|
@ -327,6 +331,13 @@ class DomClickScraper(BaseScraper):
|
|||
всем корзинам вместо одной. Работает при единственном свободном узле —
|
||||
в отличие от ротации lease, которой сейчас упираться некуда: в пуле 3
|
||||
включённых узла, 2 забанены Домкликом.
|
||||
skip_buckets: корзины, уже собранные предыдущим оборванным прогоном
|
||||
(#3118, чекпоинт done_buckets) — по ним ни одного запроса. Вместе
|
||||
со сдвигом #2854 это превращает случайную ротацию в систематический
|
||||
обход: шесть прогонов закрывают шесть корзин вместо повторов.
|
||||
buckets_total при этом = числу корзин В ЭТОМ прогоне (без
|
||||
скипнутых) — иначе honest-status читал бы возобновлённый прогон
|
||||
как вечно-частичный.
|
||||
|
||||
Returns:
|
||||
Дедуплицированный по source_id список ScrapedLot.
|
||||
|
|
@ -353,7 +364,23 @@ class DomClickScraper(BaseScraper):
|
|||
# run_id, а он ничем не ограничен.
|
||||
offset = start_bucket_index % len(ROOM_BUCKETS)
|
||||
buckets = ROOM_BUCKETS[offset:] + ROOM_BUCKETS[:offset]
|
||||
if skip_buckets:
|
||||
skipped = [b for b in buckets if b in skip_buckets]
|
||||
buckets = tuple(b for b in buckets if b not in skip_buckets)
|
||||
self.buckets_total = len(buckets)
|
||||
logger.info(
|
||||
"domklik: чекпоинт #3118 — корзины %r уже собраны предшественником, "
|
||||
"в этом прогоне %d из %d",
|
||||
skipped,
|
||||
len(buckets),
|
||||
len(ROOM_BUCKETS),
|
||||
)
|
||||
self.bucket_start_index = offset
|
||||
if not buckets:
|
||||
# Цепочка resume накопила ВСЕ корзины — работать не с чем; честный
|
||||
# пустой результат (вызывающий финализирует done с нулём новых).
|
||||
logger.info("domklik: все корзины уже в чекпоинте предшественника — no-op")
|
||||
return out_lots
|
||||
logger.info(
|
||||
"domklik: обход начинается с корзины rooms=%r (сдвиг %d из %d, #2854)",
|
||||
buckets[0],
|
||||
|
|
@ -409,6 +436,7 @@ class DomClickScraper(BaseScraper):
|
|||
)
|
||||
continue
|
||||
self.buckets_completed += 1
|
||||
self.completed_buckets.append(bucket)
|
||||
|
||||
logger.info(
|
||||
"domklik: fetch_city done city_id=%d total=%d buckets=%d/%d "
|
||||
|
|
@ -684,7 +712,9 @@ class DomClickScraper(BaseScraper):
|
|||
#
|
||||
# Риска нет: newbuilding_id бэкендом не читается нигде (проверено),
|
||||
# в upsert защищён COALESCE, на проде заполнен только у avito.
|
||||
newbuilding_id: str | None = (flat_complex.get("slug") or None) if flat_complex else None
|
||||
newbuilding_id: str | None = (
|
||||
(flat_complex.get("slug") or None) if flat_complex else None
|
||||
)
|
||||
raw_payload: dict[str, Any] = {
|
||||
"isRosreestrApproved": item.get("isRosreestrApproved"),
|
||||
"squarePrice": square_price_raw,
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue