Merge pull request 'fix(scraper-kit): yandex _leaf под CAP не пишется в done-леджер как complete — capped_buckets отдельным счётчиком' (#3373) from fix/3368-yandex-leaf-complete-flag into main
Some checks failed
Deploy Trade-In / changes (push) Successful in 11s
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 3m59s
Deploy Trade-In / deploy (push) Successful in 2m9s
Deploy Trade-In / deploy-status (push) Successful in 1s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-backend (push) Successful in 1m39s
Deploy Trade-In / perimeter-smoke (push) Failing after 11s
Some checks failed
Deploy Trade-In / changes (push) Successful in 11s
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 3m59s
Deploy Trade-In / deploy (push) Successful in 2m9s
Deploy Trade-In / deploy-status (push) Successful in 1s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-backend (push) Successful in 1m39s
Deploy Trade-In / perimeter-smoke (push) Failing after 11s
This commit is contained in:
commit
f129d52cbe
4 changed files with 315 additions and 15 deletions
|
|
@ -89,6 +89,10 @@ class _DrainAtFirstBucket:
|
|||
self._browser = None
|
||||
self._cffi = None
|
||||
self.request_delay_sec = 0.0
|
||||
# Данные-атрибуты объявляем явно: catch-all __getattr__ ниже отдаёт корутину
|
||||
# на ЛЮБОЕ имя, и счётчик #3368 (читается в finally, ревью #3373) уехал бы
|
||||
# в counters функцией — падало бы сериализацией, а не смыслом.
|
||||
self.capped_buckets = 0
|
||||
|
||||
async def __aenter__(self) -> _DrainAtFirstBucket:
|
||||
return self
|
||||
|
|
|
|||
243
tradein-mvp/backend/tests/test_3368_yandex_leaf_cap_complete.py
Normal file
243
tradein-mvp/backend/tests/test_3368_yandex_leaf_cap_complete.py
Normal file
|
|
@ -0,0 +1,243 @@
|
|||
"""Leaf-бакет яндекса, обрезанный потолком страниц, не идёт в чекпоинт (#3368).
|
||||
|
||||
Довесок к #3362: там честной сделали degraded-ветку (`on_bucket(..., complete=False)`),
|
||||
а `_leaf` продолжал писать бакет как полный, даже когда его пагинация упиралась в
|
||||
`max_pages_per_bucket`/`_GATE_MAX_PAGES_CAP`. С containment-гейтом (#3358/#3359) такой
|
||||
ключ покрывает СВОЙ интервал целиком → резюм больше не заходит в полосу, чей хвост не
|
||||
читали ни разу. У cian тот же случай считается честно (`complete = pages_needed <=
|
||||
max_pages`, cian/serp.py).
|
||||
|
||||
Переполненный leaf возможен только там, где бисекции дробить нечем: размах меньше
|
||||
`_YANDEX_SPLIT_MIN_BRACKET`, открытый верхний брекет или потолок глубины. Здесь берётся
|
||||
первый случай: плотная выдача (totalItems=2000 > cap=500) делится до размаха 499 999 и
|
||||
дальше делиться не может.
|
||||
|
||||
Сеть не нужна: `_fetch_page_json` подменяется счётчиком. Проверки ПО ЗНАЧЕНИЮ — что
|
||||
попало в done-леджер и сколько запросов сделал резюм.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
import os
|
||||
import types
|
||||
from typing import Any
|
||||
from unittest.mock import MagicMock, patch
|
||||
|
||||
import pytest
|
||||
|
||||
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
|
||||
|
||||
from scraper_kit.base import ScrapedLot
|
||||
from scraper_kit.providers.yandex.serp import YandexRealtyScraper
|
||||
|
||||
from app.services.scraper_adapters import RealScraperConfig
|
||||
|
||||
_ROOMS = "2"
|
||||
|
||||
|
||||
def _page(total_items: int) -> dict[str, Any]:
|
||||
return {
|
||||
"response": {
|
||||
"search": {
|
||||
"offers": {
|
||||
"entities": [],
|
||||
"pager": {
|
||||
"totalItems": total_items,
|
||||
"totalPages": max(1, total_items // 20),
|
||||
"page": 0,
|
||||
},
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
def _walk(
|
||||
total_items: int,
|
||||
max_pages_per_bucket: int,
|
||||
skip_buckets: set[str] | None = None,
|
||||
fail_pages: set[int] | None = None,
|
||||
) -> tuple[list[tuple[str, bool]], int, int]:
|
||||
"""Прогнать бисекцию [4M, 5M) и вернуть (отметки бакетов, запросы, capped_buckets).
|
||||
|
||||
Колбэк — той же формы, что pipeline._on_bucket: третий позиционный аргумент =
|
||||
признак полноты; в done-леджер `_mark_bucket` кладёт только complete=True.
|
||||
"""
|
||||
s = YandexRealtyScraper(RealScraperConfig())
|
||||
s.request_delay_sec = 0.0
|
||||
calls = [0]
|
||||
|
||||
async def fake_fetch(
|
||||
rooms: str | None,
|
||||
page: int,
|
||||
price_min: int | None,
|
||||
price_max: int | None,
|
||||
new_flat: str = "NO",
|
||||
) -> dict[str, Any] | None:
|
||||
calls[0] += 1
|
||||
if fail_pages and page in fail_pages:
|
||||
return None # капча/тарпит/сеть — страница выпала из пагинации
|
||||
return _page(total_items)
|
||||
|
||||
s._fetch_page_json = fake_fetch # type: ignore[method-assign]
|
||||
marked: list[tuple[str, bool]] = []
|
||||
|
||||
def on_bucket(key: str, _count: int, complete: bool = True) -> None:
|
||||
marked.append((key, complete))
|
||||
|
||||
seen: dict[str, ScrapedLot] = {}
|
||||
asyncio.run(
|
||||
s._walk_price_range(
|
||||
rooms=_ROOMS,
|
||||
lo=4_000_000,
|
||||
hi=4_999_999,
|
||||
seen=seen,
|
||||
price_cap_per_bucket=500,
|
||||
max_pages_per_bucket=max_pages_per_bucket,
|
||||
on_bucket=on_bucket,
|
||||
skip_buckets=skip_buckets,
|
||||
)
|
||||
)
|
||||
# getattr, а не атрибут напрямую: без счётчика тест обязан краснеть НЕВЕРНЫМ
|
||||
# ЗНАЧЕНИЕМ (ключ в леджере / ноль запросов на резюме), а не AttributeError'ом —
|
||||
# «возможности нет» неотличимо от «проверка не проведена».
|
||||
return marked, calls[0], getattr(s, "capped_buckets", 0)
|
||||
|
||||
|
||||
def test_capped_leaf_stays_out_of_done_ledger() -> None:
|
||||
"""Приёмка: leaf, которому нужно больше страниц, чем потолок, — не в леджере."""
|
||||
marked, _, capped = _walk(total_items=2000, max_pages_per_bucket=1)
|
||||
assert marked, "leaf не вызвал on_bucket — тест ничего не проверяет"
|
||||
ledger = {key for key, complete in marked if complete}
|
||||
assert not ledger, (
|
||||
f"бакеты {sorted(ledger)} прочитаны на 1 страницу из 100 (totalItems=2000), "
|
||||
"но помечены complete — их интервал зачтётся containment-гейтом целиком"
|
||||
)
|
||||
assert capped == len(marked), (
|
||||
f"capped_buckets={capped} при {len(marked)} обрезанных бакетах — "
|
||||
"счётчик прогона не покажет, что полоса недобрана по построению"
|
||||
)
|
||||
|
||||
|
||||
def test_fully_read_leaf_goes_into_done_ledger() -> None:
|
||||
"""Контроль честности: дочитанный до конца бакет по-прежнему чекпоинтится."""
|
||||
marked, _, capped = _walk(total_items=15, max_pages_per_bucket=1)
|
||||
ledger = {key for key, complete in marked if complete}
|
||||
assert ledger, (
|
||||
"бакет из одной страницы (totalItems=15) не попал в леджер — резюм будет "
|
||||
"перечитывать уже собранную территорию"
|
||||
)
|
||||
assert capped == 0, f"полный бакет посчитан обрезанным (capped_buckets={capped})"
|
||||
|
||||
|
||||
def test_capped_band_is_rewalked_on_resume() -> None:
|
||||
"""По значению: полоса с обрезанным leaf'ом на резюме снова обходится."""
|
||||
marked, _, _ = _walk(total_items=2000, max_pages_per_bucket=1)
|
||||
ledger = {key for key, complete in marked if complete}
|
||||
_, calls, _ = _walk(total_items=2000, max_pages_per_bucket=1, skip_buckets=ledger or None)
|
||||
assert calls > 0, (
|
||||
"резюм не сделал ни одного запроса по полосе, у которой прочитана 1 страница "
|
||||
"из 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]}"
|
||||
)
|
||||
|
|
@ -3807,6 +3807,11 @@ class YandexFullLoadCounters:
|
|||
# Бакеты, собранные ЧАСТИЧНО (probe провалился → degraded-пагинация): в
|
||||
# чекпоинт не пишутся, следующий прогон перечитает их целиком.
|
||||
partial_buckets: int = 0
|
||||
# Подмножество partial_buckets: бакет неполон ПО ПОСТРОЕНИЮ (#3368) — total
|
||||
# требует больше страниц, чем потолок, а бисекции делить его уже нечем.
|
||||
# Отдельный счётчик, потому что лечение другое: не «повторить прогон», а
|
||||
# снизить min_bracket / поднять потолок страниц.
|
||||
capped_buckets: int = 0
|
||||
|
||||
def to_dict(self) -> dict[str, int]:
|
||||
return {f.name: getattr(self, f.name) for f in fields(self)}
|
||||
|
|
@ -3943,13 +3948,20 @@ async def run_yandex_full_load(
|
|||
async with YandexRealtyScraper(config, proxy_provider=proxy_provider) as scraper:
|
||||
scraper.request_delay_sec = request_delay_sec
|
||||
|
||||
await scraper.fetch_all_secondary(
|
||||
price_cap_per_bucket=price_cap_per_bucket,
|
||||
concurrency=concurrency,
|
||||
on_bucket=_on_bucket,
|
||||
on_progress=_on_progress,
|
||||
skip_buckets=skip_set if skip_set else None,
|
||||
)
|
||||
try:
|
||||
await scraper.fetch_all_secondary(
|
||||
price_cap_per_bucket=price_cap_per_bucket,
|
||||
concurrency=concurrency,
|
||||
on_bucket=_on_bucket,
|
||||
on_progress=_on_progress,
|
||||
skip_buckets=skip_set if skip_set else None,
|
||||
)
|
||||
finally:
|
||||
# #3368: leaf'ы, обрезанные потолком страниц (в чекпоинт не попали).
|
||||
# В finally, а не после await: cancel/shutdown прилетают сюда
|
||||
# RuntimeError'ом из _on_bucket, и присваивание после вызова
|
||||
# пропускалось ровно в тех прогонах, чей итог и надо объяснить.
|
||||
counters.capped_buckets = getattr(scraper, "capped_buckets", 0)
|
||||
|
||||
logger.info(
|
||||
"yandex-full-load run_id=%d: fetch done — unique=%d ins=%d upd=%d",
|
||||
|
|
|
|||
|
|
@ -531,6 +531,11 @@ class YandexRealtyScraper(BaseScraper):
|
|||
# блокировки — см. _track_gate_result.
|
||||
self.gate_fetch_attempts: int = 0
|
||||
self.gate_fetch_failures: int = 0
|
||||
# #3368: leaf-бакеты, чья пагинация упёрлась в потолок страниц (дробить
|
||||
# бисекции уже нечем). Не пишутся в чекпоинт → читаются вызывающим для
|
||||
# counters прогона. Атрибут, а не аргумент on_bucket — тот же довод, что у
|
||||
# cian.last_dropped_nb: у колбэка есть внешние реализации.
|
||||
self.capped_buckets: int = 0
|
||||
|
||||
def _track_gate_result(self, ok: bool) -> None:
|
||||
"""Учёт исхода одного top-level gate-API запроса (#2625).
|
||||
|
|
@ -1220,13 +1225,37 @@ class YandexRealtyScraper(BaseScraper):
|
|||
logger.debug("yandex gate: skip bucket %s (checkpoint)", bucket_key)
|
||||
return
|
||||
|
||||
total_pages = min(
|
||||
math.ceil(total / 20),
|
||||
_GATE_MAX_PAGES_CAP,
|
||||
max_pages_per_bucket,
|
||||
)
|
||||
pages_needed = math.ceil(total / _GATE_PAGE_SIZE)
|
||||
max_pages = min(_GATE_MAX_PAGES_CAP, max_pages_per_bucket)
|
||||
total_pages = min(pages_needed, max_pages)
|
||||
# Полнота бакета (#3368): у обрезанного потолком страниц ключ ТОТ ЖЕ, что у
|
||||
# собранного целиком. Паритет с cian — там `complete = dropped_pages == 0 and
|
||||
# pages_needed <= max_pages` (cian/serp.py), т.е. ДВА пути частичности: потолок
|
||||
# страниц и выпавшая страница пагинации (payload=None — капча/тарпит/сеть).
|
||||
# Ниже считаются оба. complete=False → бакет не идёт в done-леджер, иначе его
|
||||
# интервал в containment-гейте (#3359) склеился бы с соседними и резюм пропустил
|
||||
# бы полосу, чей хвост никогда не читали. Бакет доходит сюда переполненным
|
||||
# только там, где бисекции делить больше нечем (размах < min_bracket,
|
||||
# открытый верхний брекет, потолок глубины) — см. walk_price_range.
|
||||
capped = pages_needed > max_pages
|
||||
if capped:
|
||||
self.capped_buckets += 1
|
||||
logger.warning(
|
||||
"yandex gate: leaf bucket %s НЕПОЛОН по построению — total=%d требует "
|
||||
"%d страниц при потолке %d; дробить дальше нечем, в чекпоинт не пишем "
|
||||
"(capped_buckets=%d)",
|
||||
bucket_key,
|
||||
total,
|
||||
pages_needed,
|
||||
max_pages,
|
||||
self.capped_buckets,
|
||||
)
|
||||
logger.info(
|
||||
"yandex gate: leaf bucket %s total=%d pages=%d", bucket_key, total, total_pages
|
||||
"yandex gate: leaf bucket %s total=%d pages=%d capped=%s",
|
||||
bucket_key,
|
||||
total,
|
||||
total_pages,
|
||||
capped,
|
||||
)
|
||||
|
||||
# Add probe lots (page 1 already fetched)
|
||||
|
|
@ -1236,17 +1265,20 @@ class YandexRealtyScraper(BaseScraper):
|
|||
|
||||
if total_pages <= 1:
|
||||
if on_bucket is not None:
|
||||
on_bucket(bucket_key, len(seen))
|
||||
on_bucket(bucket_key, len(seen), not capped)
|
||||
return
|
||||
|
||||
# Paginate pages 2..total_pages with concurrency
|
||||
sem = asyncio.Semaphore(concurrency)
|
||||
dropped_pages = 0
|
||||
|
||||
async def _fetch_leaf_page(pg: int) -> list[ScrapedLot]:
|
||||
nonlocal dropped_pages
|
||||
async with sem:
|
||||
payload = await self._fetch_page_json(rooms, pg, lo_param, phi)
|
||||
await asyncio.sleep(self.request_delay_sec)
|
||||
if payload is None:
|
||||
dropped_pages += 1
|
||||
return []
|
||||
return _parse_gate_json(payload, page_param=pg)
|
||||
|
||||
|
|
@ -1257,8 +1289,17 @@ class YandexRealtyScraper(BaseScraper):
|
|||
if lot.source_id and lot.source_id not in seen:
|
||||
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:
|
||||
on_bucket(bucket_key, len(seen))
|
||||
on_bucket(bucket_key, len(seen), complete)
|
||||
|
||||
# #3359: гейт ПЕРЕД probe (как у avito, #3315). Ключи yandex'а — `_combo_label`,
|
||||
# т.е. «rooms:lo-hi» с «None» вместо открытого потолка: другой разделитель и
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue