feat(tradein/scheduler): планировщик подхватывает чекпоинт оборванного прогона (#930) #2845
6 changed files with 588 additions and 32 deletions
|
|
@ -0,0 +1,295 @@
|
|||
"""#930 добивка: планировщик не подхватывал чекпоинт оборванного прогона.
|
||||
|
||||
#930 сделал обе половины механизма — запись точки (`counters.done_buckets`, per-bucket
|
||||
heartbeat) и её чтение (`run_*_full_load(resume_run_id=...)`, skip-set в SERP-слое), —
|
||||
но единственным входом оставил админку. У avito full-load админского эндпоинта нет
|
||||
вовсе, а планировщик передавал `resume_run_id=None` ЛИТЕРАЛОМ (scheduler.py 708/728/799
|
||||
на origin/main). То есть боевой путь возобновления не существовал ни одного дня.
|
||||
|
||||
Цена на проде (замер 2026-08-12, 90 суток, read-only): 433 корзины в 30 оборванных
|
||||
прогонах с ЖИВОЙ незабранной точкой — avito_full_load 242, cian_full_load 134,
|
||||
avito_full_load_exhaustive 57. Прогон 3547 (09.08, убит деплоем на третьем часу, 35 из
|
||||
84 корзин дерева) лежит до сих пор и будет подхвачен расписанием 139 16.08.
|
||||
|
||||
Красный прогон на origin/main:
|
||||
1. `test_scheduler_hands_checkpoint_to_pipeline` — планировщик отдаёт в пайплайн
|
||||
resume_run_id=None вместо id прошлого прогона (AssertionError на 3 источниках);
|
||||
2. `test_partial_bucket_is_not_complete` — бакет с выпавшей страницей приезжает в
|
||||
on_bucket неотличимым от целого (у колбэка нет аргумента полноты вообще);
|
||||
3. `test_pipeline_keeps_partial_bucket_out_of_checkpoint` — TypeError: `_on_bucket`
|
||||
на main принимает два аргумента, признаку полноты некуда приехать.
|
||||
Тесты ладдера (`_resume_decision`) на main падают с AttributeError — функции нет.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
from types import SimpleNamespace
|
||||
from typing import Any
|
||||
from unittest.mock import AsyncMock, MagicMock, patch
|
||||
|
||||
import pytest
|
||||
|
||||
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db")
|
||||
|
||||
from scraper_kit.orchestration import scheduler as sched
|
||||
from scraper_kit.orchestration.pipeline import run_avito_full_load
|
||||
from scraper_kit.providers.avito.serp import AvitoScraper
|
||||
|
||||
PFX = "scraper_kit.orchestration.pipeline"
|
||||
|
||||
# Прод-слепок расписания 139 (avito_full_load_exhaustive) на 2026-08-12: params прогона
|
||||
# 3547 совпадают с default_params расписания байт-в-байт — это и есть «то же задание».
|
||||
_PARAMS = {
|
||||
"concurrency": 1,
|
||||
"interval_days": 7,
|
||||
"secondary_only": True,
|
||||
"request_delay_sec": 7.0,
|
||||
"price_cap_per_bucket": 1400,
|
||||
}
|
||||
|
||||
|
||||
def _candidate(**over: Any) -> SimpleNamespace:
|
||||
"""Строка-кандидат из _RESUME_CANDIDATE_SQL: прогон 3547 как он лежит на проде."""
|
||||
base = {
|
||||
"prev_id": 3547,
|
||||
"prev_status": "cancelled",
|
||||
"prev_counters": {
|
||||
"unique_fetched": 5496,
|
||||
"done_buckets": [f"room_1_komn:{i}:0" for i in range(35)],
|
||||
},
|
||||
"same_params": True,
|
||||
"age_h": 164.6, # 6.86 суток — столько будет точке 3547 к подхвату 16.08
|
||||
"interval_days": "7",
|
||||
}
|
||||
base.update(over)
|
||||
return SimpleNamespace(**base)
|
||||
|
||||
|
||||
class _FakeDb:
|
||||
"""Двойник сессии: отдаёт ОДНУ строку-кандидата на любой SELECT, глотает UPDATE."""
|
||||
|
||||
def __init__(self, row: Any) -> None:
|
||||
self.row = row
|
||||
self.written: list[dict[str, Any]] = []
|
||||
|
||||
def execute(self, _stmt: Any, params: dict[str, Any] | None = None) -> Any:
|
||||
if params and "counters" in params: # update_heartbeat пишет вердикт
|
||||
self.written.append(params)
|
||||
return MagicMock()
|
||||
return MagicMock(fetchone=lambda: self.row)
|
||||
|
||||
def commit(self) -> None:
|
||||
pass
|
||||
|
||||
|
||||
# ── 1. Главное: планировщик обязан отдать точку в пайплайн ───────────────────
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("job", "pipeline_fn"),
|
||||
[
|
||||
(sched._job_avito_full_load, "run_avito_full_load"),
|
||||
(sched._job_avito_full_load_exhaustive, "run_avito_full_load"),
|
||||
(sched._job_cian_full_load, "run_cian_full_load"),
|
||||
],
|
||||
)
|
||||
async def test_scheduler_hands_checkpoint_to_pipeline(job: Any, pipeline_fn: str) -> None:
|
||||
"""Оборванный прогон с валидной точкой → новый прогон продолжает его, а не с нуля.
|
||||
|
||||
Падает на origin/main: планировщик передаёт литеральный None — 433 корзины за 90
|
||||
суток перебирались заново, включая 35 корзин прогона 3547.
|
||||
"""
|
||||
db = _FakeDb(_candidate())
|
||||
captured: dict[str, Any] = {}
|
||||
|
||||
async def _spy(*_a: Any, **kw: Any) -> None:
|
||||
captured.update(kw)
|
||||
|
||||
with patch.object(sched, pipeline_fn, _spy):
|
||||
await job(db, 4000, dict(_PARAMS), MagicMock())
|
||||
|
||||
assert captured["resume_run_id"] == 3547
|
||||
|
||||
|
||||
async def test_verdict_lands_in_counters_of_new_run() -> None:
|
||||
"""Подхватили или нет — видно В СЧЁТЧИКАХ прогона, а не только в docker-логах.
|
||||
|
||||
Логи теряются при редеплое (контейнер tradein-scraper пересоздаётся), поэтому
|
||||
молчаливый отказ подхватить неотличим от отсутствия правки.
|
||||
"""
|
||||
db = _FakeDb(_candidate(prev_status="zombie"))
|
||||
with patch.object(sched, "run_avito_full_load", AsyncMock()):
|
||||
await sched._job_avito_full_load(db, 4000, dict(_PARAMS), MagicMock())
|
||||
|
||||
assert db.written, "вердикт о подхвате не записан в counters нового прогона"
|
||||
written = db.written[-1]["counters"]
|
||||
assert '"resume_reason": "status_zombie"' in written
|
||||
assert '"resume_candidate": 3547' in written
|
||||
|
||||
|
||||
# ── 2. Ладдер отказов: у каждого нуля своя причина ───────────────────────────
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("row", "reason"),
|
||||
[
|
||||
(None, "no_prev_run"),
|
||||
(_candidate(prev_status="done"), "status_done"),
|
||||
(_candidate(prev_status="zombie"), "status_zombie"),
|
||||
(_candidate(same_params=False), "params_changed"),
|
||||
(_candidate(prev_counters={"unique_fetched": 2977}), "no_checkpoint"),
|
||||
(_candidate(age_h=200.0), "checkpoint_stale"),
|
||||
(_candidate(prev_counters={"done_buckets": ["a"], "resume_chain": 2}), "chain_limit"),
|
||||
],
|
||||
)
|
||||
def test_resume_refusals_are_named(row: Any, reason: str) -> None:
|
||||
"""«Не подхватили» — это семь РАЗНЫХ фактов, и в counters они различимы."""
|
||||
resume_id, verdict = sched._resume_decision(row)
|
||||
assert resume_id is None
|
||||
assert verdict["resume_reason"] == reason
|
||||
assert verdict["resume_from"] is None
|
||||
|
||||
|
||||
def test_resume_chain_is_bounded() -> None:
|
||||
"""Цепочка возобновлений считается и упирается в потолок, а не тянется вечно.
|
||||
|
||||
Потолок выведен из STALE_DIGEST_INTERVAL_FACTOR (см. scheduler.py): полный обход
|
||||
обязан начаться раньше, чем сводка объявит источник просроченным.
|
||||
"""
|
||||
assert sched._MAX_RESUME_CHAIN == sched.STALE_DIGEST_INTERVAL_FACTOR - 1
|
||||
_id, first = sched._resume_decision(_candidate())
|
||||
assert first["resume_chain"] == 1
|
||||
_id2, second = sched._resume_decision(
|
||||
_candidate(prev_counters={"done_buckets": ["a"], "resume_chain": 1})
|
||||
)
|
||||
assert second["resume_chain"] == sched._MAX_RESUME_CHAIN
|
||||
third_id, third = sched._resume_decision(
|
||||
_candidate(prev_counters={"done_buckets": ["a"], "resume_chain": 2})
|
||||
)
|
||||
assert third_id is None and third["resume_reason"] == "chain_limit"
|
||||
|
||||
|
||||
def test_stale_threshold_follows_the_source_tick() -> None:
|
||||
"""Срок годности точки считается от такта ИСТОЧНИКА, а не общей константой.
|
||||
|
||||
cian ходит раз в 3 суток, avito — раз в 7; одна и та же точка возрастом 100 ч для
|
||||
первого просрочена, для второго свежая. Плюс сутки — сетка запуска (см.
|
||||
_resume_decision): 164.6 ч прогона 3547 при такте 7 суток обязаны пройти, иначе
|
||||
точку отвергал бы jitter расписания, а пропущенный цикл (13 суток) — нет.
|
||||
"""
|
||||
assert sched._resume_decision(_candidate(age_h=100.0, interval_days="3"))[0] is None
|
||||
assert sched._resume_decision(_candidate(age_h=100.0, interval_days="7"))[0] == 3547
|
||||
assert sched._resume_decision(_candidate(age_h=164.6, interval_days="7"))[0] == 3547
|
||||
assert sched._resume_decision(_candidate(age_h=13 * 24.0, interval_days="7"))[0] is None
|
||||
|
||||
|
||||
# ── 3. Недособранный бакет не имеет права попасть в чекпоинт ─────────────────
|
||||
|
||||
|
||||
def _serp_config() -> SimpleNamespace:
|
||||
return SimpleNamespace(
|
||||
scraper_fetch_mode="curl_cffi",
|
||||
browser_http_endpoint="http://browser.test/fetch",
|
||||
scraper_proxy_url=None,
|
||||
avito_proxy_max_rotations=0,
|
||||
avito_serp_ok_not_banned=True,
|
||||
avito_proxy_rotate_settle_s=0.0,
|
||||
proxy_rotate_attempts=1,
|
||||
proxy_rotate_attempt_timeout_s=1.0,
|
||||
scraper_skip_seen_today=False,
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("page2_html", "expected_complete"),
|
||||
[(None, False), ("<page2/>", True)],
|
||||
)
|
||||
async def test_partial_bucket_is_not_complete(
|
||||
page2_html: str | None, expected_complete: bool
|
||||
) -> None:
|
||||
"""Страница 2 из 3 выпала → бакет НЕ «сделан»; все три пришли → «сделан».
|
||||
|
||||
Контрольная половина обязательна: реализация «всегда False» тоже прошла бы
|
||||
одностороннюю проверку, но убила бы возобновление целиком.
|
||||
|
||||
Падает на origin/main: `on_bucket` вызывается двумя аргументами, признака полноты
|
||||
в протоколе нет — частичный бакет неотличим от целого и попадает в done_buckets.
|
||||
"""
|
||||
scraper = AvitoScraper(_serp_config())
|
||||
scraper.request_delay_sec = 0.0
|
||||
calls: list[tuple[str, bool]] = []
|
||||
|
||||
def _on_bucket(key: str, lots: list, complete: bool = True) -> None: # type: ignore[type-arg]
|
||||
calls.append((key, complete))
|
||||
|
||||
async def _fetch_page(_self: Any, _slug: str, page: int, *_a: Any, **_k: Any) -> str | None:
|
||||
return page2_html if page == 2 else f"<page{page}/>"
|
||||
|
||||
with (
|
||||
patch.object(AvitoScraper, "_fetch_rooms_page_html", _fetch_page),
|
||||
patch.object(
|
||||
AvitoScraper,
|
||||
"_parse_html",
|
||||
lambda _self, html, **_k: [MagicMock(source_id=html, listing_segment="secondary")],
|
||||
),
|
||||
):
|
||||
await scraper._paginate_leaf_bucket(
|
||||
room_slug="kvartiry_1_komnatnye",
|
||||
room_label="room_1_komn",
|
||||
lo=0,
|
||||
hi=3999999,
|
||||
html="<page1/>",
|
||||
max_pages=3,
|
||||
seen={},
|
||||
price_cap_per_bucket=1400,
|
||||
max_pages_per_bucket=100,
|
||||
concurrency=2,
|
||||
secondary_only=False,
|
||||
on_bucket=_on_bucket,
|
||||
skip_buckets=None,
|
||||
expected_total=3 * 50,
|
||||
)
|
||||
|
||||
assert [c[1] for c in calls] == [expected_complete]
|
||||
|
||||
|
||||
async def test_pipeline_keeps_partial_bucket_out_of_checkpoint() -> None:
|
||||
"""Пайплайн: лоты частичного бакета СОХРАНЕНЫ, но в чекпоинт он не попал.
|
||||
|
||||
Именно здесь «видимая потеря» (перескрап) не превращается в «невидимую»: пропустить
|
||||
частичный бакет на следующем прогоне значит не перечитать его страницы уже никогда.
|
||||
"""
|
||||
finals: list[dict[str, Any]] = []
|
||||
|
||||
class _Recorder:
|
||||
def is_cancelled(self, *_a: Any, **_k: Any) -> bool:
|
||||
return False
|
||||
|
||||
def update_heartbeat(self, *_a: Any, **_k: Any) -> None:
|
||||
pass
|
||||
|
||||
def mark_done(self, _db: Any, _rid: int, counters: dict[str, Any]) -> None:
|
||||
finals.append(dict(counters))
|
||||
|
||||
async def _fetch(*_a: Any, on_bucket: Any = None, **_k: Any) -> None:
|
||||
on_bucket("room_1_komn:0:3999999", [MagicMock(source_id="a1")], True)
|
||||
on_bucket("room_1_komn:4000000:4999999", [MagicMock(source_id="a2")], False)
|
||||
|
||||
scraper = MagicMock()
|
||||
scraper.__aenter__ = AsyncMock(return_value=scraper)
|
||||
scraper.__aexit__ = AsyncMock(return_value=None)
|
||||
scraper.fetch_all_secondary = _fetch
|
||||
|
||||
with (
|
||||
patch(f"{PFX}.AvitoScraper", return_value=scraper),
|
||||
patch(f"{PFX}.save_listings", MagicMock(return_value=(1, 0))),
|
||||
patch(f"{PFX}.runs", _Recorder()),
|
||||
):
|
||||
counters = await run_avito_full_load(
|
||||
MagicMock(), run_id=1, config=_serp_config(), matcher=MagicMock()
|
||||
)
|
||||
|
||||
assert finals[0]["done_buckets"] == ["room_1_komn:0:3999999"]
|
||||
assert finals[0]["partial_buckets"] == 1
|
||||
assert counters.unique_fetched == 2, "лоты частичного бакета обязаны быть сохранены"
|
||||
|
|
@ -2953,6 +2953,11 @@ class CianFullLoadCounters:
|
|||
detail_enriched: int = 0
|
||||
detail_failed: int = 0
|
||||
errors_count: int = 0
|
||||
# Бакеты, отданные SERP-слоем как НЕполные (страница выпала / исключение
|
||||
# проглочено / hard-cap): лоты сохранены, но в done_buckets бакет не попал.
|
||||
# Без этого счётчика «сколько бакетов чекпоинт не покрывает» видно только грепом
|
||||
# логов, которые теряются при редеплое.
|
||||
partial_buckets: int = 0
|
||||
|
||||
def to_dict(self) -> dict[str, int]:
|
||||
return {f.name: getattr(self, f.name) for f in fields(self)}
|
||||
|
|
@ -3007,8 +3012,32 @@ async def run_cian_full_load(
|
|||
|
||||
done: set[str] = set(skip_set) # накапливаем завершённые бакеты этого прогона
|
||||
|
||||
def _on_bucket(bucket_key: str, lots: list) -> None: # type: ignore[type-arg]
|
||||
"""Инкрементальный save после каждого leaf-бакета. Дописывает bucket_key в done."""
|
||||
def _mark_bucket(bucket_key: str, complete: bool) -> None:
|
||||
"""В чекпоинт — только ПОЛНОСТЬЮ собранный бакет; частичный лишь считаем.
|
||||
|
||||
Ключ done_buckets ("room:lo:hi") одинаков у целого и у недособранного бакета,
|
||||
поэтому признак полноты приезжает отдельным аргументом из SERP-слоя. Пропуск
|
||||
частичного бакета на следующем прогоне означал бы, что его непрочитанные
|
||||
страницы не перечитает уже никто, а счётчики покажут успех.
|
||||
"""
|
||||
if complete:
|
||||
done.add(bucket_key)
|
||||
return
|
||||
counters.partial_buckets += 1
|
||||
logger.warning(
|
||||
"cian-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 → лоты сохраняем, бакет в чекпоинт не пишем
|
||||
(см. _mark_bucket). Дефолт True — для вызывающих без пагинации.
|
||||
"""
|
||||
nonlocal done
|
||||
if runs.is_cancelled(db, run_id):
|
||||
logger.info("cian-full-load run_id=%d: cancel detected in on_bucket", run_id)
|
||||
|
|
@ -3023,7 +3052,7 @@ async def run_cian_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)}
|
||||
)
|
||||
|
|
@ -3042,7 +3071,7 @@ async def run_cian_full_load(
|
|||
counters.saved_inserted += inserted
|
||||
counters.saved_updated += updated
|
||||
counters.unique_fetched += len(lots)
|
||||
done.add(bucket_key)
|
||||
_mark_bucket(bucket_key, complete)
|
||||
runs.update_heartbeat(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)})
|
||||
logger.info(
|
||||
"cian-full-load run_id=%d: bucket %s saved ins=%d upd=%d total_unique=%d",
|
||||
|
|
@ -3225,7 +3254,8 @@ async def run_cian_full_load(
|
|||
|
||||
except NoProxyAvailableError as exc:
|
||||
# #2687, см. ту же ветку в run_avito_full_load: наш отказ (пул пуст) обязан
|
||||
# сохранять done_buckets, а generic-RuntimeError → mark_failed его теряет.
|
||||
# сохранять done_buckets, а generic-RuntimeError уходил в mark_failed, который
|
||||
# его терял (теперь counters мержатся — runs.update_heartbeat/mark_*, #930).
|
||||
logger.error("cian-full-load run_id=%d: no proxy available — %s", run_id, exc)
|
||||
counters.errors_count += 1
|
||||
runs.mark_banned(
|
||||
|
|
@ -3467,7 +3497,8 @@ async def run_yandex_full_load(
|
|||
|
||||
except NoProxyAvailableError as exc:
|
||||
# #2687, см. ту же ветку в run_avito_full_load: наш отказ (пул пуст) обязан
|
||||
# сохранять done_buckets, а generic-RuntimeError → mark_failed его теряет.
|
||||
# сохранять done_buckets, а generic-RuntimeError уходил в mark_failed, который
|
||||
# его терял (теперь counters мержатся — runs.update_heartbeat/mark_*, #930).
|
||||
logger.error("yandex-full-load run_id=%d: no proxy available — %s", run_id, exc)
|
||||
counters.errors_count += 1
|
||||
runs.mark_banned(
|
||||
|
|
@ -3527,6 +3558,8 @@ class AvitoFullLoadCounters:
|
|||
saved_inserted: int = 0
|
||||
saved_updated: int = 0
|
||||
errors_count: int = 0
|
||||
# Бакеты, отданные SERP-слоем как НЕполные — см. CianFullLoadCounters.
|
||||
partial_buckets: int = 0
|
||||
|
||||
def to_dict(self) -> dict[str, int]:
|
||||
return {f.name: getattr(self, f.name) for f in fields(self)}
|
||||
|
|
@ -3599,8 +3632,32 @@ async def run_avito_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_buckets ("room:lo:hi") одинаков у целого и у недособранного бакета,
|
||||
поэтому признак полноты приезжает отдельным аргументом из SERP-слоя. Пропуск
|
||||
частичного бакета на следующем прогоне означал бы, что его непрочитанные
|
||||
страницы не перечитает уже никто, а счётчики покажут успех.
|
||||
"""
|
||||
if complete:
|
||||
done.add(bucket_key)
|
||||
return
|
||||
counters.partial_buckets += 1
|
||||
logger.warning(
|
||||
"avito-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 → лоты сохраняем, бакет в чекпоинт не пишем
|
||||
(см. _mark_bucket). Дефолт True — для вызывающих без пагинации.
|
||||
"""
|
||||
nonlocal done
|
||||
if runs.is_cancelled(db, run_id):
|
||||
logger.info("avito-full-load run_id=%d: cancel detected in on_bucket", run_id)
|
||||
|
|
@ -3613,7 +3670,7 @@ async def run_avito_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)}
|
||||
)
|
||||
|
|
@ -3632,7 +3689,7 @@ async def run_avito_full_load(
|
|||
counters.saved_inserted += inserted
|
||||
counters.saved_updated += updated
|
||||
counters.unique_fetched += len(lots)
|
||||
done.add(bucket_key)
|
||||
_mark_bucket(bucket_key, complete)
|
||||
runs.update_heartbeat(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)})
|
||||
logger.info(
|
||||
"avito-full-load run_id=%d: bucket=%s saved ins=%d upd=%d total_unique=%d",
|
||||
|
|
@ -3699,9 +3756,10 @@ async def run_avito_full_load(
|
|||
except NoProxyAvailableError as exc:
|
||||
# #2687: пул опустел mid-run. Ветка стоит ДО generic-RuntimeError намеренно —
|
||||
# NoProxyAvailableError его подкласс, и без неё отказ уходил в mark_failed,
|
||||
# который (в отличие от mark_banned) НЕ пишет done_buckets. То есть чекпоинт
|
||||
# который (в отличие от mark_banned) НЕ передавал done_buckets. То есть чекпоинт
|
||||
# терялся ровно на НАШЕМ отказе — том исходе, для которого #2686 требовал его
|
||||
# сохранять наравне с блокировкой площадкой.
|
||||
# сохранять наравне с блокировкой площадкой. Диагноз ban_kind='infra' ветка
|
||||
# даёт по-прежнему; сам чекпоинт с #930-мержем counters переживает и mark_failed.
|
||||
logger.error("avito-full-load run_id=%d: no proxy available — %s", run_id, exc)
|
||||
counters.errors_count += 1
|
||||
runs.mark_banned(
|
||||
|
|
|
|||
|
|
@ -100,8 +100,8 @@ def _leading_streak(rows: list[Any], is_bad: Callable[[Any], bool]) -> int:
|
|||
# РЯДОМ со status='banned', а не ВМЕСТО него — сознательный выбор между «новый
|
||||
# статус» и «явное поле причины»:
|
||||
# 1. Побочная функция 'banned' — сохранение done_buckets-чекпоинта (mark_failed
|
||||
# его теряет) — нужна ОБОИМ исходам. Оставив статус, получаем её даром; расщепив
|
||||
# статус, пришлось бы дублировать её в каждом потребителе.
|
||||
# его тогда терял) — нужна ОБОИМ исходам. Оставив статус, получаем её даром;
|
||||
# расщепив статус, пришлось бы дублировать её в каждом потребителе.
|
||||
# 2. Новое значение статуса пришлось бы доучить пяти местам, каждое из которых
|
||||
# молча даёт неверный ответ, если про него забыть: CHECK-констрейнт схемы,
|
||||
# IN-списки обоих сторожей (_alert_if_consecutive_failures / _zero_results),
|
||||
|
|
@ -570,12 +570,21 @@ def mark_skipped(db: Session, *, source: str, reason: str, details: str | None =
|
|||
return int(row.id)
|
||||
|
||||
|
||||
def update_heartbeat(db: Session, run_id: int, counters: dict[str, int]) -> None:
|
||||
"""UPDATE heartbeat_at + counters=:counters + total_seen/new_count колонки.
|
||||
def update_heartbeat(db: Session, run_id: int, counters: dict[str, Any]) -> None:
|
||||
"""UPDATE heartbeat_at + counters (МЕРЖ, не замена) + total_seen/new_count колонки.
|
||||
|
||||
total_seen/new_count извлекаются из counters (lots_fetched/lots_inserted) и
|
||||
пишутся в выделенные колонки, чтобы observability не показывала 0 (audit #1926).
|
||||
COALESCE: если ключа нет в counters — старое значение колонки сохраняется.
|
||||
|
||||
`counters || :counters` вместо замены (#930 добивка): чекпоинт `done_buckets`
|
||||
писали ТОЛЬКО сайты, знающие о нём (`_on_bucket`), а heartbeat'ы, которые о нём не
|
||||
знают, целиком перезаписывали объект и СТИРАЛИ точку. На проде это давало
|
||||
немонотонную точку: `_on_progress` (после каждой комнатности) и фоновый heartbeat
|
||||
cian'а (каждые 60 с) отправляли `counters.to_dict()` без ключа — то есть у cian
|
||||
точка в БД жила лишь от сохранения бакета до ближайшего тика. Прод-след: 15
|
||||
оборванных прогонов с доказанной работой (35 706 + 9 222 fetched) и БЕЗ ключа
|
||||
вообще. Мерж делает точку монотонной для любого писателя, а не только для знающих.
|
||||
"""
|
||||
total_seen, new_count = _column_counts(counters)
|
||||
db.execute(
|
||||
|
|
@ -583,7 +592,7 @@ def update_heartbeat(db: Session, run_id: int, counters: dict[str, int]) -> None
|
|||
"""
|
||||
UPDATE scrape_runs
|
||||
SET heartbeat_at = clock_timestamp(),
|
||||
counters = CAST(:counters AS jsonb),
|
||||
counters = COALESCE(counters, CAST('{}' AS jsonb)) || CAST(:counters AS jsonb),
|
||||
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
|
||||
new_count = COALESCE(CAST(:new_count AS int), new_count)
|
||||
WHERE id = :run_id
|
||||
|
|
@ -641,7 +650,7 @@ def mark_done(db: Session, run_id: int, counters: dict[str, int]) -> None:
|
|||
UPDATE scrape_runs
|
||||
SET status = 'done',
|
||||
finished_at = clock_timestamp(), heartbeat_at = clock_timestamp(),
|
||||
counters = CAST(:counters AS jsonb),
|
||||
counters = COALESCE(counters, CAST('{}' AS jsonb)) || CAST(:counters AS jsonb),
|
||||
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
|
||||
new_count = COALESCE(CAST(:new_count AS int), new_count)
|
||||
WHERE id = :run_id AND status = 'running'
|
||||
|
|
@ -681,7 +690,8 @@ def mark_failed(db: Session, run_id: int, error: str, counters: dict[str, int])
|
|||
UPDATE scrape_runs
|
||||
SET status = 'failed',
|
||||
finished_at = clock_timestamp(), heartbeat_at = clock_timestamp(),
|
||||
error = :error, counters = CAST(:counters AS jsonb),
|
||||
error = :error,
|
||||
counters = COALESCE(counters, CAST('{}' AS jsonb)) || CAST(:counters AS jsonb),
|
||||
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
|
||||
new_count = COALESCE(CAST(:new_count AS int), new_count)
|
||||
WHERE id = :run_id AND status = 'running'
|
||||
|
|
@ -715,8 +725,9 @@ def mark_banned(
|
|||
|
||||
Per migration 015 — 'banned' задокументирован как 'Avito вернул 403/captcha'.
|
||||
Отличается от 'failed': прогон оборван внешним/блокирующим условием, а не нашим
|
||||
багом, и — важно — СОХРАНЯЕТ done_buckets-чекпоинт в counters (mark_failed его
|
||||
теряет). Cooldown 2-4 часа.
|
||||
багом. Чекпоинт done_buckets раньше сохранял только этот финализатор — теперь
|
||||
counters мержатся во всех (см. update_heartbeat), и точка переживает любой из них.
|
||||
Cooldown 2-4 часа.
|
||||
|
||||
`ban_kind` разводит два исхода, которые раньше схлопывались в один статус:
|
||||
- BAN_KIND_PLATFORM — площадка нас заблокировала (firewall/403/captcha);
|
||||
|
|
@ -742,7 +753,8 @@ def mark_banned(
|
|||
UPDATE scrape_runs
|
||||
SET status = 'banned',
|
||||
finished_at = clock_timestamp(), heartbeat_at = clock_timestamp(),
|
||||
error = :error, counters = CAST(:counters AS jsonb),
|
||||
error = :error,
|
||||
counters = COALESCE(counters, CAST('{}' AS jsonb)) || CAST(:counters AS jsonb),
|
||||
ban_kind = :ban_kind,
|
||||
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
|
||||
new_count = COALESCE(CAST(:new_count AS int), new_count)
|
||||
|
|
|
|||
|
|
@ -497,6 +497,132 @@ def _claim_run(db: Session, schedule_row: dict[str, Any], ctx: SchedulerContext)
|
|||
return run_id
|
||||
|
||||
|
||||
# ── подхват чекпоинта оборванного прогона (#930, вторая половина) ────────────
|
||||
# #930 сделал чекпоинт (`counters.done_buckets`) и приёмную сторону
|
||||
# (`run_*_full_load(resume_run_id=...)`), но единственным входом оставил админку. У
|
||||
# avito full-load её нет вовсе, поэтому 299 из 433 «впустую перебранных» корзин за 90
|
||||
# суток не имели НИКАКОГО пути возобновления, даже ручного. Планировщик передавал
|
||||
# resume_run_id=None литералом.
|
||||
#
|
||||
# Точка берётся только когда выполнены ВСЕ условия ниже; иначе прогон честно начинает с
|
||||
# нуля, а ПРИЧИНА пишется в его counters (молчаливый отказ неотличим от отсутствия
|
||||
# правки — см. _resume_decision).
|
||||
|
||||
# 'zombie' НЕ в списке НАМЕРЕННО. reap_zombies снимает пометку 'running', но НЕ убивает
|
||||
# процесс (прямо задокументировано в app/tasks/listing_source_snapshot.py) — а
|
||||
# has_running_run гейтит claim именно по статусу. То есть после reap'а старый сборщик
|
||||
# может продолжать писать в ту же строку: подхват читал бы ДВИЖУЩУЮСЯ точку и запускал
|
||||
# второй сборщик на ту же площадку. 147 корзин в 10 zombie-прогонах за 90 суток
|
||||
# остаются несобранными сознательно — это цена, а не недосмотр.
|
||||
# 'failed' в списке: его чекпоинт больше не стирается финализатором (runs.py, мерж
|
||||
# counters), а причина отказа («наш баг») ничего не говорит о полноте УЖЕ записанных
|
||||
# корзин — они записаны тем же heartbeat'ом, что и у banned.
|
||||
_RESUME_STATUSES = frozenset({"banned", "cancelled", "failed"})
|
||||
|
||||
# Длина цепочки возобновлений. Не круглое число: STALE_DIGEST_INTERVAL_FACTOR (=3) —
|
||||
# уже существующий в этом файле порог «источник не собирал дольше 3× своего такта =
|
||||
# сломан». Цепочка не имеет права отодвинуть полный обход дальше этой же черты,
|
||||
# поэтому подряд идущих подхватов допускается на один меньше: прогоны 1 и 2 могут
|
||||
# продолжать предшественника, третий обязан пойти с нуля. При такте avito 7 суток это
|
||||
# гарантирует попытку полного обхода не реже, чем раз в 21 сутки — ровно в тот момент,
|
||||
# когда сводка объявляет источник просроченным.
|
||||
_MAX_RESUME_CHAIN = STALE_DIGEST_INTERVAL_FACTOR - 1
|
||||
|
||||
# Кандидат — ПОСЛЕДНИЙ прогон источника, а не последний подходящий: если после обрыва
|
||||
# уже прошёл полный ('done') прогон, дерево обойдено и возобновлять нечего. Строки
|
||||
# 'skipped' — бухгалтерия планировщика, а не прогоны, поэтому не в счёт.
|
||||
_RESUME_CANDIDATE_SQL = text("""
|
||||
WITH cur AS (
|
||||
SELECT source, params FROM scrape_runs WHERE id = CAST(:rid AS bigint)
|
||||
)
|
||||
SELECT r.id AS prev_id,
|
||||
r.status AS prev_status,
|
||||
r.counters AS prev_counters,
|
||||
(r.params IS NOT DISTINCT FROM cur.params) AS same_params,
|
||||
EXTRACT(EPOCH FROM (clock_timestamp() - r.heartbeat_at)) / 3600.0 AS age_h,
|
||||
cur.params ->> 'interval_days' AS interval_days
|
||||
FROM scrape_runs r, cur
|
||||
WHERE r.source = cur.source
|
||||
AND r.id <> CAST(:rid AS bigint)
|
||||
AND r.status <> 'skipped'
|
||||
ORDER BY r.started_at DESC
|
||||
LIMIT 1
|
||||
""")
|
||||
|
||||
|
||||
def _resume_decision(row: Any) -> tuple[int | None, dict[str, Any]]:
|
||||
"""Решение «подхватывать ли точку» + счётчики-объяснение. Чистая функция.
|
||||
|
||||
Возвращает (resume_run_id | None, counters-заготовка нового прогона). Причина
|
||||
отказа — машиночитаемый слаг в `resume_reason`, по нему «предыдущего прогона не
|
||||
было» отличается от «параметры разъехались» ЗАПРОСОМ, а не грепом логов.
|
||||
|
||||
Срок годности точки — такт источника ПЛЮС сутки (`interval_days` из его же params).
|
||||
Ни одно из слагаемых не выбрано произвольно. Такт — объявленный самим расписанием
|
||||
срок, в течение которого собранное считается свежим; точка старше него пережила
|
||||
цикл, в котором источник обязан был обойти дерево целиком. Сутки — сетка запуска:
|
||||
`compute_next_run_at` выбирает ДЕНЬ (сегодня+interval_days) и случайное время внутри
|
||||
окна, поэтому два соседних запуска отстоят друг от друга на interval_days ± меньше
|
||||
суток. Без этого слагаемого точку отвергал бы jitter расписания, а не устаревание:
|
||||
прогон 3547 убит деплоем 09.08 16:53, расписание 139 подхватит его 16.08 13:37 —
|
||||
164.6 ч при такте 168 ч, запас 3.4 ч при ширине окна 2 ч. Пропущенный цикл в окно
|
||||
всё равно не влезает: для avito это 13 суток против порога 8.
|
||||
"""
|
||||
if row is None:
|
||||
return None, {"resume_from": None, "resume_reason": "no_prev_run", "resume_chain": 0}
|
||||
|
||||
prev_counters = row.prev_counters if isinstance(row.prev_counters, dict) else {}
|
||||
done_buckets = prev_counters.get("done_buckets")
|
||||
done_n = len(done_buckets) if isinstance(done_buckets, list) else 0
|
||||
chain_raw = prev_counters.get("resume_chain")
|
||||
prev_chain = chain_raw if isinstance(chain_raw, int) else 0
|
||||
verdict: dict[str, Any] = {
|
||||
"resume_from": None,
|
||||
"resume_candidate": int(row.prev_id),
|
||||
"resume_buckets": done_n,
|
||||
"resume_chain": 0,
|
||||
}
|
||||
|
||||
if row.prev_status not in _RESUME_STATUSES:
|
||||
verdict["resume_reason"] = f"status_{row.prev_status}"
|
||||
elif not row.same_params:
|
||||
verdict["resume_reason"] = "params_changed"
|
||||
elif done_n == 0:
|
||||
verdict["resume_reason"] = "no_checkpoint"
|
||||
elif row.age_h is None or float(row.age_h) > 24.0 * (
|
||||
_schedule_interval_days(row.interval_days) + 1
|
||||
):
|
||||
verdict["resume_reason"] = "checkpoint_stale"
|
||||
elif prev_chain >= _MAX_RESUME_CHAIN:
|
||||
verdict["resume_reason"] = "chain_limit"
|
||||
else:
|
||||
verdict["resume_from"] = int(row.prev_id)
|
||||
verdict["resume_reason"] = "ok"
|
||||
verdict["resume_chain"] = prev_chain + 1
|
||||
return int(row.prev_id), verdict
|
||||
return None, verdict
|
||||
|
||||
|
||||
def _pick_resume(db: Session, run_id: int) -> int | None:
|
||||
"""Чекпоинт какого прогона наследует `run_id` (или None) + запись вердикта.
|
||||
|
||||
Тождество задания сверяется РОВНО по тем полям, которыми задание задаётся: source
|
||||
(кандидат ищется в пределах одного source) и params целиком, побайтово. Кандидат
|
||||
сравнивается с ТЕКУЩИМ прогоном, а не с расписанием, потому что именно params
|
||||
прогона поехали в pipeline. На проде за 90 суток 89 корзин из 433 (21%) лежат в
|
||||
прогонах, чьи params отличаются от следующего — по ним пропуск был бы неверным:
|
||||
`incremental_days` меняет СМЫСЛ ключа (дочитано до watermark ≠ бакет перебран), а
|
||||
`price_cap_per_bucket` меняет само дерево бисекции, то есть какие ключи существуют.
|
||||
"""
|
||||
row = db.execute(_RESUME_CANDIDATE_SQL, {"rid": run_id}).fetchone()
|
||||
resume_run_id, verdict = _resume_decision(row)
|
||||
# Вердикт кладём в counters НОВОГО прогона: `update_heartbeat` мержит jsonb, поэтому
|
||||
# последующие heartbeat'ы пайплайна его не затрут и он доживёт до финализатора.
|
||||
_kit_runs.update_heartbeat(db, run_id, verdict)
|
||||
logger.info("scheduler: resume run_id=%d — %s", run_id, verdict)
|
||||
return resume_run_id
|
||||
|
||||
|
||||
def _defer_next_run_at(db: Session, schedule_row: dict[str, Any]) -> None:
|
||||
"""Сдвинуть next_run_at на следующее окно БЕЗ создания run (#1522).
|
||||
|
||||
|
|
@ -705,7 +831,7 @@ async def _job_avito_full_load(
|
|||
concurrency=int(params.get("concurrency", 5)),
|
||||
request_delay_sec=float(params.get("request_delay_sec", 7.0)),
|
||||
secondary_only=bool(params.get("secondary_only", True)),
|
||||
resume_run_id=None,
|
||||
resume_run_id=_pick_resume(db, run_id),
|
||||
incremental_days=incremental_days,
|
||||
)
|
||||
|
||||
|
|
@ -725,7 +851,7 @@ async def _job_avito_full_load_exhaustive(
|
|||
concurrency=int(params.get("concurrency", 5)),
|
||||
request_delay_sec=float(params.get("request_delay_sec", 7.0)),
|
||||
secondary_only=bool(params.get("secondary_only", True)),
|
||||
resume_run_id=None,
|
||||
resume_run_id=_pick_resume(db, run_id),
|
||||
incremental_days=None,
|
||||
)
|
||||
|
||||
|
|
@ -796,7 +922,7 @@ async def _job_cian_full_load(
|
|||
request_delay_sec=float(params.get("request_delay_sec", 4.0)),
|
||||
enrich_detail=bool(params.get("enrich_detail", False)),
|
||||
detail_top_n=int(params.get("detail_top_n", 0)),
|
||||
resume_run_id=None,
|
||||
resume_run_id=_pick_resume(db, run_id),
|
||||
)
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -1095,7 +1095,9 @@ class AvitoScraper(BaseScraper):
|
|||
leaf-бакета. Может быть async или sync. Исключение прерывает прогон.
|
||||
on_progress: опциональный callback(unique_count) для heartbeat (per room-bucket).
|
||||
skip_buckets: множество ключей «room_label:lo:hi» уже завершённых бакетов —
|
||||
пагинация и on_bucket для них пропускаются. Probe-запросы выполняются.
|
||||
пагинация и on_bucket для них пропускаются. В exhaustive-режиме probe-запросы
|
||||
всё равно выполняются (skip проверяется уже в листе, после probe), в
|
||||
инкрементальном probe'а нет — там пропускается весь бакет целиком.
|
||||
since: если None (default) — EXHAUSTIVE bisection-обход (поведение без
|
||||
изменений). Если задана date — INCREMENTAL: на каждый (комнатность ×
|
||||
seed-брекет) последовательная пагинация newest-first с ранней остановкой,
|
||||
|
|
@ -1141,6 +1143,7 @@ class AvitoScraper(BaseScraper):
|
|||
max_pages_per_bucket=max_pages_per_bucket,
|
||||
secondary_only=secondary_only,
|
||||
on_bucket=on_bucket,
|
||||
skip_buckets=skip_buckets,
|
||||
)
|
||||
else:
|
||||
await self._walk_price_range(
|
||||
|
|
@ -1390,14 +1393,17 @@ class AvitoScraper(BaseScraper):
|
|||
# иначе фетчим её как обычную страницу (открытый брекет без probe-html).
|
||||
sem = asyncio.Semaphore(concurrency)
|
||||
first_url = self._build_rooms_url(room_slug, 1, _lo_param, _hi_param)
|
||||
dropped_pages = 0 # страницы, не отдавшие карточки по отказу (не по пустоте)
|
||||
|
||||
async def _one_page(p: int) -> list[ScrapedLot]:
|
||||
nonlocal dropped_pages
|
||||
if p == 1 and html is not None:
|
||||
return self._parse_html(html, source_url_base=first_url)
|
||||
async with sem:
|
||||
page_html = await self._fetch_rooms_page_html(room_slug, p, _lo_param, _hi_param)
|
||||
await asyncio.sleep(self.request_delay_sec)
|
||||
if page_html is None:
|
||||
dropped_pages += 1
|
||||
logger.warning(
|
||||
"avito: page_html=None %s [%d, %s] page=%d — skipping page",
|
||||
room_label,
|
||||
|
|
@ -1420,6 +1426,7 @@ class AvitoScraper(BaseScraper):
|
|||
# Блокировки пробрасываем наверх (mark_banned в pipeline-обёртке).
|
||||
if isinstance(res, AvitoBlockedError | AvitoRateLimitedError):
|
||||
raise res
|
||||
dropped_pages += 1
|
||||
logger.warning(
|
||||
"avito: page exception %s [%d, %s] page=%d — %r",
|
||||
room_label,
|
||||
|
|
@ -1464,8 +1471,25 @@ class AvitoScraper(BaseScraper):
|
|||
if key:
|
||||
seen[key] = lot
|
||||
|
||||
# ── Полнота бакета: ключ у целиком и частично собранного ОДИНАКОВ ─────
|
||||
# `bucket_key` — это "room:lo:hi" и больше ничего, поэтому «сделано» едет
|
||||
# отдельным аргументом. Три пути частичности сходятся здесь:
|
||||
# 1. страница молча выпала (page_html=None выше);
|
||||
# 2. исключение страницы проглочено gather'ом (return_exceptions=True);
|
||||
# 3. признанный tail-loss — probe провалился (expected_total=None, открытый
|
||||
# брекет) или страниц нужно больше, чем max_pages.
|
||||
# Ни один из них не оставлял следа в чекпоинте: следующий прогон видел ключ и
|
||||
# пропускал бакет. Пока признак не доехал до done_buckets, включать подхват
|
||||
# нельзя — видимая потеря (перескрап) стала бы невидимой (пропуск страниц).
|
||||
pages_needed = (
|
||||
math.ceil(expected_total / _AVITO_OFFERS_PER_PAGE)
|
||||
if expected_total is not None
|
||||
else None
|
||||
)
|
||||
complete = dropped_pages == 0 and pages_needed is not None and pages_needed <= max_pages
|
||||
logger.info(
|
||||
"avito: %s [%d, %s] paginated=%d pages collected=%d dropped_nb=%d unique_total=%d",
|
||||
"avito: %s [%d, %s] paginated=%d pages collected=%d dropped_nb=%d unique_total=%d "
|
||||
"complete=%s dropped_pages=%d",
|
||||
room_label,
|
||||
lo,
|
||||
_hi_repr,
|
||||
|
|
@ -1473,11 +1497,13 @@ class AvitoScraper(BaseScraper):
|
|||
collected_this_bucket,
|
||||
dropped_nb,
|
||||
len(seen),
|
||||
complete,
|
||||
dropped_pages,
|
||||
)
|
||||
|
||||
# ── on_bucket callback: инкрементальный save ──────────────────────────
|
||||
if on_bucket is not None and bucket_lots:
|
||||
res_cb = on_bucket(bucket_key, bucket_lots)
|
||||
res_cb = on_bucket(bucket_key, bucket_lots, complete)
|
||||
if inspect.isawaitable(res_cb):
|
||||
await res_cb
|
||||
|
||||
|
|
@ -1493,6 +1519,7 @@ class AvitoScraper(BaseScraper):
|
|||
max_pages_per_bucket: int,
|
||||
secondary_only: bool,
|
||||
on_bucket: Callable[..., Any] | None,
|
||||
skip_buckets: set[str] | None = None,
|
||||
) -> None:
|
||||
"""INCREMENTAL пагинация одного (комнатность × seed-брекет) с ранней остановкой.
|
||||
|
||||
|
|
@ -1514,16 +1541,37 @@ class AvitoScraper(BaseScraper):
|
|||
bucket_key, secondary_only-фильтр и дедуп в seen — идентичны _paginate_leaf_bucket.
|
||||
on_bucket вызывается один раз для собранного бакета (async/sync-aware).
|
||||
AvitoBlockedError/AvitoRateLimitedError из page-фетчей пробрасываются наверх.
|
||||
|
||||
skip_buckets: ключи, дочитанные ПРЕДЫДУЩИМ прогоном с ТЕМИ ЖЕ params (тождество
|
||||
задания проверяет планировщик, `_pick_resume`). В инкрементальном режиме
|
||||
«сделано» значит «дочитал до watermark `since`», а не «перебрал бакет целиком»,
|
||||
поэтому смешивать такой ключ с exhaustive-ключом нельзя — они дословно совпадают
|
||||
(6 из 11 seed-ключей), но означают разное. Пропуск здесь безопасен по покрытию
|
||||
ровно потому, что окно ретроспективы не уже такта (#2674, гарантируется
|
||||
_job_avito_full_load): бакет, дочитанный до watermark N суток назад, следующий
|
||||
плановый прогон перечитает со своим since = сегодня−N и увидит всё, что успело
|
||||
появиться. Раньше аргумент сюда не передавался вовсе — resume в боевом
|
||||
(инкрементальном) режиме avito_full_load был чистым no-op: 134 из 242 корзин.
|
||||
|
||||
complete=False (см. `_paginate_leaf_bucket`) отдаётся, когда бакет НЕ дочитан до
|
||||
watermark: страница выпала, или страниц не хватило (max_pages), или остановка
|
||||
произошла по grace-эвристике «2 подряд недатированные страницы» — там watermark
|
||||
не доказан, а не достигнут.
|
||||
"""
|
||||
_lo_param = lo if lo > 0 else None
|
||||
_hi_param = hi # None → _build_rooms_url не ставит pmax
|
||||
_hi_repr = "open" if hi is None else str(hi)
|
||||
bucket_key = f"{room_label}:{lo}:{_hi_repr}"
|
||||
|
||||
if skip_buckets and bucket_key in skip_buckets:
|
||||
logger.info("avito: skip bucket %s — already read to watermark (resume)", bucket_key)
|
||||
return
|
||||
|
||||
bucket_lots: list[ScrapedLot] = []
|
||||
pages_fetched = 0
|
||||
not_fresh_streak = 0 # подряд идущие не-свежие (all None-или-старые) страницы
|
||||
stop_reason = "max-pages" # перетирается ниже на реальную причину
|
||||
complete = False # дочитан ли бакет до watermark; True только на честных стопах
|
||||
|
||||
for p in range(1, max_pages_per_bucket + 1):
|
||||
# Последовательный фетч (НЕ asyncio.gather): early-stop требует читать
|
||||
|
|
@ -1540,6 +1588,7 @@ class AvitoScraper(BaseScraper):
|
|||
page_lots = self._parse_html(page_html, source_url_base=page_url)
|
||||
if not page_lots:
|
||||
stop_reason = "end-of-pages (empty parse)"
|
||||
complete = True # выдача кончилась — читать в этом брекете больше нечего
|
||||
break
|
||||
|
||||
bucket_lots.extend(page_lots)
|
||||
|
|
@ -1556,6 +1605,7 @@ class AvitoScraper(BaseScraper):
|
|||
# Есть даты, но все < since → newest-first гарантирует, что дальше
|
||||
# только старее → стоп немедленно.
|
||||
stop_reason = "early-stop (page all older than since)"
|
||||
complete = True # watermark достигнут — ровно то, что значит «сделано»
|
||||
break
|
||||
# Все карточки undated (None) → не стопим сразу (могут быть свежие без
|
||||
# даты), но копим streak; 2 подряд недатированные страницы → grace-стоп.
|
||||
|
|
@ -1581,7 +1631,7 @@ class AvitoScraper(BaseScraper):
|
|||
|
||||
logger.info(
|
||||
"avito: %s [%d, %s] paginated=%d pages collected=%d dropped_nb=%d "
|
||||
"unique_total=%d incremental since=%s stop=%s",
|
||||
"unique_total=%d incremental since=%s stop=%s complete=%s",
|
||||
room_label,
|
||||
lo,
|
||||
_hi_repr,
|
||||
|
|
@ -1591,11 +1641,12 @@ class AvitoScraper(BaseScraper):
|
|||
len(seen),
|
||||
since.isoformat(),
|
||||
stop_reason,
|
||||
complete,
|
||||
)
|
||||
|
||||
# ── on_bucket callback: инкрементальный save ──────────────────────────
|
||||
if on_bucket is not None and bucket_lots:
|
||||
res_cb = on_bucket(bucket_key, bucket_lots)
|
||||
res_cb = on_bucket(bucket_key, bucket_lots, complete)
|
||||
if inspect.isawaitable(res_cb):
|
||||
await res_cb
|
||||
|
||||
|
|
|
|||
|
|
@ -621,8 +621,10 @@ class CianScraper(BaseScraper):
|
|||
|
||||
# Страница 1 уже есть (html из probe выше); остальные — параллельно.
|
||||
sem = asyncio.Semaphore(concurrency)
|
||||
dropped_pages = 0 if html else 1 # пустой probe-html = страница 1 не собрана
|
||||
|
||||
async def _one_page(p: int) -> list[ScrapedLot]:
|
||||
nonlocal dropped_pages
|
||||
if p == 1:
|
||||
# Используем уже полученный HTML от probe
|
||||
return self._parse_serp_html(html) if html else []
|
||||
|
|
@ -630,6 +632,7 @@ class CianScraper(BaseScraper):
|
|||
page_html = await self._fetch_page_html(rooms, p, _lo_param, hi)
|
||||
await asyncio.sleep(self.request_delay_sec)
|
||||
if page_html is None:
|
||||
dropped_pages += 1
|
||||
logger.warning(
|
||||
"cian: page_html=None %s [%d, %s] page=%d — skipping page",
|
||||
room_label,
|
||||
|
|
@ -648,6 +651,7 @@ class CianScraper(BaseScraper):
|
|||
bucket_lots: list[ScrapedLot] = []
|
||||
for p_idx, res in enumerate(page_results, start=1):
|
||||
if isinstance(res, BaseException):
|
||||
dropped_pages += 1
|
||||
logger.warning(
|
||||
"cian: page exception %s [%d, %s] page=%d — %r",
|
||||
room_label,
|
||||
|
|
@ -674,8 +678,16 @@ class CianScraper(BaseScraper):
|
|||
if key:
|
||||
seen[key] = lot
|
||||
|
||||
# ── Полнота бакета: ключ у целиком и частично собранного ОДИНАКОВ ─────
|
||||
# См. avito/serp.py — те же два пути частичности (выпавшая страница,
|
||||
# проглоченное gather'ом исключение) плюс признанный hard-cap выше. У cian это
|
||||
# тяжелее: в SERP-слое НЕТ класса блок-исключения вообще (ban-детект #2625
|
||||
# агрегатный, на выходе из скрапера), поэтому капча посреди бакета приходит сюда
|
||||
# как page_html=None и раньше давала «сделанный» бакет из уцелевших страниц.
|
||||
complete = dropped_pages == 0 and pages_needed <= max_pages
|
||||
logger.info(
|
||||
"cian: %s [%d, %s] paginated=%d pages collected=%d dropped_nb=%d unique_total=%d",
|
||||
"cian: %s [%d, %s] paginated=%d pages collected=%d dropped_nb=%d unique_total=%d "
|
||||
"complete=%s dropped_pages=%d",
|
||||
room_label,
|
||||
lo,
|
||||
_hi_repr,
|
||||
|
|
@ -683,11 +695,13 @@ class CianScraper(BaseScraper):
|
|||
collected_this_bucket,
|
||||
dropped_nb,
|
||||
len(seen),
|
||||
complete,
|
||||
dropped_pages,
|
||||
)
|
||||
|
||||
# ── on_bucket callback: инкрементальный save ──────────────────────────
|
||||
if on_bucket is not None and bucket_lots:
|
||||
res_cb = on_bucket(bucket_key, bucket_lots)
|
||||
res_cb = on_bucket(bucket_key, bucket_lots, complete)
|
||||
if inspect.isawaitable(res_cb):
|
||||
await res_cb
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue