fix(tradein/ingest): rosreestr_dkp_import курсор переживает рестарт (#3168) #3176
2 changed files with 316 additions and 9 deletions
|
|
@ -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",
|
||||
|
|
|
|||
208
tradein-mvp/backend/tests/test_3168_backfill_cursor_resume.py
Normal file
208
tradein-mvp/backend/tests/test_3168_backfill_cursor_resume.py
Normal file
|
|
@ -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
|
||||
Loading…
Add table
Reference in a new issue