fix(tradein/ingest): rosreestr_dkp_import курсор переживает рестарт (#3168)
All checks were successful
CI Trade-In / changes (pull_request) Successful in 10s
CI / changes (pull_request) Successful in 11s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 4m45s

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 тестов).
This commit is contained in:
bot-backend 2026-08-28 01:05:54 +03:00
parent bdb9b64b03
commit 51027d8b02
2 changed files with 316 additions and 9 deletions

View file

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

View 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