feat(tradein/scheduler): планировщик подхватывает чекпоинт оборванного прогона (#930) (#2845)
Some checks failed
Deploy Trade-In / test (push) Blocked by required conditions
Deploy Trade-In / build-backend (push) Blocked by required conditions
Deploy Trade-In / build-frontend (push) Blocked by required conditions
Deploy Trade-In / build-browser (push) Blocked by required conditions
Deploy Trade-In / deploy (push) Blocked by required conditions
Deploy Trade-In / changes (push) Has been cancelled

This commit is contained in:
bot-backend 2026-08-12 18:51:02 +00:00
parent 081319c290
commit 1ff6699b95
6 changed files with 588 additions and 32 deletions

View file

@ -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, "лоты частичного бакета обязаны быть сохранены"

View file

@ -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(

View file

@ -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)

View file

@ -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),
)

View file

@ -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

View file

@ -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