fix(yandex): degraded-бакет не зачитывается как пройденный
All checks were successful
CI Trade-In / changes (pull_request) Successful in 8s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI / changes (pull_request) Successful in 12s
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 5m11s
All checks were successful
CI Trade-In / changes (pull_request) Successful in 8s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI / changes (pull_request) Successful in 12s
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 5m11s
Ревью #3362: в _degraded (probe провалился → пагинация до пустоты, обрезаемая max_pages_per_bucket) yandex звал on_bucket БЕЗ признака полноты, и ключ падал в done-леджер наравне с честно добранным листом. До containment-гейта это был точечный skip одного ключа; теперь ключ покрывает ИНТЕРВАЛ и сливается со смежными — недобранная после отказа полоса больше никогда не переобходится. Тот же путь, что у cian: третий позиционный аргумент complete, в чекпоинт пишет только _mark_bucket(..., True); partial_buckets вынесен в счётчики прогона (виден в heartbeat), лоты и cancel/shutdown-проверки не трогаем. Плюс комментарий смежности в bisection.py приводил полуоткрытый пример, споря с «hi ВКЛЮЧИТЕЛЬНА» в том же докстринге.
This commit is contained in:
parent
35c7ea8492
commit
4184f751bb
4 changed files with 93 additions and 6 deletions
|
|
@ -121,3 +121,63 @@ def test_fully_done_room_resumes_with_zero_requests() -> None:
|
|||
)
|
||||
)
|
||||
assert calls[0] == 0, f"резюм готовой комнатности сделал {calls[0]} запрос(ов) вместо нуля"
|
||||
|
||||
|
||||
def _walk_degraded(skip_buckets: set[str] | None) -> tuple[list[tuple[str, bool]], int]:
|
||||
"""Прогон с мёртвой сетью (probe провалился → degraded-ветка).
|
||||
|
||||
Возвращает (что бакет-колбэк отметил, число сетевых запросов). Колбэк — той же
|
||||
формы, что pipeline._on_bucket: третий позиционный аргумент = признак полноты.
|
||||
"""
|
||||
s, calls = _scraper()
|
||||
|
||||
async def dead_fetch(*_a: Any, **_k: Any) -> None:
|
||||
calls[0] += 1
|
||||
return None
|
||||
|
||||
async def no_rotate() -> bool:
|
||||
return False
|
||||
|
||||
s._fetch_page_json = dead_fetch # type: ignore[method-assign]
|
||||
s._rotate_ip = no_rotate # type: ignore[method-assign]
|
||||
marked: list[tuple[str, bool]] = []
|
||||
|
||||
def on_bucket(key: str, _count: int, complete: bool = True) -> None:
|
||||
marked.append((key, complete))
|
||||
|
||||
asyncio.run(
|
||||
s._walk_price_range(
|
||||
rooms=_ROOMS,
|
||||
lo=4_000_000,
|
||||
hi=4_999_999,
|
||||
seen={},
|
||||
price_cap_per_bucket=500,
|
||||
max_pages_per_bucket=1,
|
||||
on_bucket=on_bucket,
|
||||
skip_buckets=skip_buckets,
|
||||
)
|
||||
)
|
||||
return marked, calls[0]
|
||||
|
||||
|
||||
def test_degraded_bucket_stays_out_of_done_ledger() -> None:
|
||||
"""Бакет с провалившимся probe НЕ зачитывается как пройденный."""
|
||||
marked, _ = _walk_degraded(None)
|
||||
assert marked, "degraded-ветка не вызвала on_bucket — тест ничего не проверяет"
|
||||
ledger = {key for key, complete in marked if complete}
|
||||
assert not ledger, (
|
||||
f"бакеты {sorted(ledger)} собраны degraded-пагинацией (полнота неизвестна), "
|
||||
"но помечены complete — в леджере их интервал сольётся с соседними и резюм "
|
||||
"не переобойдёт недобор"
|
||||
)
|
||||
|
||||
|
||||
def test_degraded_band_is_rewalked_on_resume() -> None:
|
||||
"""По значению: леджер после degraded-прогона не гасит эту полосу на резюме."""
|
||||
marked, _ = _walk_degraded(None)
|
||||
ledger = {key for key, complete in marked if complete}
|
||||
_, calls = _walk_degraded(ledger or None)
|
||||
assert calls > 0, (
|
||||
"резюм не сделал ни одного запроса по полосе, собранной лишь частично — "
|
||||
"best-effort территория после отказа потеряна навсегда"
|
||||
)
|
||||
|
|
|
|||
|
|
@ -3791,6 +3791,9 @@ class YandexFullLoadCounters:
|
|||
saved_updated: int = 0
|
||||
price_history_rows: int = 0
|
||||
errors_count: int = 0
|
||||
# Бакеты, собранные ЧАСТИЧНО (probe провалился → degraded-пагинация): в
|
||||
# чекпоинт не пишутся, следующий прогон перечитает их целиком.
|
||||
partial_buckets: int = 0
|
||||
|
||||
def to_dict(self) -> dict[str, int]:
|
||||
return {f.name: getattr(self, f.name) for f in fields(self)}
|
||||
|
|
@ -3843,8 +3846,27 @@ async def run_yandex_full_load(
|
|||
)
|
||||
done: set[str] = set(skip_set)
|
||||
|
||||
def _on_bucket(bucket_key: str, lots: list) -> None: # type: ignore[type-arg]
|
||||
"""Инкрементальный save после каждого leaf-бакета."""
|
||||
def _mark_bucket(bucket_key: str, complete: bool) -> None:
|
||||
"""В done-леджер пишем ТОЛЬКО полностью собранный бакет."""
|
||||
if complete:
|
||||
done.add(bucket_key)
|
||||
return
|
||||
counters.partial_buckets += 1
|
||||
logger.warning(
|
||||
"yandex-full-load run_id=%d: bucket=%s собран ЧАСТИЧНО — в чекпоинт НЕ пишем, "
|
||||
"следующий прогон перечитает его целиком (partial_buckets=%d)",
|
||||
run_id,
|
||||
bucket_key,
|
||||
counters.partial_buckets,
|
||||
)
|
||||
|
||||
def _on_bucket(bucket_key: str, lots: list, complete: bool = True) -> None: # type: ignore[type-arg]
|
||||
"""Инкрементальный save после каждого leaf-бакета.
|
||||
|
||||
complete=False (degraded-ветка бисекции) → в чекпоинт бакет не пишем,
|
||||
следующий прогон перечитает полосу целиком. Дефолт True — для вызывающих
|
||||
без пагинации. Тот же контракт, что у cian (см. run_cian_full_load).
|
||||
"""
|
||||
nonlocal done
|
||||
if runs.is_cancelled(db, run_id):
|
||||
logger.info("yandex-full-load run_id=%d: cancel detected in on_bucket", run_id)
|
||||
|
|
@ -3857,7 +3879,7 @@ async def run_yandex_full_load(
|
|||
)
|
||||
raise RuntimeError("shutdown")
|
||||
if not lots:
|
||||
done.add(bucket_key)
|
||||
_mark_bucket(bucket_key, complete)
|
||||
runs.update_heartbeat(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)})
|
||||
return
|
||||
inserted, updated = save_listings(
|
||||
|
|
@ -3889,7 +3911,7 @@ async def run_yandex_full_load(
|
|||
db.rollback()
|
||||
except Exception:
|
||||
pass
|
||||
done.add(bucket_key)
|
||||
_mark_bucket(bucket_key, complete)
|
||||
runs.update_heartbeat(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)})
|
||||
logger.info(
|
||||
"yandex-full-load run_id=%d: bucket=%s saved ins=%d upd=%d total_unique=%d",
|
||||
|
|
|
|||
|
|
@ -186,7 +186,7 @@ def done_range_skipper(
|
|||
|
||||
merged: list[tuple[int, float]] = []
|
||||
for lo, hi in sorted(parsed):
|
||||
if merged and lo <= merged[-1][1] + 1: # пересечение ИЛИ смежность ([4М,5М)+[5М,6М))
|
||||
if merged and lo <= merged[-1][1] + 1: # пересечение ИЛИ смежность ([4М,5М]+[5М+1,6М])
|
||||
prev_lo, prev_hi = merged[-1]
|
||||
merged[-1] = (prev_lo, max(prev_hi, hi))
|
||||
else:
|
||||
|
|
|
|||
|
|
@ -1203,7 +1203,12 @@ class YandexRealtyScraper(BaseScraper):
|
|||
page += 1
|
||||
pages_fetched += 1
|
||||
if on_bucket is not None:
|
||||
on_bucket(bucket_key, len(seen))
|
||||
# Полноты не знаем: probe провалился, пагинация оборвана пустой
|
||||
# страницей ИЛИ потолком max_pages_per_bucket. complete=False →
|
||||
# бакет НЕ пишется в done-леджер (как у cian), иначе его интервал
|
||||
# слился бы с соседними в containment-гейте (#3359) и резюм уже не
|
||||
# переобошёл бы недобранную полосу.
|
||||
on_bucket(bucket_key, len(seen), False)
|
||||
|
||||
async def _leaf(plo: int | None, phi: int | None, result: ProbeResult) -> None:
|
||||
total = result.count
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue