From 51027d8b02b4676a55bbc6d3e4a9019108e33173 Mon Sep 17 00:00:00 2001 From: bot-backend Date: Fri, 28 Aug 2026 01:05:54 +0300 Subject: [PATCH] =?UTF-8?q?fix(tradein/ingest):=20rosreestr=5Fdkp=5Fimport?= =?UTF-8?q?=20=D0=BA=D1=83=D1=80=D1=81=D0=BE=D1=80=20=D0=BF=D0=B5=D1=80?= =?UTF-8?q?=D0=B5=D0=B6=D0=B8=D0=B2=D0=B0=D0=B5=D1=82=20=D1=80=D0=B5=D1=81?= =?UTF-8?q?=D1=82=D0=B0=D1=80=D1=82=20(#3168)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit last_id жил только в памяти процесса (import_rosreestr_dkp, scheduler.py): heartbeat писал его в scrape_runs.counters каждый батч (комментарий рядом прямо называл это чекпоинтом), но при старте last_id всегда инициализировался литералом 0 — обрыв (деплой/OOM/рестарт хоста) откатывал прогресс и заставлял пере-сканировать источник с начала. Разведка: из пяти backfill-циклов issue (avito_detail_backfill, house_imv_backfill, cian_history_backfill, yandex_detail_backfill, geocode_missing_listings) ни один не имеет этого дефекта — все устроены как WHERE ... IS NULL/NOT EXISTS ... LIMIT, естественно резюмируемы без курсора. Единственный код, буквально описанный в issue (строки/SQL/комментарий), это шестой, не входящий в таблицу backfill — rosreestr_dkp_import. Фикс — _resume_dkp_cursor(db, run_id): - кандидат — последний прогон source='rosreestr_dkp_import'; - резюмится только незавершённый штатно прогон: status running/zombie, либо done с counters.interrupted=1 (SIGTERM-drain — эта ветка раньше считала последующий full rescan штатным поведением, теперь помечает себя как прерванную и резюмится наравне с zombie); - потолок возраста чекпоинта — 24ч, старше — 'checkpoint_stale', старт с 0; - чистый 'done' (полный проход) не резюмится — иначе ON CONFLICT DO UPDATE перестанет ловить правки уже импортированных сделок при следующем проходе. Вердикт и per-batch чекпоинт пишутся через kit_runs.update_heartbeat (merge `counters || :counters`) вместо локального runs_mod.update_heartbeat (полная замена) — иначе resume-вердикт стирался первым же heartbeat'ом батча. Тесты: tests/test_3168_backfill_cursor_resume.py — резюм с сохранённого last_id, резюм после SIGTERM-drain, отказ резюмить чистый done, отказ резюмить протухший (>24ч) чекпоинт, merge не стирает посторонние ключи. Обратимость проверена вручную (временный откат _resume_dkp_cursor красил 6 из 8 тестов). --- tradein-mvp/backend/app/services/scheduler.py | 117 +++++++++- .../tests/test_3168_backfill_cursor_resume.py | 208 ++++++++++++++++++ 2 files changed, 316 insertions(+), 9 deletions(-) create mode 100644 tradein-mvp/backend/tests/test_3168_backfill_cursor_resume.py diff --git a/tradein-mvp/backend/app/services/scheduler.py b/tradein-mvp/backend/app/services/scheduler.py index b1a43824..d1f7fe7d 100644 --- a/tradein-mvp/backend/app/services/scheduler.py +++ b/tradein-mvp/backend/app/services/scheduler.py @@ -28,6 +28,12 @@ from __future__ import annotations import logging from typing import Any +# kit_runs.update_heartbeat (в отличие от локального runs_mod.update_heartbeat) мержит +# counters (`counters || :counters`) вместо замены — нужен для чекпоинта курсора +# import_rosreestr_dkp (issue #3168), чтобы resume-вердикт, записанный на старте, не +# затирался последующими per-batch heartbeat'ами того же прогона. +from scraper_kit.orchestration import runs as kit_runs + # compute_next_run_at жил здесь ВТОРОЙ, побайтово одинаковой копией kit-версии (#2674). # Обе копии одинаково умели interval_days — но такт доезжал до next_run_at только через # kit (_claim_run/_defer_next_run_at читают default_params["interval_days"]); admin.py @@ -159,6 +165,85 @@ async def _execute_cian_backfill( raise +_DKP_SOURCE = "rosreestr_dkp_import" + +# Потолок возраста чекпоинта: старше — last_id прошлого прогона не подхватываем, прогон +# стартует с id=0 (issue #3168). У предиката `id > last_id` нет протухания в смысле +# свипов (он остаётся корректным сколь угодно долго), но апстрим +# (gendesign_rosreestr_deals) наполняется НЕЗАВИСИМЫМ ETL, который может дописать более +# старую сделку под новым id уже ПОСЛЕ того, как наш курсор её обогнал — курсор недельной +# давности унёс бы с собой всё, что апстрим добавил за эту неделю ниже last_id. Порог +# того же порядка, что суточная сетка резюма в scraper_kit (_resume_decision), и с +# запасом перекрывает самый долгий соседний backfill в семье (avito_detail_backfill, +# 150 мин максимума heartbeat-gap за 30 суток — issue #3168). +_DKP_CHECKPOINT_STALE_HOURS = 24.0 + +_DKP_RESUME_CANDIDATE_SQL = text(""" + SELECT id AS prev_id, + status AS prev_status, + counters AS prev_counters, + EXTRACT(EPOCH FROM (clock_timestamp() - heartbeat_at)) / 3600.0 AS age_h + FROM scrape_runs + WHERE source = CAST(:source AS text) + AND id <> CAST(:rid AS bigint) + AND status <> 'skipped' + ORDER BY started_at DESC + LIMIT 1 +""") + + +def _resume_dkp_cursor(db: Session, run_id: int) -> tuple[int, dict[str, Any]]: + """Продолжить last_id прошлого прогона или начать с 0 — решение + explain (issue #3168). + + last_id раньше жил только в памяти процесса (init на 0 при каждом запуске): heartbeat + писал его в counters КАЖДЫЙ батч (checkpoint), но на старте никто это не читал + обратно. Обрыв (деплой/OOM/рестарт хоста) откатывал прогресс на 0 и заставлял + пере-сканировать источник с начала. + + Кандидат — ПОСЛЕДНИЙ прогон этого source (тот же принцип, что и + scraper_kit.orchestration.scheduler._pick_resume, локальная копия ладдера — контракт + другой: нет params/interval_days, курсор числовой, а не bucket-set): + - 'running' / 'zombie' — прогон, которого не завершили штатно. + - 'done' с counters.interrupted=1 — SIGTERM-drain (см. комментарий у mark_done + ниже по коду): статус 'done', но это НЕ полный проход, прогресс оборван. + - Чистое 'done' без этого флага — полный проход завершился сам, резюмить нечего: + следующий прогон обязан пере-сканировать с 0, иначе ON CONFLICT DO UPDATE + перестанет ловить правки уже импортированных сделок (см. докстринг функции ниже). + + Возвращает (last_id, verdict) — verdict пишется в scrape_runs.counters вызывающим + кодом (kit_runs.update_heartbeat — merge, не замена), чтобы решение было видно в + scrape_runs, а не только в логе. + """ + row = db.execute(_DKP_RESUME_CANDIDATE_SQL, {"source": _DKP_SOURCE, "rid": run_id}).fetchone() + verdict: dict[str, Any] = {"resume_from": None} + + if row is None: + verdict["resume_reason"] = "no_prev_run" + return 0, verdict + + verdict["resume_candidate"] = int(row.prev_id) + prev_counters = row.prev_counters if isinstance(row.prev_counters, dict) else {} + prev_last_id = prev_counters.get("last_id") + interrupted = bool(prev_counters.get("interrupted")) + resumable_status = row.prev_status in ("running", "zombie") or ( + row.prev_status == "done" and interrupted + ) + + if not resumable_status: + verdict["resume_reason"] = f"status_{row.prev_status}" + elif not isinstance(prev_last_id, int) or prev_last_id <= 0: + verdict["resume_reason"] = "no_checkpoint" + elif row.age_h is None or float(row.age_h) > _DKP_CHECKPOINT_STALE_HOURS: + verdict["resume_reason"] = "checkpoint_stale" + else: + verdict["resume_from"] = int(row.prev_id) + verdict["resume_reason"] = "ok" + verdict["last_id"] = prev_last_id + return int(prev_last_id), verdict + + return 0, verdict + + def import_rosreestr_dkp( db: Session, run_id: int, @@ -195,7 +280,10 @@ def import_rosreestr_dkp( Rooms: выводятся из площади (Росреестр не отдаёт кол-во комнат). Batch-процессинг: читаем из FDW батчами по batch_size через cursor-based пагинацию - (WHERE id > last_id ORDER BY id). Heartbeat обновляется каждый батч (= checkpoint). + (WHERE id > last_id ORDER BY id). Heartbeat обновляется каждый батч (= checkpoint), + мержем (kit_runs.update_heartbeat), а не заменой. На старте _resume_dkp_cursor решает + продолжить с last_id прошлого прогона или начать с 0 — чекпоинт переживает рестарт + процесса (деплой/OOM/SIGTERM), пока не старше суток (issue #3168). SAVEPOINT per row — один сбойный row не откатывает батч. Координаты: NULL после импорта — геокодинг остаётся follow-up (geocode-deals). @@ -241,8 +329,14 @@ def import_rosreestr_dkp( ) db.rollback() - last_id: int = 0 + last_id, resume_verdict = _resume_dkp_cursor(db, run_id) total_batches = 0 + kit_runs.update_heartbeat(db, run_id, resume_verdict) + logger.info( + "rosreestr_dkp_import run_id=%d: resume decision — %s", + run_id, + resume_verdict, + ) try: while True: @@ -255,11 +349,14 @@ def import_rosreestr_dkp( return # #1182 Phase 2: кооперативный SIGTERM-drain (деплой recreate scraper). # Это НЕ user-cancel → mark_done (partial), не mark_cancelled. Курсор - # уже зафиксирован update_heartbeat'ом каждый батч; следующий run - # пере-сканирует с id=0 — INSERT ... ON CONFLICT DO UPDATE идемпотентен - # для неизменных строк (WHERE ... IS DISTINCT FROM пропускает совпадающие - # факты), поэтому повторный проход не плодит дубликаты и не трогает данные. + # уже зафиксирован heartbeat'ом каждый батч; следующий run подхватит + # last_id через _resume_dkp_cursor (issue #3168), если чекпоинт не старше + # суток — раньше здесь безусловно пере-сканировали с id=0 при каждом + # обрыве. counters['interrupted']=1 помечает это 'done' как НЕ полный + # проход, чтобы _resume_dkp_cursor не спутал его со штатным завершением + # (которое обязано пере-сканировать с 0 ради ON CONFLICT DO UPDATE правок). # mark_done выводит run из 'running' → reap_zombies его не тронет. + counters["interrupted"] = 1 logger.info( "rosreestr_dkp_import run_id=%d: SIGTERM-drain — committing partial " "(last_id=%d, batches=%d) and exiting", @@ -267,7 +364,7 @@ def import_rosreestr_dkp( last_id, total_batches, ) - runs_mod.update_heartbeat(db, run_id, counters) + kit_runs.update_heartbeat(db, run_id, counters) runs_mod.mark_done(db, run_id, counters) return @@ -441,8 +538,10 @@ def import_rosreestr_dkp( counters["batches_done"] = total_batches counters["last_id"] = last_id # type: ignore[assignment] - # Heartbeat = checkpoint: allows zombie detection + resume visibility - runs_mod.update_heartbeat(db, run_id, counters) + # Heartbeat = checkpoint: allows zombie detection + resume visibility. + # kit_runs (merge, не замена) — иначе этот REPLACE стёр бы resume_verdict, + # записанный _resume_dkp_cursor'ом перед циклом (issue #3168). + kit_runs.update_heartbeat(db, run_id, counters) logger.info( "rosreestr_dkp_import run_id=%d: batch=%d fetched=%d " "inserted=%d updated=%d skipped=%d errored=%d last_id=%d", diff --git a/tradein-mvp/backend/tests/test_3168_backfill_cursor_resume.py b/tradein-mvp/backend/tests/test_3168_backfill_cursor_resume.py new file mode 100644 index 00000000..25e7f499 --- /dev/null +++ b/tradein-mvp/backend/tests/test_3168_backfill_cursor_resume.py @@ -0,0 +1,208 @@ +"""Issue #3168: курсорные backfill-циклы не переживают рестарт (rosreestr_dkp_import). + +Разведка (см. PR): среди пяти backfill-циклов из issue (avito_detail_backfill, +house_imv_backfill, cian_history_backfill, yandex_detail_backfill, +geocode_missing_listings) НИ ОДИН не использует in-memory числовой курсор — все +устроены как WHERE ... IS NULL / NOT EXISTS ... LIMIT (естественно резюмируемы без +персистентной точки, см. cian_history_backfill.py и комментарий в scheduler.py про +"Checkpoint/resume semantics"). Единственное место с буквально описанным в issue дефектом +(курсор `WHERE id > last_id ORDER BY id`, последняя строчка heartbeat = "checkpoint", но +last_id инициализируется литералом 0 при каждом запуске) — app/services/scheduler.py :: +import_rosreestr_dkp (source='rosreestr_dkp_import'), шестой backfill, не входящий в +таблицу issue. Фикс и тесты — про него; остальные пять на отдельный заход НЕ нужны, у них +нет этого дефекта. + +Фикс — _resume_dkp_cursor(db, run_id) в scheduler.py: + - Кандидат — последний прогон source='rosreestr_dkp_import' (кроме себя, кроме + 'skipped'), тот же принцип, что и scraper_kit.orchestration.scheduler._pick_resume + (локальная копия ладдера — контракт другой: нет params/interval_days, курсор + числовой, а не bucket-set, поэтому общий _pick_resume не переиспользован). + - Резюмится только если прошлый прогон НЕ завершился штатно: status in + ('running', 'zombie'), либо 'done' с counters.interrupted=1 (SIGTERM-drain — см. + ветку в import_rosreestr_dkp, которая раньше называла последующий full rescan + штатным поведением). + - Потолок возраста чекпоинта — 24ч (_DKP_CHECKPOINT_STALE_HOURS): старше — курсор + считается негодным, прогон стартует с last_id=0, причина 'checkpoint_stale'. + - Вердикт пишется в scrape_runs.counters через kit_runs.update_heartbeat + (`counters || :counters` — merge), а не локальный runs_mod.update_heartbeat + (`CAST(:counters AS jsonb)` — полная замена): иначе первый же per-batch heartbeat + после старта стирает resume-вердикт. + +Обратимость (см. PR summary): временный откат last_id на литерал 0 красит +test_resume_continues_from_saved_last_id (last_id == 123456 не совпадает с 0). +""" + +from __future__ import annotations + +import inspect +import json +import os +from typing import Any +from unittest.mock import MagicMock + +os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db") + +from app.services import scheduler as sched + + +class _FakeRow: + def __init__(self, **kw: Any) -> None: + self.__dict__.update(kw) + + +class _FakeDb: + """Двойник сессии: SELECT отдаёт фиксированную строку-кандидата, UPDATE копится. + + Тот же паттерн, что test_930_scheduler_resume_checkpoint.py::_FakeDb — различает + вызовы по наличию ключа 'counters' в params (heartbeat-запись), а не по SQL-тексту. + """ + + 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 is not None and "counters" in params: + self.written.append(params) + return MagicMock() + return MagicMock(fetchone=lambda: self.row) + + def commit(self) -> None: + pass + + +def _candidate(**over: Any) -> _FakeRow: + """Строка-кандидат из _DKP_RESUME_CANDIDATE_SQL: прогон 777, убитый на 20-м батче.""" + base: dict[str, Any] = { + "prev_id": 777, + "prev_status": "zombie", + "prev_counters": {"rows_fetched": 40000, "last_id": 123456, "batches_done": 20}, + "age_h": 2.0, + } + base.update(over) + return _FakeRow(**base) + + +# ── 1. Оборванный прогон продолжает с сохранённого last_id, а не с нуля ───────────── + + +def test_resume_continues_from_saved_last_id() -> None: + """status='zombie' (реап после обрыва) + свежий чекпоинт → last_id прошлого прогона. + + Красный до фикса: import_rosreestr_dkp инициализировал `last_id: int = 0` литералом — + эквивалент этого пути (см. docstring модуля § Обратимость и PR summary с реальным + выводом pytest до/после отката). + """ + db = _FakeDb(_candidate(prev_status="zombie", age_h=2.0)) + last_id, verdict = sched._resume_dkp_cursor(db, run_id=999) + + assert last_id == 123456 + assert verdict["resume_from"] == 777 + assert verdict["resume_reason"] == "ok" + assert verdict["last_id"] == 123456 + + +def test_resume_continues_after_sigterm_drain_despite_done_status() -> None: + """SIGTERM-drain коммитит status='done' (не 'zombie') — резюмится по counters.interrupted.""" + db = _FakeDb( + _candidate( + prev_status="done", + prev_counters={"last_id": 999000, "interrupted": 1}, + age_h=0.5, + ) + ) + last_id, verdict = sched._resume_dkp_cursor(db, run_id=1000) + + assert last_id == 999000 + assert verdict["resume_reason"] == "ok" + + +def test_clean_done_run_is_not_resumed() -> None: + """Полный проход (status='done', БЕЗ interrupted) — резюмить нечего, старт с 0. + + Иначе следующий прогон перестал бы пере-сканировать источник целиком и ON CONFLICT + DO UPDATE перестал бы ловить правки уже импортированных сделок. + """ + db = _FakeDb(_candidate(prev_status="done", prev_counters={"last_id": 5000}, age_h=1.0)) + last_id, verdict = sched._resume_dkp_cursor(db, run_id=1001) + + assert last_id == 0 + assert verdict["resume_reason"] == "status_done" + + +def test_no_previous_run_starts_from_zero() -> None: + """Первый прогон источника вообще — резюмить нечего, явная причина 'no_prev_run'.""" + db = _FakeDb(None) + last_id, verdict = sched._resume_dkp_cursor(db, run_id=1) + + assert last_id == 0 + assert verdict["resume_reason"] == "no_prev_run" + assert verdict["resume_from"] is None + + +# ── 2. Курсор старше потолка не подхватывается — прогон стартует заново ───────────── + + +def test_stale_checkpoint_is_rejected() -> None: + """Чекпоинту больше суток (issue: «курсор недельной давности пропустит правки + апстрима ниже last_id») → прогон стартует с 0, причина видна в вердикте.""" + db = _FakeDb(_candidate(prev_status="running", prev_counters={"last_id": 42}, age_h=24.01)) + last_id, verdict = sched._resume_dkp_cursor(db, run_id=1002) + + assert last_id == 0 + assert verdict["resume_reason"] == "checkpoint_stale" + + +def test_checkpoint_just_under_ceiling_is_still_accepted() -> None: + """Граница потолка (23.99ч) — ещё подхватывается, не off-by-one.""" + db = _FakeDb(_candidate(prev_status="running", prev_counters={"last_id": 42}, age_h=23.99)) + last_id, verdict = sched._resume_dkp_cursor(db, run_id=1003) + + assert last_id == 42 + assert verdict["resume_reason"] == "ok" + + +# ── 3. Запись курсора не затирает посторонние ключи в counters (мерж, не замена) ──── + + +def test_cursor_write_uses_merge_not_replace_heartbeat() -> None: + """import_rosreestr_dkp обязан писать чекпоинт через kit_runs.update_heartbeat + (merge: `counters || :counters`), а не локальный runs_mod.update_heartbeat (замена: + `CAST(:counters AS jsonb)`) — иначе resume-вердикт, записанный ДО цикла, стирается + первым же per-batch heartbeat'ом того же прогона. + """ + src = inspect.getsource(sched.import_rosreestr_dkp) + assert "kit_runs.update_heartbeat" in src + assert "runs_mod.update_heartbeat" not in src + + +def test_kit_runs_update_heartbeat_merges_into_existing_counters() -> None: + """Сама примитива слияния: второй write добавляет rows_fetched/last_id, НЕ стирая + resume_from/resume_reason, записанные первым write'ом. + + SQL-текст реального update_heartbeat (scraper_kit.orchestration.runs) содержит + `||` — проверяем это как контракт, а сам мерж эмулируем на стороне двойника (нет + живого Postgres в юнит-тесте). + """ + from scraper_kit.orchestration import runs as kit_runs + + state: dict[str, Any] = {} + + class _MergeDb: + def execute(self, stmt: Any, params: dict[str, Any]) -> Any: + sql = str(stmt) + assert "||" in sql, "update_heartbeat должен МЕРЖИТЬ (counters || :counters)" + state.update(json.loads(params["counters"])) + return MagicMock() + + def commit(self) -> None: + pass + + db = _MergeDb() + kit_runs.update_heartbeat(db, 1, {"resume_from": 777, "resume_reason": "ok"}) + kit_runs.update_heartbeat(db, 1, {"rows_fetched": 2000, "last_id": 123456}) + + assert state["resume_from"] == 777 # не стёрлось вторым write'ом + assert state["resume_reason"] == "ok" + assert state["rows_fetched"] == 2000 + assert state["last_id"] == 123456 -- 2.45.3