feat(tradein/scraper): чекпоинты для avito_newbuilding_sweep — страница как единица (#3074)
All checks were successful
CI Trade-In / changes (pull_request) Successful in 8s
CI / changes (pull_request) Successful in 10s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 4m41s

Последний длинный свип без возобновления: при обрыве прогон начинался с
первой страницы, а собранное терялось целиком — save_listings вызывался
один раз на весь sweep.

Единица возобновления — страница выдачи, по образцу якорей в city sweep.
`_paginate_sweep`/`fetch_newbuildings` получили `start_page` (уже собранные
страницы не запрашиваются) и колбэк `on_page`, который вызывается только
после того, как страница пройдена до конца. Сохранение стало постраничным,
номера пройденных страниц копятся в `scrape_runs.counters.done_buckets`
мержем через `update_heartbeat`.

Два инварианта, без которых фича вредна:

1. В чекпоинт попадает только страница, чьи лоты СОХРАНЕНЫ. Отказ
   save_listings перехвачен и прогон продолжается, но отметить такую
   страницу пройденной значило бы, что следующий прогон её пропустит и
   объявления оттуда не соберутся никогда — молча, потому что прогон
   завершится штатно.

2. Подхват начинается с ПЕРВОЙ несобранной страницы, а не с max+1. Дыра в
   чекпоинте возможна ровно из-за п.1, и max+1 перепрыгнул бы её навсегда.
   Страницы после дыры перечитаются — это дешевле потери и безопасно,
   повторная запись схлопывается по dedup_hash.

Оба инварианта закрыты тестами, которые падают при их нарушении.
This commit is contained in:
bot-backend 2026-08-28 00:10:56 +03:00
parent ba57cf7c05
commit 5a410687ac
5 changed files with 400 additions and 28 deletions

View file

@ -0,0 +1,266 @@
"""Чекпоинт по страницам для avito_newbuilding_sweep (#3074).
Последний длинный свип без чекпоинтов: 27.08 прогон убит деплоем на 343-й
минуте, собранное потеряно целиком save_listings был ОДИН на весь sweep,
ждать в конце было нечего. Единица возобновления СТРАНИЦА выдачи (у функции
уже есть `pages`; цикл по страницам живёт в `AvitoScraper._paginate_sweep`).
Ключ чекпоинта номер страницы. В отличие от якорей/combo, страницы строго
последовательны (break-on-empty), поэтому resume это `start_page =
max(done_buckets) + 1`, без skip-набора произвольных элементов.
ИНВАРИАНТЫ, РАДИ КОТОРЫХ ТЕСТ:
1. Резюм пропускает уже собранные страницы БЕЗ единого запроса к источнику.
2. В чекпоинт попадает только страница, пройденная до конца оборванная
исключением/break страница НЕ фиксируется, иначе следующий прогон
пропустил бы её навсегда, причём молча (прогон завершится штатно).
3. Запись идёт мержем (`counters || :counters`), а не заменой посторонний
ключ в `counters` текущего run'а переживает запись чекпоинта.
"""
from __future__ import annotations
import os
# Settings собирается автофикстурой conftest'а и требует database_url. Выставляем
# до остальных импортов — так же, как в test_3074_avito_anchor_checkpoint.py.
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
class _FakeDb:
"""Резюм-SELECT отдаёт counters предшественника; heartbeat'ы записываются."""
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 _MergingFakeDb(_FakeDb):
"""Как _FakeDb, но heartbeat-запись мержит в текущую строку (jsonb `||`),
а не просто копится списком для теста инварианта #3 (посторонний ключ)."""
def __init__(self, prev_counters: dict[str, Any] | None = None) -> None:
super().__init__(prev_counters)
self.row: dict[str, Any] = {}
def execute(self, _stmt: Any, params: dict[str, Any] | None = None) -> Any:
if params and "counters" in params:
incoming = json.loads(params["counters"])
self.row = {**self.row, **incoming}
self.heartbeats.append(dict(self.row))
return MagicMock()
return super().execute(_stmt, params)
class _FakeAsyncSession:
def __init__(self, *_a: Any, **_kw: Any) -> None: ...
async def __aenter__(self) -> _FakeAsyncSession:
return self
async def __aexit__(self, *_e: Any) -> None:
return None
def _lot(tag: str) -> MagicMock:
m = MagicMock()
m.source_id = tag
return m
class _FakeScraper:
"""Двойник AvitoScraper: помнит реально пройденные страницы, умеет упасть."""
visited_pages: list[int] = [] # noqa: RUF012 — тестовый сборник
fail_on: int | None = None
def __init__(self, *_a: Any, **_kw: Any) -> None:
self._browser = None
self._cffi = None
async def fetch_newbuildings(
self,
*,
pages: int,
start_page: int = 1,
on_page: Any = None,
delay_override_sec: float | None = None,
) -> list[Any]:
all_lots: list[Any] = []
for page in range(start_page, pages + 1):
_FakeScraper.visited_pages.append(page)
if _FakeScraper.fail_on == page:
# Страница оборвана (сеть/парсинг) — break, on_page НЕ зовём.
break
new_lots = [_lot(f"p{page}-{i}") for i in range(2)]
all_lots.extend(new_lots)
if on_page is not None:
on_page(page, new_lots)
return all_lots
def _config() -> types.SimpleNamespace:
return types.SimpleNamespace(
scraper_fetch_mode="cffi",
scraper_proxy_url=None,
use_proxy_pool_browser=False,
browser_http_endpoint=None,
environment="test",
)
def _make_save(fail_on_call: int | None):
"""save_listings, падающий на N-м вызове — имитация отказа БД на одной странице."""
calls = {"n": 0}
def _save(*_a, **_kw):
calls["n"] += 1
if fail_on_call is not None and calls["n"] == fail_on_call:
raise RuntimeError(f"save_listings упал на странице {calls['n']}")
return (2, 0)
return _save
async def _run(
db: _FakeDb,
*,
resume_run_id: int | None,
pages: int = 4,
fail_on: int | None = None,
save_fail_on_call: int | None = None,
):
from scraper_kit.orchestration import pipeline as pl
_FakeScraper.visited_pages = []
_FakeScraper.fail_on = fail_on
with (
patch.object(pl, "AvitoScraper", _FakeScraper),
patch.object(pl, "AsyncSession", _FakeAsyncSession),
patch.object(pl, "save_listings", _make_save(save_fail_on_call)),
):
await pl.run_avito_newbuilding_sweep(
db, # type: ignore[arg-type]
run_id=8001,
config=_config(),
matcher=MagicMock(),
pages=pages,
request_delay_sec=0.0,
resume_run_id=resume_run_id,
)
return db
def _last_checkpoint(db: _FakeDb) -> list[int]:
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_resume_skips_done_pages_and_fetches_the_rest() -> None:
"""Резюм с чекпоинтом [1,2] на pages=4 обходит только страницы 3 и 4 —
ни одного запроса к уже собранным."""
db = _FakeDb({"done_buckets": [1, 2]})
await _run(db, resume_run_id=7999, pages=4)
assert _FakeScraper.visited_pages == [3, 4], (
"уже собранные страницы 1/2 не должны опрашиваться повторно"
)
assert _last_checkpoint(db) == [1, 2, 3, 4], "чекпоинт дописывается поверх унаследованного"
@pytest.mark.asyncio
async def test_without_resume_starts_from_page_one() -> None:
"""Без resume_run_id поведение прежнее — обход с первой страницы."""
db = _FakeDb(None)
await _run(db, resume_run_id=None, pages=3)
assert _FakeScraper.visited_pages == [1, 2, 3]
assert _last_checkpoint(db) == [1, 2, 3]
@pytest.mark.asyncio
async def test_failed_page_does_not_enter_checkpoint() -> None:
"""Страница, оборванная на середине (fail_on=2), НЕ считается пройденной.
Иначе следующий прогон пропустит её навсегда молча sweep завершится
штатно, просто часть выдачи не соберётся никогда.
"""
db = _FakeDb(None)
await _run(db, resume_run_id=None, pages=3, fail_on=2)
# Страница 2 была АТАКОВАНА (попытка была — visited), но не пройдена до конца.
assert _FakeScraper.visited_pages == [1, 2]
ckpt = _last_checkpoint(db)
assert ckpt == [1], "упавшая страница 2 (и не начатая 3) не должны попасть в чекпоинт"
@pytest.mark.asyncio
async def test_checkpoint_write_merges_and_preserves_foreign_key() -> None:
"""Запись чекпоинта — мерж (`counters || :counters`), не замена.
Симулируем текущую строку run'а с посторонним ключом (например, оставленным
другим писателем/предыдущим heartbeat'ом) — он обязан пережить наши записи.
"""
db = _MergingFakeDb(None)
db.row = {"foreign_key": "survives-me"}
await _run(db, resume_run_id=None, pages=1)
assert db.row.get("foreign_key") == "survives-me", (
"посторонний ключ в counters не пережил запись чекпоинта — запись была заменой, не мержем"
)
assert "done_buckets" in db.row
@pytest.mark.asyncio
async def test_page_whose_save_failed_does_not_enter_checkpoint() -> None:
"""Страница собрана, но её save упал — в чекпоинт она попасть не должна.
Отказ save_listings перехватывается и прогон продолжается (это осознанно:
одна упавшая страница не должна ронять весь sweep). Но если отметить её
пройденной, следующий прогон её пропустит, и объявления оттуда не соберутся
НИКОГДА молча, потому что прогон завершится штатно.
"""
db = _FakeDb(None)
await _run(db, resume_run_id=None, pages=3, save_fail_on_call=2)
assert _FakeScraper.visited_pages == [1, 2, 3]
ckpt = _last_checkpoint(db)
assert 2 not in ckpt, "страница с упавшим save попала в чекпоинт — покрытие потеряно молча"
assert ckpt == [1, 3]
@pytest.mark.asyncio
async def test_resume_starts_at_first_gap_not_after_the_last_page() -> None:
"""Дыра в чекпоинте перечитывается, а не перепрыгивается.
`max(done)+1` пропустил бы страницу 2 навсегда. Продолжаем с первой
несобранной: страницы после дыры перечитаются, что дешевле потери и
безопасно повторная запись схлопывается по dedup_hash.
"""
db = _FakeDb({"done_buckets": [1, 3]})
await _run(db, resume_run_id=7777, pages=4)
assert _FakeScraper.visited_pages[0] == 2, (
"подхват начался не с дыры — пропущенная страница не соберётся никогда"
)

View file

@ -495,10 +495,21 @@ async def _drive_nb_sweep(
recorder = _RunsRecorder()
db = MagicMock()
lots = [MagicMock() for _ in range(6)]
async def _fake_fetch_newbuildings(
*, pages: int, start_page: int = 1, on_page: Any = None, delay_override_sec: Any = None
) -> list[Any]:
# #3074: реальный fetch_newbuildings зовёт on_page ПОСЛЕ каждой пройденной
# страницы (инкрементальный save) — двойник имитирует одну страницу с
# ВСЕМИ 6 лотами, чтобы save_mock/счётчики остались как раньше.
if on_page is not None:
on_page(start_page, lots)
return lots
scraper = MagicMock()
scraper._cffi = None
scraper._browser = None
scraper.fetch_newbuildings = AsyncMock(return_value=lots)
scraper.fetch_newbuildings = AsyncMock(side_effect=_fake_fetch_newbuildings)
save_mock = MagicMock(side_effect=[(5, 1)])
avito_scraper_cls = MagicMock(return_value=scraper)
if capture is not None:

View file

@ -1947,6 +1947,7 @@ async def run_avito_newbuilding_sweep(
pages: int = 20,
request_delay_sec: float = 7.0,
region_code: int = DEFAULT_REGION_CODE,
resume_run_id: int | None = None,
) -> NewbuildingSweepCounters:
"""Citywide-обход ЕКБ-выборки только новостроек (novostroyka-filter) → save.
@ -1962,9 +1963,91 @@ async def run_avito_newbuilding_sweep(
Инжекция (#2135 F2): config/matcher/shutdown_requested приходят снаружи вместо
прямых импортов app.* (см. scraper_kit.contracts).
Checkpoint/resume (#3074): единица обхода — СТРАНИЦА выдачи ────────────
Последний длинный свип без чекпоинтов: прогон убивался деплоем посреди
обхода (страница 300+ при pages в разы больше дефолта), и всё собранное
терялось целиком, потому что save_listings раньше был ОДИН на весь sweep
ждать в конце было нечего. Здесь (как и у combo в yandex-свипе) save
происходит инкрементально в on_page-callback: сразу после того, как
страница пройдена до конца, её лоты сохраняются и её номер уходит в
чекпоинт `done_buckets` (список номеров страниц). Страницы единственная и
строго последовательная единица обхода (break-on-empty), поэтому resume
это просто `start_page = max(done_buckets) + 1`, без skip-набора: страница,
оборванная исключением, callback не получает и в чекпоинт не попадает
иначе следующий прогон пропустил бы её навсегда, причём молча.
"""
counters = NewbuildingSweepCounters()
_done_pages: set[int] = 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 {}
_done_pages = {int(p) for p in _pc.get("done_buckets", [])}
logger.info(
"nb-sweep run_id=%d: resuming from run %s%d страниц уже собрано",
run_id,
resume_run_id,
len(_done_pages),
)
# #3074: продолжаем с ПЕРВОЙ несобранной страницы, а не с max+1. Дыра в
# чекпоинте возможна (страница собрана, но её save упал — она не отмечена,
# а следующие отмечены), и max+1 перепрыгнул бы её навсегда: молча, потому
# что прогон завершится штатно. Уже собранные страницы после дыры при этом
# перечитаются — это дешевле потери, а повторная запись идемпотентна
# (save_listings схлопывает по dedup_hash).
_start_page = 1
while _start_page in _done_pages:
_start_page += 1
def _on_page(page: int, new_lots: list[ScrapedLot]) -> None:
"""Страница собрана и СОХРАНЕНА — только тогда чекпоинт (мерж, не замена)."""
counters.lots_fetched += len(new_lots)
_saved_ok = True
if new_lots:
try:
# #2594: citywide novostroyka-обход — только ЕКБ (см. docstring).
ins, upd = save_listings(
db,
new_lots,
matcher=matcher,
region_code=region_code,
run_id=run_id,
city=EKATERINBURG_CITY_NAME,
)
counters.lots_inserted += ins
counters.lots_updated += upd
except Exception as save_exc:
logger.exception(
"nb-sweep run_id=%d page=%d: save_listings failed: %s",
run_id,
page,
save_exc,
)
counters.errors_count += 1
_saved_ok = False
try:
db.rollback()
except Exception:
pass
# #3074: страница попадает в чекпоинт ТОЛЬКО если её лоты сохранены.
# Отказ save_listings перехвачен выше и прогон продолжается — но отметить
# страницу пройденной значило бы, что следующий прогон её пропустит, а
# объявления оттуда не соберутся НИКОГДА, причём молча: прогон завершится
# штатно. Тот же инвариант, что у якоря в avito city sweep и у бакета в
# cian full-load (_mark_bucket): в чекпоинт — только полностью собранная
# единица. Heartbeat обновляем в любом случае, иначе reap_zombies посчитает
# живой прогон мёртвым.
if _saved_ok:
_done_pages.add(page)
runs.update_heartbeat(
db, run_id, {**counters.to_dict(), "done_buckets": sorted(_done_pages)}
)
browser_mode = config.scraper_fetch_mode == "browser"
async with AsyncExitStack() as stack:
session: AsyncSession | None = None
@ -2016,8 +2099,12 @@ async def run_avito_newbuilding_sweep(
scraper._cffi = session
try:
lots: list[ScrapedLot] = await scraper.fetch_newbuildings(
# #3074: save + чекпоинт идут инкрементально внутри _on_page —
# aggregate `lots` ниже используется только для лог-сообщения.
await scraper.fetch_newbuildings(
pages=pages,
start_page=_start_page,
on_page=_on_page,
delay_override_sec=request_delay_sec,
)
except (AvitoBlockedError, AvitoRateLimitedError) as e:
@ -2027,30 +2114,6 @@ async def run_avito_newbuilding_sweep(
)
return counters
counters.lots_fetched += len(lots)
if lots:
try:
# #2594: citywide novostroyka-обход — только ЕКБ (см. docstring).
ins, upd = save_listings(
db,
lots,
matcher=matcher,
region_code=region_code,
run_id=run_id,
city=EKATERINBURG_CITY_NAME,
)
counters.lots_inserted += ins
counters.lots_updated += upd
except Exception as save_exc:
logger.exception(
"nb-sweep run_id=%d: save_listings failed: %s", run_id, save_exc
)
counters.errors_count += 1
try:
db.rollback()
except Exception:
pass
runs.update_heartbeat(db, run_id, counters.to_dict())
runs.mark_done(db, run_id, counters.to_dict())
logger.info(

View file

@ -860,6 +860,9 @@ async def _job_avito_newbuilding_sweep(
proxy_provider=ctx.proxy_provider,
pages=int(params.get("pages", 20)),
request_delay_sec=float(params.get("request_delay_sec", 7.0)),
# #3074: подхват страниц у оборванного предшественника — см.
# _job_avito_city_sweep выше.
resume_run_id=_pick_resume(db, run_id),
)

View file

@ -298,6 +298,7 @@ def _avito_bisection_config(cap: int) -> BisectionConfig:
probe_fail_policy=ProbeFailPolicy.SPLIT_OR_SKIP,
)
# HTTP 429 в curl_cffi-режиме через backconnect-прокси (mproxy.site) — НЕ IP-ban, а
# transient «слишком много одновременных соединений» (лимит 5). Проходит на коротком
# retry без ротации IP. Делаем до _AVITO_429_MAX_RETRIES коротких пауз; если они
@ -1803,13 +1804,28 @@ class AvitoScraper(BaseScraper):
*,
label: str,
delay_override_sec: float | None = None,
start_page: int = 1,
on_page: Callable[[int, list[ScrapedLot]], None] | None = None,
) -> list[ScrapedLot]:
"""Общий paginated-обход ЕКБ (citywide / novostroyka) с break-on-empty.
url_builder(page) функция построения URL страницы (citywide или
novostroyka). Сохраняет anti-block pipeline (_fetch_serp_html: firewall,
IP rotation, ретраи), дедуп по source_id, break-on-empty. label только
для логов. Не меняет наблюдаемое поведение fetch_city_wide.
для логов. start_page/on_page по умолчанию не меняют поведение
fetch_city_wide вызывает без них.
start_page (#3074): страница, с которой начать обход — страницы до
неё уже собраны предыдущим оборванным прогоном (checkpoint/resume), по
ним не делается ни одного HTTP-запроса.
on_page (#3074): опциональный callback(page: int, new_lots: list[ScrapedLot])
-> None, вызывается СРАЗУ после того, как страница пройдена до конца
(в т.ч. с пустым new_lots иначе последнюю пустую страницу
перечитывали бы вечно при resume). Страница, оборванная исключением
или break, callback не получает вызов означает «страница пройдена
до конца». Позволяет инкрементальный save за пределами этого метода:
без него собранное часами обхода терялось целиком при убийстве
процесса (SIGKILL/деплой) единственный save в конце ждать было нечем.
"""
if delay_override_sec is not None:
self.request_delay_sec = delay_override_sec
@ -1817,7 +1833,7 @@ class AvitoScraper(BaseScraper):
all_lots: list[ScrapedLot] = []
seen_ids: set[str] = set()
for page in range(1, pages + 1):
for page in range(start_page, pages + 1):
url = url_builder(page)
try:
html = await self._fetch_serp_html_with_retry(url, page)
@ -1864,6 +1880,11 @@ class AvitoScraper(BaseScraper):
len(lots) - len(new_lots),
len(all_lots),
)
if on_page is not None:
try:
on_page(page, new_lots)
except Exception:
logger.exception("avito %s: on_page callback failed page=%d", label, page)
if page < pages:
await self.sleep_between_requests()
return all_lots
@ -1873,6 +1894,8 @@ class AvitoScraper(BaseScraper):
pages: int = 30,
*,
delay_override_sec: float | None = None,
start_page: int = 1,
on_page: Callable[[int, list[ScrapedLot]], None] | None = None,
) -> list[ScrapedLot]:
"""Обход ЕКБ-выборки только новостроек (novostroyka-filter), paginated.
@ -1886,6 +1909,10 @@ class AvitoScraper(BaseScraper):
pages: максимальное число страниц (default 30).
delay_override_sec: если задан переопределяет request_delay_sec для
этого вызова.
start_page (#3074): страница, с которой продолжить обход после
checkpoint/resume см. _paginate_sweep.
on_page (#3074): callback(page, new_lots) после каждой пройденной
страницы см. _paginate_sweep.
Returns:
Список ScrapedLot новостроек (все страницы, дедуп по source_id).
@ -1895,6 +1922,8 @@ class AvitoScraper(BaseScraper):
self._build_newbuilding_url,
label="newbuilding",
delay_override_sec=delay_override_sec,
start_page=start_page,
on_page=on_page,
)
logger.info("avito fetch_newbuildings pages=%d total_lots=%d", pages, len(all_lots))
return all_lots