Compare commits
No commits in common. "993b44c7e0084076dedcfe1b55d35ce6c112878f" and "929eff11d7db02a89822b60820340d416aedf378" have entirely different histories.
993b44c7e0
...
929eff11d7
5 changed files with 7 additions and 386 deletions
|
|
@ -1,63 +0,0 @@
|
||||||
"""#3074: чекпоинт предшественника наследуется в counters НОВОГО прогона при claim.
|
|
||||||
|
|
||||||
Прод-факт (Beget, 23.08): прогон 4707 подхватил чекпоинт 4117 (42 корзины,
|
|
||||||
resume_reason='ok'), был убит деплоем на 26-й минуте — ДО завершения первой
|
|
||||||
НОВОЙ корзины — и не успел ни разу написать heartbeat с done_buckets. Его
|
|
||||||
собственный чекпоинт остался пуст: следующий кандидат увидел бы no_checkpoint,
|
|
||||||
и 42 корзины пропали бы, хотя цепочка (resume_chain=2) ещё позволяла подхват.
|
|
||||||
|
|
||||||
Фикс: _resume_decision при вердикте 'ok' кладёт done_buckets предшественника в
|
|
||||||
counters-заготовку нового прогона — она пишется в БД прямо при claim
|
|
||||||
(_pick_resume → update_heartbeat), до старта пайплайна. Heartbeat мержит jsonb,
|
|
||||||
первый настоящий bucket-heartbeat перезапишет ключ надмножеством.
|
|
||||||
|
|
||||||
Красный на origin/main по ЗНАЧЕНИЮ: 'done_buckets' отсутствует в вердикте.
|
|
||||||
"""
|
|
||||||
|
|
||||||
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 scraper_kit.orchestration import scheduler as sched
|
|
||||||
|
|
||||||
|
|
||||||
def _candidate(**over: Any) -> SimpleNamespace:
|
|
||||||
"""Кандидат из _RESUME_CANDIDATE_SQL: прогон 4117 как он лежал на проде 23.08."""
|
|
||||||
base = {
|
|
||||||
"prev_id": 4117,
|
|
||||||
"prev_status": "banned",
|
|
||||||
"prev_counters": {
|
|
||||||
"resume_chain": 1,
|
|
||||||
"done_buckets": [f"room_1_komn:{i}:0" for i in range(42)],
|
|
||||||
},
|
|
||||||
"same_params": True,
|
|
||||||
"age_h": 164.8,
|
|
||||||
"interval_days": "7",
|
|
||||||
}
|
|
||||||
base.update(over)
|
|
||||||
return SimpleNamespace(**base)
|
|
||||||
|
|
||||||
|
|
||||||
def test_ok_verdict_carries_predecessor_checkpoint() -> None:
|
|
||||||
"""Вердикт 'ok' несёт done_buckets предшественника — чекпоинт персистится при
|
|
||||||
claim и переживает обрыв до первой новой корзины (кейс 4707)."""
|
|
||||||
resume_id, verdict = sched._resume_decision(_candidate())
|
|
||||||
assert resume_id == 4117
|
|
||||||
assert verdict["resume_reason"] == "ok"
|
|
||||||
assert verdict.get("done_buckets") == sorted(f"room_1_komn:{i}:0" for i in range(42)), (
|
|
||||||
"counters-заготовка нового прогона не содержит чекпоинт предшественника — "
|
|
||||||
"обрыв до первой новой корзины снова потеряет всю цепочку"
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
def test_refusal_verdict_carries_no_checkpoint() -> None:
|
|
||||||
"""Отказ от подхвата чекпоинт не наследует: прогон честно идёт с нуля, и его
|
|
||||||
counters не должны врать, будто корзины предшественника уже собраны."""
|
|
||||||
_id, verdict = sched._resume_decision(_candidate(same_params=False))
|
|
||||||
assert verdict["resume_reason"] == "params_changed"
|
|
||||||
assert "done_buckets" not in verdict
|
|
||||||
|
|
@ -1,241 +0,0 @@
|
||||||
"""#3074: чекпоинты для yandex_city_sweep — combo как единица возобновления.
|
|
||||||
|
|
||||||
Из таблицы убитых деплоем прогонов (#3074): yandex_city_sweep 15.08 прожил 65 мин,
|
|
||||||
12.08 — 2 ч 33 мин; оба потеряны целиком, потому что у свипа не было чекпоинтов
|
|
||||||
вовсе (avito/cian full-load обзавелись ими в #930/#2845).
|
|
||||||
|
|
||||||
Единица чекпоинта — combo (сегмент × комнатность × ценовой диапазон): ровно то,
|
|
||||||
чем цикл обхода уже итерируется, метка та же, что у бакетов avito
|
|
||||||
("vtorichka/room_1:0-3000000"). Три слоя:
|
|
||||||
провайдер — skip_combos (ни одного HTTP по собранным) + on_combo для каждого
|
|
||||||
ПРОЙДЕННОГО combo, включая пустые (иначе пустой combo не попадал
|
|
||||||
бы в чекпоинт и перечитывался бы вечно);
|
|
||||||
пайплайн — done_combos → heartbeat c done_buckets (мерж jsonb);
|
|
||||||
планировщик — resume_run_id=_pick_resume(...) в диспатче.
|
|
||||||
|
|
||||||
Красные на origin/main: планировщик передаёт None литералом (по значению),
|
|
||||||
on_combo не вызывается для пустых combo (по значению).
|
|
||||||
"""
|
|
||||||
|
|
||||||
from __future__ import annotations
|
|
||||||
|
|
||||||
import os
|
|
||||||
|
|
||||||
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
|
|
||||||
from scraper_kit.orchestration import scheduler as sched
|
|
||||||
from scraper_kit.providers.yandex.serp import YandexRealtyScraper
|
|
||||||
|
|
||||||
# ── провайдер: заглушка gate-API ─────────────────────────────────────────────
|
|
||||||
|
|
||||||
|
|
||||||
def _gate_payload(entities: list[dict[str, Any]], total_pages: int = 1) -> str:
|
|
||||||
offers = {"entities": entities, "pager": {"totalPages": total_pages}}
|
|
||||||
return json.dumps({"response": {"search": {"offers": offers}}})
|
|
||||||
|
|
||||||
|
|
||||||
class _RecordingBrowser:
|
|
||||||
"""BrowserFetcher-заглушка: отвечает одним и тем же телом, считает вызовы."""
|
|
||||||
|
|
||||||
def __init__(self, body: str) -> None:
|
|
||||||
self.body = body
|
|
||||||
self.urls: list[str] = []
|
|
||||||
|
|
||||||
async def fetch(self, url: str) -> str:
|
|
||||||
self.urls.append(url)
|
|
||||||
return self.body
|
|
||||||
|
|
||||||
|
|
||||||
def _scraper(body: str) -> YandexRealtyScraper:
|
|
||||||
s = YandexRealtyScraper(types.SimpleNamespace())
|
|
||||||
s._browser = _RecordingBrowser(body) # type: ignore[assignment]
|
|
||||||
s.sleep_between_requests = _no_sleep # type: ignore[method-assign]
|
|
||||||
return s
|
|
||||||
|
|
||||||
|
|
||||||
async def _no_sleep() -> None:
|
|
||||||
pass
|
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.asyncio
|
|
||||||
async def test_on_combo_fires_for_complete_empty_combo() -> None:
|
|
||||||
"""Пройденный до конца combo с ПУСТОЙ выдачей обязан дойти до on_combo —
|
|
||||||
иначе он не попадёт в чекпоинт и будет перечитываться каждым продолжением.
|
|
||||||
На origin/main on_combo вызывался только при непустых new_lots."""
|
|
||||||
scraper = _scraper(_gate_payload([]))
|
|
||||||
seen_labels: list[str] = []
|
|
||||||
|
|
||||||
await scraper.fetch_around_multi_room(
|
|
||||||
56.84,
|
|
||||||
60.60,
|
|
||||||
1000,
|
|
||||||
max_pages=1,
|
|
||||||
rooms_list=["room_1"],
|
|
||||||
price_ranges=[(None, 3_000_000)],
|
|
||||||
segments=["NO"],
|
|
||||||
on_combo=lambda label, lots: seen_labels.append(label),
|
|
||||||
)
|
|
||||||
|
|
||||||
assert seen_labels == ["vtorichka/room_1:None-3000000"], (
|
|
||||||
"пустой, но полностью пройденный combo не дошёл до on_combo — "
|
|
||||||
"в чекпоинт он не попадёт никогда"
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.asyncio
|
|
||||||
async def test_skip_combos_makes_zero_http_requests() -> None:
|
|
||||||
"""Combo из чекпоинта не порождает ни одного HTTP-запроса и не зовёт on_combo."""
|
|
||||||
scraper = _scraper(_gate_payload([]))
|
|
||||||
seen_labels: list[str] = []
|
|
||||||
|
|
||||||
await scraper.fetch_around_multi_room(
|
|
||||||
56.84,
|
|
||||||
60.60,
|
|
||||||
1000,
|
|
||||||
max_pages=1,
|
|
||||||
rooms_list=["room_1", "room_2"],
|
|
||||||
price_ranges=[(None, 3_000_000)],
|
|
||||||
segments=["NO"],
|
|
||||||
on_combo=lambda label, lots: seen_labels.append(label),
|
|
||||||
skip_combos={"vtorichka/room_1:None-3000000"},
|
|
||||||
)
|
|
||||||
|
|
||||||
browser: _RecordingBrowser = scraper._browser # type: ignore[assignment]
|
|
||||||
assert seen_labels == ["vtorichka/room_2:None-3000000"]
|
|
||||||
assert len(browser.urls) == 1, f"скипнутый combo всё равно ходил в сеть: {browser.urls}"
|
|
||||||
|
|
||||||
|
|
||||||
# ── планировщик: диспатч отдаёт точку в пайплайн (зеркало test_930) ─────────
|
|
||||||
|
|
||||||
|
|
||||||
def _candidate(**over: Any) -> types.SimpleNamespace:
|
|
||||||
base = {
|
|
||||||
"prev_id": 4117,
|
|
||||||
"prev_status": "banned",
|
|
||||||
"prev_counters": {"done_buckets": ["vtorichka/room_1:None-3000000"]},
|
|
||||||
"same_params": True,
|
|
||||||
"age_h": 20.0,
|
|
||||||
"interval_days": "1",
|
|
||||||
}
|
|
||||||
base.update(over)
|
|
||||||
return types.SimpleNamespace(**base)
|
|
||||||
|
|
||||||
|
|
||||||
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_yandex_sweep() -> None:
|
|
||||||
"""_job_yandex_city_sweep передаёт resume_run_id, а не None литералом.
|
|
||||||
Красный на origin/main по значению (ключа в kwargs нет → None != 4117)."""
|
|
||||||
db = _FakeDb(_candidate())
|
|
||||||
captured: dict[str, Any] = {}
|
|
||||||
|
|
||||||
async def _spy(*_a: Any, **kw: Any) -> None:
|
|
||||||
captured.update(kw)
|
|
||||||
|
|
||||||
with patch.object(sched, "run_yandex_city_sweep", _spy):
|
|
||||||
await sched._job_yandex_city_sweep(db, 5000, {}, MagicMock())
|
|
||||||
|
|
||||||
assert captured.get("resume_run_id") == 4117, (
|
|
||||||
"планировщик не отдал чекпоинт яндекс-свипу — прогон пойдёт с нуля"
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
# ── пайплайн: проводка skip→scraper и done_buckets→heartbeat ────────────────
|
|
||||||
|
|
||||||
|
|
||||||
class _SweepFakeDb:
|
|
||||||
"""Резюм-SELECT отдаёт counters предшественника; UPDATE'ы записываются."""
|
|
||||||
|
|
||||||
def __init__(self, prev_counters: dict[str, Any]) -> None:
|
|
||||||
self.prev_counters = prev_counters
|
|
||||||
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() # is_cancelled/сторожа — best-effort, глотаем
|
|
||||||
|
|
||||||
def commit(self) -> None:
|
|
||||||
pass
|
|
||||||
|
|
||||||
def rollback(self) -> None:
|
|
||||||
pass
|
|
||||||
|
|
||||||
|
|
||||||
class _FakeSweepScraper:
|
|
||||||
"""Двойник YandexRealtyScraper для пайплайна: фиксирует kwargs fetch'а,
|
|
||||||
отдаёт один пройденный пустой combo через on_combo."""
|
|
||||||
|
|
||||||
captured: dict[str, Any] = {} # noqa: RUF012 — тестовый сборник kwargs
|
|
||||||
gate_fetch_attempts = 1
|
|
||||||
gate_fetch_failures = 0
|
|
||||||
|
|
||||||
def __init__(self, *_a: Any, **_kw: Any) -> None:
|
|
||||||
pass
|
|
||||||
|
|
||||||
async def __aenter__(self) -> _FakeSweepScraper:
|
|
||||||
return self
|
|
||||||
|
|
||||||
async def __aexit__(self, *_exc: Any) -> None:
|
|
||||||
return None
|
|
||||||
|
|
||||||
async def fetch_around_multi_room(self, *_a: Any, **kw: Any) -> list:
|
|
||||||
_FakeSweepScraper.captured = dict(kw)
|
|
||||||
on_combo = kw.get("on_combo")
|
|
||||||
if on_combo is not None:
|
|
||||||
on_combo("vtorichka/room_2:None-3000000", [])
|
|
||||||
return []
|
|
||||||
|
|
||||||
|
|
||||||
async def test_pipeline_resumes_and_checkpoints() -> None:
|
|
||||||
"""run_yandex_city_sweep: чекпоинт предшественника уезжает в scraper как
|
|
||||||
skip_combos, пройденный combo дописывается в done_buckets heartbeat'а."""
|
|
||||||
from scraper_kit.orchestration import pipeline as pl
|
|
||||||
|
|
||||||
prev = {"done_buckets": ["vtorichka/room_1:None-3000000"]}
|
|
||||||
db = _SweepFakeDb(prev)
|
|
||||||
|
|
||||||
with (
|
|
||||||
patch.object(pl, "YandexRealtyScraper", _FakeSweepScraper),
|
|
||||||
patch.object(pl.runs, "is_cancelled", lambda *_a: False),
|
|
||||||
):
|
|
||||||
await pl.run_yandex_city_sweep(
|
|
||||||
db, # type: ignore[arg-type]
|
|
||||||
run_id=5001,
|
|
||||||
config=types.SimpleNamespace(scraper_proxy_url=None),
|
|
||||||
matcher=MagicMock(),
|
|
||||||
enrichment=MagicMock(),
|
|
||||||
enrich_address=False,
|
|
||||||
resume_run_id=4999,
|
|
||||||
)
|
|
||||||
|
|
||||||
assert _FakeSweepScraper.captured.get("skip_combos") == {"vtorichka/room_1:None-3000000"}, (
|
|
||||||
"чекпоинт предшественника не доехал до scraper'а"
|
|
||||||
)
|
|
||||||
|
|
||||||
with_ckpt = [hb for hb in db.heartbeats if "done_buckets" in hb]
|
|
||||||
assert with_ckpt, "ни один heartbeat не унёс done_buckets — чекпоинт не персистится"
|
|
||||||
assert with_ckpt[-1]["done_buckets"] == [
|
|
||||||
"vtorichka/room_1:None-3000000",
|
|
||||||
"vtorichka/room_2:None-3000000",
|
|
||||||
], "done_buckets не аккумулирует пройденный combo поверх унаследованных"
|
|
||||||
|
|
@ -2078,7 +2078,6 @@ async def run_yandex_city_sweep(
|
||||||
price_ranges: list[tuple[int | None, int | None]] | None = None,
|
price_ranges: list[tuple[int | None, int | None]] | None = None,
|
||||||
segments: list[str] | None = None,
|
segments: list[str] | None = None,
|
||||||
region_code: int = DEFAULT_REGION_CODE,
|
region_code: int = DEFAULT_REGION_CODE,
|
||||||
resume_run_id: int | None = None,
|
|
||||||
) -> YandexCitySweepCounters:
|
) -> YandexCitySweepCounters:
|
||||||
"""Yandex.Недвижимость city sweep: rooms × price combos от центра ЕКБ → save → address-enrich.
|
"""Yandex.Недвижимость city sweep: rooms × price combos от центра ЕКБ → save → address-enrich.
|
||||||
|
|
||||||
|
|
@ -2170,39 +2169,6 @@ async def run_yandex_city_sweep(
|
||||||
_proxy_url = config.scraper_proxy_url
|
_proxy_url = config.scraper_proxy_url
|
||||||
_proxies = {"http": _proxy_url, "https": _proxy_url} if _proxy_url else None
|
_proxies = {"http": _proxy_url, "https": _proxy_url} if _proxy_url else None
|
||||||
|
|
||||||
# ── Checkpoint/resume (#3074): combo — естественная единица обхода ────────
|
|
||||||
# Ключ чекпоинта — combo_label ("сегмент/комнатность:lo:hi"), якоря в нём нет,
|
|
||||||
# поэтому подхват включается ТОЛЬКО при единственном якоре (штатный прод-режим:
|
|
||||||
# и ЕКБ-центр, и city-свипы — ровно один anchor). Multi-anchor — легаси/ручной
|
|
||||||
# режим, у него один combo_label повторяется на каждом якоре и пропуск был бы
|
|
||||||
# пропуском ЧУЖИХ якорей.
|
|
||||||
skip_combos: set[str] = set()
|
|
||||||
if resume_run_id is not None and len(_anchors) == 1:
|
|
||||||
_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_combos = set(_prev_counters.get("done_buckets", []))
|
|
||||||
logger.info(
|
|
||||||
"yandex-sweep run_id=%d: resuming from run %s — %d combos already done",
|
|
||||||
run_id,
|
|
||||||
resume_run_id,
|
|
||||||
len(skip_combos),
|
|
||||||
)
|
|
||||||
elif resume_run_id is not None:
|
|
||||||
logger.warning(
|
|
||||||
"yandex-sweep run_id=%d: resume от run %s ОТКЛОНЁН — %d якорей (>1), "
|
|
||||||
"combo-чекпоинт применим только к единственному якорю",
|
|
||||||
run_id,
|
|
||||||
resume_run_id,
|
|
||||||
len(_anchors),
|
|
||||||
)
|
|
||||||
done_combos: set[str] = set(skip_combos)
|
|
||||||
|
|
||||||
try:
|
try:
|
||||||
for idx, (lat, lon, name) in enumerate(_anchors, start=1):
|
for idx, (lat, lon, name) in enumerate(_anchors, start=1):
|
||||||
# ── Cooperative cancel ───────────────────────────────────────────
|
# ── Cooperative cancel ───────────────────────────────────────────
|
||||||
|
|
@ -2265,21 +2231,7 @@ async def run_yandex_city_sweep(
|
||||||
new_lots: list[ScrapedLot],
|
new_lots: list[ScrapedLot],
|
||||||
_accumulator: list[ScrapedLot] = _al,
|
_accumulator: list[ScrapedLot] = _al,
|
||||||
) -> None:
|
) -> None:
|
||||||
"""Callback: после каждого ПРОЙДЕННОГО combo (и пустого — #3074).
|
"""Callback: вызывается scraper'ом после каждого combo с новыми лотами."""
|
||||||
|
|
||||||
Вызов означает «combo пройден до конца» — фиксируем его в
|
|
||||||
чекпоинт done_combos и пишем в heartbeat (мерж jsonb: другие
|
|
||||||
heartbeat'ы без ключа его не затирают). Пустой new_lots не
|
|
||||||
гоняет save_listings.
|
|
||||||
"""
|
|
||||||
done_combos.add(combo_label)
|
|
||||||
if not new_lots:
|
|
||||||
runs.update_heartbeat(
|
|
||||||
db,
|
|
||||||
run_id,
|
|
||||||
{**counters.to_dict(), "done_buckets": sorted(done_combos)},
|
|
||||||
)
|
|
||||||
return
|
|
||||||
_accumulator.extend(new_lots)
|
_accumulator.extend(new_lots)
|
||||||
counters.lots_fetched += len(new_lots)
|
counters.lots_fetched += len(new_lots)
|
||||||
try:
|
try:
|
||||||
|
|
@ -2307,9 +2259,7 @@ async def run_yandex_city_sweep(
|
||||||
db.rollback()
|
db.rollback()
|
||||||
except Exception:
|
except Exception:
|
||||||
pass
|
pass
|
||||||
runs.update_heartbeat(
|
runs.update_heartbeat(db, run_id, counters.to_dict())
|
||||||
db, run_id, {**counters.to_dict(), "done_buckets": sorted(done_combos)}
|
|
||||||
)
|
|
||||||
logger.debug(
|
logger.debug(
|
||||||
"yandex-sweep run_id=%d combo %s: saved %d lots "
|
"yandex-sweep run_id=%d combo %s: saved %d lots "
|
||||||
"(ins=%d upd=%d total_fetched=%d)",
|
"(ins=%d upd=%d total_fetched=%d)",
|
||||||
|
|
@ -2335,7 +2285,6 @@ async def run_yandex_city_sweep(
|
||||||
price_ranges=_price_ranges,
|
price_ranges=_price_ranges,
|
||||||
segments=_segments,
|
segments=_segments,
|
||||||
on_combo=_on_combo,
|
on_combo=_on_combo,
|
||||||
skip_combos=skip_combos or None,
|
|
||||||
)
|
)
|
||||||
# #2625: аккумулируем run-level gate-API attempts/failures.
|
# #2625: аккумулируем run-level gate-API attempts/failures.
|
||||||
yandex_gate_attempts += scraper.gate_fetch_attempts
|
yandex_gate_attempts += scraper.gate_fetch_attempts
|
||||||
|
|
|
||||||
|
|
@ -599,17 +599,6 @@ def _resume_decision(row: Any) -> tuple[int | None, dict[str, Any]]:
|
||||||
verdict["resume_from"] = int(row.prev_id)
|
verdict["resume_from"] = int(row.prev_id)
|
||||||
verdict["resume_reason"] = "ok"
|
verdict["resume_reason"] = "ok"
|
||||||
verdict["resume_chain"] = prev_chain + 1
|
verdict["resume_chain"] = prev_chain + 1
|
||||||
# Наследуем чекпоинт в counters нового прогона ПРЯМО ПРИ CLAIM (#3074).
|
|
||||||
# Прод-факт (run 4707, 23.08): возобновлённый прогон, убитый деплоем на
|
|
||||||
# 26-й минуте — ДО завершения первой новой корзины, — не успел ни разу
|
|
||||||
# написать heartbeat с done_buckets. Его собственный чекпоинт остался
|
|
||||||
# пуст, кандидат следующей субботы увидел no_checkpoint, и 42 корзины
|
|
||||||
# предшественника (4117) пропали. Пайплайн сеет `done = set(skip_set)`
|
|
||||||
# у себя в памяти, но до первого _on_bucket это знание нигде не
|
|
||||||
# персистится; heartbeat мержит jsonb (`counters || :counters`), так что
|
|
||||||
# первый же настоящий bucket-heartbeat перезапишет ключ тем же множеством
|
|
||||||
# плюс новое — двойной записи не возникает.
|
|
||||||
verdict["done_buckets"] = sorted(done_buckets) if done_buckets else []
|
|
||||||
return int(row.prev_id), verdict
|
return int(row.prev_id), verdict
|
||||||
return None, verdict
|
return None, verdict
|
||||||
|
|
||||||
|
|
@ -891,7 +880,6 @@ async def _job_yandex_city_sweep(
|
||||||
radius_m=int(params.get("radius_m", 1500)),
|
radius_m=int(params.get("radius_m", 1500)),
|
||||||
enrich_address=bool(params.get("enrich_address", True)),
|
enrich_address=bool(params.get("enrich_address", True)),
|
||||||
segments=params.get("segments"),
|
segments=params.get("segments"),
|
||||||
resume_run_id=_pick_resume(db, run_id),
|
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -848,7 +848,6 @@ class YandexRealtyScraper(BaseScraper):
|
||||||
price_ranges: list[tuple[int | None, int | None]] | None = None,
|
price_ranges: list[tuple[int | None, int | None]] | None = None,
|
||||||
segments: list[str] | None = None,
|
segments: list[str] | None = None,
|
||||||
on_combo: Any = None,
|
on_combo: Any = None,
|
||||||
skip_combos: Any = None,
|
|
||||||
**_legacy_kwargs: Any,
|
**_legacy_kwargs: Any,
|
||||||
) -> list[ScrapedLot]:
|
) -> list[ScrapedLot]:
|
||||||
"""Fetch via segment x rooms x price-range combos; paginate each combo to totalPages.
|
"""Fetch via segment x rooms x price-range combos; paginate each combo to totalPages.
|
||||||
|
|
@ -861,16 +860,9 @@ class YandexRealtyScraper(BaseScraper):
|
||||||
["NO", "YES"] to sweep both vtorichka and novostroyki in one call.
|
["NO", "YES"] to sweep both vtorichka and novostroyki in one call.
|
||||||
|
|
||||||
on_combo: опциональный callback(combo_label: str, new_lots: list[ScrapedLot]) -> None.
|
on_combo: опциональный callback(combo_label: str, new_lots: list[ScrapedLot]) -> None.
|
||||||
Вызывается после каждого НЕ скипнутого combo — в том числе с пустым
|
Вызывается после каждого combo с новыми (не повторными) лотами из него.
|
||||||
new_lots (#3074: полностью пройденный combo без новых лотов обязан
|
Позволяет инкрементальный save: вызывающий код сохраняет new_lots сразу,
|
||||||
попасть в чекпоинт, иначе его перечитывали бы вечно). Позволяет
|
|
||||||
инкрементальный save: вызывающий код сохраняет new_lots сразу,
|
|
||||||
не дожидаясь конца sweep'а. Дубли между combo НЕ передаются повторно.
|
не дожидаясь конца sweep'а. Дубли между combo НЕ передаются повторно.
|
||||||
Combo, оборванный отказом (gate/extraction failure), on_combo НЕ
|
|
||||||
вызывает — вызов означает «combo пройден до конца».
|
|
||||||
skip_combos: опциональная коллекция combo_label, которые уже собраны
|
|
||||||
предыдущим оборванным прогоном (#3074, чекпоинт done_buckets) — по
|
|
||||||
ним не делается ни одного HTTP-запроса и on_combo не вызывается.
|
|
||||||
"""
|
"""
|
||||||
seen: dict[str, ScrapedLot] = {}
|
seen: dict[str, ScrapedLot] = {}
|
||||||
_segments = segments or ["NO"]
|
_segments = segments or ["NO"]
|
||||||
|
|
@ -887,9 +879,6 @@ class YandexRealtyScraper(BaseScraper):
|
||||||
_seg = _segment_for(new_flat)
|
_seg = _segment_for(new_flat)
|
||||||
for rooms, price_min, price_max in combos:
|
for rooms, price_min, price_max in combos:
|
||||||
combo_label = f"{_seg}/{_combo_label(rooms, price_min, price_max)}"
|
combo_label = f"{_seg}/{_combo_label(rooms, price_min, price_max)}"
|
||||||
if skip_combos and combo_label in skip_combos:
|
|
||||||
# #3074: собрано предыдущим оборванным прогоном — ни одного запроса.
|
|
||||||
continue
|
|
||||||
total_pages: int | None = None
|
total_pages: int | None = None
|
||||||
combo_new_lots: list[ScrapedLot] = []
|
combo_new_lots: list[ScrapedLot] = []
|
||||||
combo_skipped = False
|
combo_skipped = False
|
||||||
|
|
@ -1023,10 +1012,9 @@ class YandexRealtyScraper(BaseScraper):
|
||||||
seen[key] = lot
|
seen[key] = lot
|
||||||
combo_new_lots.append(lot)
|
combo_new_lots.append(lot)
|
||||||
|
|
||||||
# Инкрементальный save + чекпоинт: on_combo для КАЖДОГО пройденного
|
# Инкрементальный save: вызываем on_combo если есть новые лоты.
|
||||||
# combo, включая пустые (#3074) — вызов означает «combo пройден до
|
# Скипнутые combo (combo_skipped=True) не вызывают on_combo.
|
||||||
# конца». Оборванные отказом (combo_skipped=True) не вызывают.
|
if on_combo is not None and combo_new_lots and not combo_skipped:
|
||||||
if on_combo is not None and not combo_skipped:
|
|
||||||
try:
|
try:
|
||||||
on_combo(combo_label, combo_new_lots)
|
on_combo(combo_label, combo_new_lots)
|
||||||
except Exception:
|
except Exception:
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue