fix(scraper-kit/yandex): выпавшая страница leaf'а делает бакет неполным; capped_buckets переживает дрейн
All checks were successful
CI Trade-In / changes (pull_request) Successful in 9s
CI / changes (pull_request) Successful in 12s
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 / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 5m5s

Ревью #3373, два minor:
1. `_leaf` считал полноту только по потолку страниц — упавшая страница
   (`payload=None` → `[]`) молча уходила в чекпоинт как собранная. Теперь
   паритет с cian по-настоящему: `complete = not capped and dropped_pages == 0`.
2. `counters.capped_buckets` присваивался ПОСЛЕ await — cancel/shutdown
   (RuntimeError из `_on_bucket`) уносил управление мимо строки. Перенесено в
   `finally`, счётчик виден в финальном payload дрейна.

Нит: `ceil(total / 20)` → `_GATE_PAGE_SIZE`.

Тесты по значению: бакет с 1 выпавшей страницей из 3 — не в done-леджере
(+ контроль: 3 успешных страницы по-прежнему complete); дрейн после обрезанного
бакета → `capped_buckets == 1` в финальных counters. Фейк `_DrainAtFirstBucket`
(#3355) теперь объявляет `capped_buckets` явно: его catch-all `__getattr__`
отдавал корутину на любое имя и ронял сериализацию counters.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
bot-backend 2026-09-06 01:16:14 +05:00
parent 5ad0d1a304
commit 9519c89d92
4 changed files with 151 additions and 20 deletions

View file

@ -89,6 +89,10 @@ class _DrainAtFirstBucket:
self._browser = None self._browser = None
self._cffi = None self._cffi = None
self.request_delay_sec = 0.0 self.request_delay_sec = 0.0
# Данные-атрибуты объявляем явно: catch-all __getattr__ ниже отдаёт корутину
# на ЛЮБОЕ имя, и счётчик #3368 (читается в finally, ревью #3373) уехал бы
# в counters функцией — падало бы сериализацией, а не смыслом.
self.capped_buckets = 0
async def __aenter__(self) -> _DrainAtFirstBucket: async def __aenter__(self) -> _DrainAtFirstBucket:
return self return self

View file

@ -19,8 +19,13 @@ max_pages`, cian/serp.py).
from __future__ import annotations from __future__ import annotations
import asyncio import asyncio
import json
import os import os
import types
from typing import Any from typing import Any
from unittest.mock import MagicMock, patch
import pytest
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test") os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
@ -53,6 +58,7 @@ def _walk(
total_items: int, total_items: int,
max_pages_per_bucket: int, max_pages_per_bucket: int,
skip_buckets: set[str] | None = None, skip_buckets: set[str] | None = None,
fail_pages: set[int] | None = None,
) -> tuple[list[tuple[str, bool]], int, int]: ) -> tuple[list[tuple[str, bool]], int, int]:
"""Прогнать бисекцию [4M, 5M) и вернуть (отметки бакетов, запросы, capped_buckets). """Прогнать бисекцию [4M, 5M) и вернуть (отметки бакетов, запросы, capped_buckets).
@ -69,8 +75,10 @@ def _walk(
price_min: int | None, price_min: int | None,
price_max: int | None, price_max: int | None,
new_flat: str = "NO", new_flat: str = "NO",
) -> dict[str, Any]: ) -> dict[str, Any] | None:
calls[0] += 1 calls[0] += 1
if fail_pages and page in fail_pages:
return None # капча/тарпит/сеть — страница выпала из пагинации
return _page(total_items) return _page(total_items)
s._fetch_page_json = fake_fetch # type: ignore[method-assign] s._fetch_page_json = fake_fetch # type: ignore[method-assign]
@ -133,3 +141,103 @@ def test_capped_band_is_rewalked_on_resume() -> None:
"резюм не сделал ни одного запроса по полосе, у которой прочитана 1 страница " "резюм не сделал ни одного запроса по полосе, у которой прочитана 1 страница "
"из 100 — недобранный хвост потерян навсегда" "из 100 — недобранный хвост потерян навсегда"
) )
def test_dropped_page_makes_leaf_incomplete() -> None:
"""Ревью #3373: выпавшая страница пагинации — второй путь частичности (как у cian).
`_fetch_page_json` None (капча/тарпит/сеть) молча давала пустой список, и бакет
из 2 прочитанных страниц вместо 3 уходил в чекпоинт полным. У cian оба пути
считаются вместе: `complete = dropped_pages == 0 and pages_needed <= max_pages`.
"""
ok, _, _ = _walk(total_items=45, max_pages_per_bucket=5)
assert [c for _k, c in ok] == [True], (
"контроль: бакет из 3 успешных страниц обязан быть complete — "
f"иначе тест ниже красный по любой причине, отметки={ok}"
)
marked, _, capped = _walk(total_items=45, max_pages_per_bucket=5, fail_pages={3})
assert marked, "leaf не вызвал on_bucket — тест ничего не проверяет"
ledger = {key for key, complete in marked if complete}
assert not ledger, (
f"бакеты {sorted(ledger)} прочитаны на 2 страницы из 3 (одна выпала с "
"payload=None), но помечены complete — их интервал зачтётся containment-гейтом"
)
# pipeline._mark_bucket увеличивает partial_buckets ровно на complete=False.
assert len([1 for _k, complete in marked if not complete]) == 1, (
f"ровно один бакет обязан лечь в partial_buckets, отметки={marked}"
)
assert capped == 0, (
f"выпавшая страница посчитана потолком (capped_buckets={capped}) — "
"лечение у этих случаев разное: повторить прогон vs снизить min_bracket"
)
class _FakeDb:
"""Пишет каждый UPDATE с counters (как в test_3355_drain_mark_full_loads.py)."""
def __init__(self) -> None:
self.writes: list[dict[str, Any]] = []
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> Any:
if params and "counters" in params:
self.writes.append(json.loads(params["counters"]))
return MagicMock()
def commit(self) -> None: ...
def rollback(self) -> None: ...
class _CappedThenDrain:
"""Один обрезанный потолком leaf, затем дрейн на границе бакета."""
def __init__(self, *_a: Any, **_kw: Any) -> None:
self.request_delay_sec = 0.0
self.capped_buckets = 0
self.gate_fetch_attempts = 0
self.gate_fetch_failures = 0
async def __aenter__(self) -> _CappedThenDrain:
return self
async def __aexit__(self, *_e: Any) -> None:
return None
async def fetch_all_secondary(self, **kw: Any) -> None:
self.capped_buckets += 1
kw["on_bucket"]("2:4000000-4999999", [], False)
raise AssertionError("_on_bucket не оборвал прогон при shutdown_requested()")
@pytest.mark.asyncio
async def test_capped_counter_survives_shutdown() -> None:
"""Ревью #3373: счётчик писался ПОСЛЕ await — cancel/shutdown его терял."""
from scraper_kit.orchestration import pipeline as pl
db = _FakeDb()
with (
patch.object(pl, "YandexRealtyScraper", _CappedThenDrain),
patch.object(pl.runs, "is_cancelled", lambda *_a: False),
):
await pl.run_yandex_full_load(
db, # type: ignore[arg-type]
run_id=3373,
config=types.SimpleNamespace(
scraper_fetch_mode="cffi",
scraper_proxy_url=None,
scraper_skip_seen_today=False,
use_proxy_pool_browser=False,
browser_http_endpoint=None,
environment="test",
),
matcher=MagicMock(),
enrichment=MagicMock(),
request_delay_sec=0.0,
shutdown_requested=lambda: True,
)
assert db.writes, "дрейн не оставил ни одной записи counters"
assert db.writes[-1].get("capped_buckets") == 1, (
"оборванный дрейном прогон не сообщил про обрезанный потолком бакет: "
f"финальные counters={db.writes[-1]}"
)

View file

@ -3948,15 +3948,20 @@ async def run_yandex_full_load(
async with YandexRealtyScraper(config, proxy_provider=proxy_provider) as scraper: async with YandexRealtyScraper(config, proxy_provider=proxy_provider) as scraper:
scraper.request_delay_sec = request_delay_sec scraper.request_delay_sec = request_delay_sec
await scraper.fetch_all_secondary( try:
price_cap_per_bucket=price_cap_per_bucket, await scraper.fetch_all_secondary(
concurrency=concurrency, price_cap_per_bucket=price_cap_per_bucket,
on_bucket=_on_bucket, concurrency=concurrency,
on_progress=_on_progress, on_bucket=_on_bucket,
skip_buckets=skip_set if skip_set else None, on_progress=_on_progress,
) skip_buckets=skip_set if skip_set else None,
# #3368: leaf'ы, обрезанные потолком страниц (в чекпоинт не попали). )
counters.capped_buckets = getattr(scraper, "capped_buckets", 0) finally:
# #3368: leaf'ы, обрезанные потолком страниц (в чекпоинт не попали).
# В finally, а не после await: cancel/shutdown прилетают сюда
# RuntimeError'ом из _on_bucket, и присваивание после вызова
# пропускалось ровно в тех прогонах, чей итог и надо объяснить.
counters.capped_buckets = getattr(scraper, "capped_buckets", 0)
logger.info( logger.info(
"yandex-full-load run_id=%d: fetch done — unique=%d ins=%d upd=%d", "yandex-full-load run_id=%d: fetch done — unique=%d ins=%d upd=%d",

View file

@ -1225,18 +1225,20 @@ class YandexRealtyScraper(BaseScraper):
logger.debug("yandex gate: skip bucket %s (checkpoint)", bucket_key) logger.debug("yandex gate: skip bucket %s (checkpoint)", bucket_key)
return return
pages_needed = math.ceil(total / 20) pages_needed = math.ceil(total / _GATE_PAGE_SIZE)
max_pages = min(_GATE_MAX_PAGES_CAP, max_pages_per_bucket) max_pages = min(_GATE_MAX_PAGES_CAP, max_pages_per_bucket)
total_pages = min(pages_needed, max_pages) total_pages = min(pages_needed, max_pages)
# Полнота бакета (#3368): у обрезанного потолком страниц ключ ТОТ ЖЕ, что у # Полнота бакета (#3368): у обрезанного потолком страниц ключ ТОТ ЖЕ, что у
# собранного целиком, — как у cian (`complete = pages_needed <= max_pages`). # собранного целиком. Паритет с cian — там `complete = dropped_pages == 0 and
# complete=False → бакет не идёт в done-леджер, иначе его интервал в # pages_needed <= max_pages` (cian/serp.py), т.е. ДВА пути частичности: потолок
# containment-гейте (#3359) склеился бы с соседними и резюм пропустил бы # страниц и выпавшая страница пагинации (payload=None — капча/тарпит/сеть).
# полосу, чей хвост никогда не читали. Бакет доходит сюда переполненным # Ниже считаются оба. complete=False → бакет не идёт в done-леджер, иначе его
# интервал в containment-гейте (#3359) склеился бы с соседними и резюм пропустил
# бы полосу, чей хвост никогда не читали. Бакет доходит сюда переполненным
# только там, где бисекции делить больше нечем (размах < min_bracket, # только там, где бисекции делить больше нечем (размах < min_bracket,
# открытый верхний брекет, потолок глубины) — см. walk_price_range. # открытый верхний брекет, потолок глубины) — см. walk_price_range.
complete = pages_needed <= max_pages capped = pages_needed > max_pages
if not complete: if capped:
self.capped_buckets += 1 self.capped_buckets += 1
logger.warning( logger.warning(
"yandex gate: leaf bucket %s НЕПОЛОН по построению — total=%d требует " "yandex gate: leaf bucket %s НЕПОЛОН по построению — total=%d требует "
@ -1249,11 +1251,11 @@ class YandexRealtyScraper(BaseScraper):
self.capped_buckets, self.capped_buckets,
) )
logger.info( logger.info(
"yandex gate: leaf bucket %s total=%d pages=%d complete=%s", "yandex gate: leaf bucket %s total=%d pages=%d capped=%s",
bucket_key, bucket_key,
total, total,
total_pages, total_pages,
complete, capped,
) )
# Add probe lots (page 1 already fetched) # Add probe lots (page 1 already fetched)
@ -1263,17 +1265,20 @@ class YandexRealtyScraper(BaseScraper):
if total_pages <= 1: if total_pages <= 1:
if on_bucket is not None: if on_bucket is not None:
on_bucket(bucket_key, len(seen), complete) on_bucket(bucket_key, len(seen), not capped)
return return
# Paginate pages 2..total_pages with concurrency # Paginate pages 2..total_pages with concurrency
sem = asyncio.Semaphore(concurrency) sem = asyncio.Semaphore(concurrency)
dropped_pages = 0
async def _fetch_leaf_page(pg: int) -> list[ScrapedLot]: async def _fetch_leaf_page(pg: int) -> list[ScrapedLot]:
nonlocal dropped_pages
async with sem: async with sem:
payload = await self._fetch_page_json(rooms, pg, lo_param, phi) payload = await self._fetch_page_json(rooms, pg, lo_param, phi)
await asyncio.sleep(self.request_delay_sec) await asyncio.sleep(self.request_delay_sec)
if payload is None: if payload is None:
dropped_pages += 1
return [] return []
return _parse_gate_json(payload, page_param=pg) return _parse_gate_json(payload, page_param=pg)
@ -1284,6 +1289,15 @@ class YandexRealtyScraper(BaseScraper):
if lot.source_id and lot.source_id not in seen: if lot.source_id and lot.source_id not in seen:
seen[lot.source_id] = lot seen[lot.source_id] = lot
complete = not capped and dropped_pages == 0
if dropped_pages:
logger.warning(
"yandex gate: leaf bucket %s%d стр. из %d выпали (payload=None), "
"бакет НЕПОЛОН, в чекпоинт не пишем",
bucket_key,
dropped_pages,
total_pages - 1,
)
if on_bucket is not None: if on_bucket is not None:
on_bucket(bucket_key, len(seen), complete) on_bucket(bucket_key, len(seen), complete)