"""Shared scheduler helpers (#2397 Part C — legacy scheduler_loop removed). История: до #2192 этот модуль нёс полный in-app asyncio scheduler (`scheduler_loop`, 27-веточный `if/elif` dispatch, `trigger_*_run` per-source launchers, `get_due_schedules`, `reap_zombies`, `_claim_run` advisory-lock claim, `_spawn_tracked`/`_drain_inflight` graceful-drain machinery). #2192 представил kit-путь (`scraper_kit.orchestration.scheduler` + `app.services.product_handlers`) как ship-dark alternative; прод переключился на него (`USE_KIT_SCHEDULER=true` в tradein-scraper, см. docker-compose.prod.yml). #2397 Part C убрал legacy `scheduler_loop` и все `trigger_*`/dispatch-функции — kit теперь единственный scheduling-путь (`app/scheduler_main.py` безусловно запускает `_run_kit_scheduler()`). Что осталось в этом модуле — НЕ scheduler-loop, а функции с живыми потребителями вне удалённой machinery: - `compute_next_run_at` — читается admin.py (операторский предпросмотр "next run"). С #2674 это re-export kit-версии, а не вторая копия формулы. - `has_running_run` — читается admin.py (UI-индикатор "уже бежит"). - `import_rosreestr_dkp` — job-тело, вызываемое kit-handler'ом product_handlers._job_rosreestr_dkp (lazy import). - `_execute_cian_backfill` — job-тело, вызываемое kit-handler'ом product_handlers._job_cian_history_backfill (lazy import). Zombie-reap, advisory-lock claim и tick-loop теперь целиком в `scraper_kit/orchestration/scheduler.py` (см. test_scraper_kit_scheduler_parity.py). """ 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 # звал эту копию БЕЗ аргумента, получал default=1 и сбивал любой источник на «завтра». # Копия удалена, а не подправлена: пока формула лежит в двух файлах, следующая правка # такта снова разъедется по одному из них. Re-export (а не правка импорта у вызывающих) # сохраняет `from app.services.scheduler import compute_next_run_at` в admin.py и тестах. from scraper_kit.orchestration.scheduler import compute_next_run_at from sqlalchemy import text from sqlalchemy.orm import Session from app.core.shutdown import shutdown_requested from app.services import scrape_runs as runs_mod __all__ = ["compute_next_run_at", "has_running_run"] logger = logging.getLogger(__name__) # import_rosreestr_dkp: доля per-row INSERT-ошибок (rows_errored / rows_fetched), выше # которой прогон помечается FAILED, а не silent-green (Fix C). Единичные битые строки # (редкий bad row) не валят импорт; систематический сбой (≈100% ошибок) — валит. DKP_IMPORT_ERROR_RATE_THRESHOLD = 0.05 def has_running_run(db: Session, source: str) -> bool: """Есть ли активный run для source (status='running').""" row = db.execute( text( """ SELECT 1 FROM scrape_runs WHERE source = :source AND status = 'running' LIMIT 1 """ ), {"source": source}, ).fetchone() return row is not None async def _execute_cian_backfill( db: Session, *, run_id: int, params: dict[str, Any], ) -> None: """Orchestrate Cian history backfill with heartbeat + checkpoint. Wraps backfill_cian_history(), updating scrape_runs counters (via update_heartbeat) НА КАЖДОЙ сущности батча, а не только до и после него (#2725). Раньше сигнал живости слался ровно один раз — до батча, — а `reap_zombies` меряет именно heartbeat_at с порогом 6 ч, и добивал живые прогоны строго на 6-м часу: 6 прод- прогонов этого источника помечены 'zombie' со сдвигом heartbeat 16-32 мс, при том что у пятерых внутри окна писались строки offer_price_history (у прогона 304 — до 5.4 ч после старта), а штатная длительность источника доходит до 5.06 ч (346). Цена ошибки не косметическая: mark_done апдейтит WHERE status='running', так что после ложной пометки собственный финал прогона становится no-op (отсюда нулевые counters у всех шести), а has_running_run перестаёт видеть прогон и следующий тик может запустить второй такой же батч поверх работающего. Checkpoint/resume semantics: backfill_cian_history() queries rows WHERE history IS NULL via LEFT JOIN — so re-running after a partial completion naturally skips already-processed rows (idempotent by design). Params (from default_params jsonb): batch_size: int — rows per run (listings + houses counted separately). listings_pending: str — "history" (дефолт) | "detail", см. #3284. do_houses: bool — дефолт true; у cian_detail_backfill выключен. """ from app.tasks.cian_history_backfill import CianBackfillResult, backfill_cian_history batch_size = int(params.get("batch_size", 100)) # #3284: одно тело обслуживает ДВА расписания. cian_history_backfill идёт с # дефолтами (история + дома), cian_detail_backfill — с listings_pending="detail" # и do_houses=false: дома у него уже разбирает суточный сосед, а гонять их # круглосуточно незачем. listings_pending = str(params.get("listings_pending", "history")) do_houses = bool(params.get("do_houses", True)) def _counters(result: CianBackfillResult) -> dict[str, int]: return { "listings_processed": result.listings_processed, "listings_succeeded": result.listings_succeeded, "listings_failed": result.listings_failed_fetch + result.listings_failed_save, "houses_processed": result.houses_processed, "houses_succeeded": result.houses_succeeded, "houses_failed": result.houses_failed_fetch + result.houses_failed_save, } def _heartbeat(progress: CianBackfillResult) -> None: """Сигнал живости из середины батча. Best-effort: сбой heartbeat не должен ронять уже идущую работу — прогон в худшем случае вернётся к прежнему поведению (пометка 'zombie' на 6-м часу).""" try: runs_mod.update_heartbeat(db, run_id, _counters(progress)) except Exception: logger.warning( "scheduler: cian_history_backfill run_id=%d heartbeat failed (ignored)", run_id, exc_info=True, ) counters: dict[str, int] = { "listings_processed": 0, "listings_succeeded": 0, "listings_failed": 0, "houses_processed": 0, "houses_succeeded": 0, "houses_failed": 0, } try: runs_mod.update_heartbeat(db, run_id, counters) result = await backfill_cian_history( db, batch_size=batch_size, do_listings=True, do_houses=do_houses, do_valuations=False, on_progress=_heartbeat, listings_pending=listings_pending, ) counters = {**_counters(result), "duration_sec": int(result.duration_sec)} # #3196: отказ detail-фетча теперь несёт диагноз (HTTP-статус последнего ответа # сайдкара). В 'banned' переводим ТОЛЬКО прогон, который отказы видел и не # обогатил НИЧЕГО, — частичный успех остаётся 'done', как и был. if result.ban_kinds and (result.listings_succeeded + result.houses_succeeded) == 0: counters["blocked"] = result.listings_blocked # Полная перепись диагнозов, а не только доминирующий вид (#3196) — иначе # запись прогона теряет, например, единичный infra среди platform. counters["ban_kinds"] = dict(result.ban_kinds) runs_mod.mark_banned( db, run_id, f"cian detail: {result.listings_blocked} отказов, ни одного обогащения", counters, ban_kind=result.ban_kind, ) else: runs_mod.mark_done(db, run_id, counters) logger.info( "scheduler: cian_history_backfill run_id=%d done — listings=%d/%d houses=%d/%d %.1fs", run_id, result.listings_succeeded, result.listings_total, result.houses_succeeded, result.houses_total, result.duration_sec, ) except Exception as exc: logger.exception("scheduler: cian_history_backfill run_id=%d failed", run_id) runs_mod.mark_failed(db, run_id, str(exc)[:1000], counters) 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, params: dict[str, Any], ) -> None: """Import ДКП-сделок из gendesign rosreestr_deals через postgres_fdw. Python-порт import-rosreestr.sh (Variant C из #563). Источник: foreign table gendesign_rosreestr_deals (создана в migration 072). SERVER gendesign_remote настроен в 060_postgres_fdw_extension.sql. USER MAPPING создаётся при startup в core/fdw.py (tradein_fdw_reader). Область покрытия: вся Свердловская область (region_code=66), не только Екатеринбург — прежний ILIKE-фильтр по подстроке города (ограничивавший импорт одним Екатеринбургом) снят (Mera trade-in расширяется на весь регион, unlocks +47183 сделок вне ЕКБ уже сидящих в source foreign table). address и deals.city строятся из реального city источника (не хардкод "Екатеринбург"), deals.region_code заполняется из строки источника (= 66 при текущем фильтре). Фильтры (совпадают с import-rosreestr.sh + Fix_Rosreestr_Dkp_Filter_May24): - region_code = 66 (вся Свердловская область, все города) - city IS NOT NULL AND trim(city) != '' (непустой город → корректный address) - realestate_type_code = '002001003000' (квартира) - area BETWEEN 18 AND 200 - deal_price BETWEEN 1000000 AND 100000000 - street IS NOT NULL AND trim(street) != '' - doc_type = 'ДКП' (только вторичка — #549 / Fix_Rosreestr_Dkp_Filter_May24) - period_start_date >= since (default '2024-01-01') dedup_hash: 'ros:dkp:' || id — плоский натуральный ключ (инъективный, без коллизий, human-readable). До #576 здесь был md5('ros:dkp:' || id); миграция 077 конвертировала существующие строки. source_id хранит исходный rosreestr id (дедуп переустанавливаем). Rooms: выводятся из площади (Росреестр не отдаёт кол-во комнат). Batch-процессинг: читаем из FDW батчами по batch_size через cursor-based пагинацию (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). TODO (follow-up): запустить geocode backfill после import. Cleanup: удаляет legacy строки address='Екатеринбург, реальная сделка' (pre-#549). """ since: str = str(params.get("since", "2024-01-01")) batch_size: int = int(params.get("batch_size", 2000)) counters: dict[str, int] = { "rows_fetched": 0, "rows_inserted": 0, # rows_updated: ON CONFLICT DO UPDATE обновил существующую строку # (исправленный/переопубликованный квартал — Fix D). "rows_updated": 0, # rows_skipped: ТОЛЬКО легитимный dedup-пропуск (строка уже есть, факты # идентичны — DO UPDATE ... WHERE distinct не сработал). "rows_skipped": 0, # rows_errored: реальные per-row INSERT-ошибки, отделены от dedup-skip (Fix C), # раньше обе категории клались в rows_skipped → систематический сбой выглядел # как обычный дедуп и прогон рапортовал success. "rows_errored": 0, "batches_done": 0, } # Cleanup legacy synthetic rows (pre-#549, idempotent) try: deleted = db.execute( text("DELETE FROM deals WHERE address = CAST(:addr AS text) RETURNING id"), {"addr": "Екатеринбург, реальная сделка"}, ).fetchall() db.commit() if deleted: logger.info( "rosreestr_dkp_import run_id=%d: removed %d legacy synthetic rows", run_id, len(deleted), ) except Exception as exc: logger.warning( "rosreestr_dkp_import run_id=%d: cleanup failed (non-fatal): %s", run_id, exc ) db.rollback() 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: cancelled = runs_mod.is_cancelled(db, run_id) if cancelled or shutdown_requested(): if cancelled: # User-cancel: семантика без изменений — mark_cancelled. logger.info("rosreestr_dkp_import run_id=%d: cancelled by user", run_id) runs_mod.mark_cancelled(db, run_id) return # #1182 Phase 2: кооперативный SIGTERM-drain (деплой recreate scraper). # Это НЕ user-cancel → mark_done (partial), не mark_cancelled. Курсор # уже зафиксирован 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", run_id, last_id, total_batches, ) kit_runs.update_heartbeat(db, run_id, counters) runs_mod.mark_done(db, run_id, counters) return # Cursor-based pagination via foreign table gendesign_rosreestr_deals. # FDW pushes WHERE + ORDER BY + LIMIT to gendesign-postgres automatically # (postgres_fdw 'use_remote_estimate' is off by default but clause push-down # still happens for simple predicates on the remote side). batch_rows = ( db.execute( text(""" SELECT id, id AS source_id_src, 'ros:dkp:' || CAST(id AS text) AS dedup_hash, trim(city) || ', ' || trim(street) AS address, region_code, trim(city) AS city, CASE WHEN area < 30 THEN 0 WHEN area < 44 THEN 1 WHEN area < 62 THEN 2 WHEN area < 85 THEN 3 ELSE 4 END AS rooms, round(area, 2) AS area_m2, LEAST( NULLIF( -- #1525: '-?[0-9]+' сохраняет знак минус, иначе цоколь/подвал -- '-1' импортируется как 1. Для '5/9' по-прежнему берётся 5. substring(floor FROM '-?[0-9]+'), '' )::int, 100 ) AS floor_num, year_build AS year_built, round(deal_price)::bigint AS price_rub, round(price_per_sqm)::int AS price_per_m2, period_start_date AS deal_date FROM gendesign_rosreestr_deals WHERE region_code = 66 AND city IS NOT NULL AND trim(city) <> '' AND realestate_type_code = '002001003000' AND area BETWEEN 18 AND 200 AND deal_price BETWEEN 1000000 AND 100000000 AND street IS NOT NULL AND trim(street) <> '' AND doc_type = 'ДКП' AND period_start_date >= CAST(:since AS date) AND id > CAST(:last_id AS bigint) ORDER BY id LIMIT CAST(:batch_size AS int) """), {"since": since, "last_id": last_id, "batch_size": batch_size}, ) .mappings() .all() ) if not batch_rows: break total_batches += 1 batch_inserted = 0 batch_updated = 0 batch_skipped = 0 batch_errored = 0 batch_max_id = last_id for row in batch_rows: row_id: int = int(row["id"]) if row_id > batch_max_id: batch_max_id = row_id try: with db.begin_nested(): # SAVEPOINT per row # Fix D: ON CONFLICT DO UPDATE вместо прежнего no-op-дедупа — при # повторном импорте ИСПРАВЛЕННОГО/переопубликованного квартала обновляем # сырые факты Росреестра. WHERE ... IS DISTINCT FROM оставляет # неизменные строки нетронутыми (идемпотентность resume: строка # без изменений → 0 returned → legit dedup-skip). source/source_id # (identity) и dedup_hash (ключ конфликта) стабильны, не трогаем. # Обогащение (lat/lon/geom/geocode_tried_at, cadastral_number, # total_floors, house_type ...) НЕ в EXCLUDED-списке → сохраняется. # RETURNING (xmax = 0): freshly-inserted → xmax=0 (was_inserted), # обновлённая по ON CONFLICT → xmax<>0; отличаем insert от update. result = db.execute( text(""" INSERT INTO deals ( source, dedup_hash, source_id, address, region_code, city, rooms, area_m2, floor, year_built, price_rub, price_per_m2, deal_date ) VALUES ( 'rosreestr', CAST(:dedup_hash AS text), CAST(:source_id AS text), CAST(:address AS text), CAST(:region_code AS int), CAST(:city AS text), CAST(:rooms AS int), CAST(:area_m2 AS numeric), CAST(:floor_num AS int), CAST(:year_built AS int), CAST(:price_rub AS bigint), CAST(:price_per_m2 AS int), CAST(:deal_date AS date) ) ON CONFLICT (dedup_hash) DO UPDATE SET address = EXCLUDED.address, region_code = EXCLUDED.region_code, city = EXCLUDED.city, rooms = EXCLUDED.rooms, area_m2 = EXCLUDED.area_m2, floor = EXCLUDED.floor, year_built = EXCLUDED.year_built, price_rub = EXCLUDED.price_rub, price_per_m2 = EXCLUDED.price_per_m2, deal_date = EXCLUDED.deal_date WHERE deals.address IS DISTINCT FROM EXCLUDED.address OR deals.region_code IS DISTINCT FROM EXCLUDED.region_code OR deals.city IS DISTINCT FROM EXCLUDED.city OR deals.rooms IS DISTINCT FROM EXCLUDED.rooms OR deals.area_m2 IS DISTINCT FROM EXCLUDED.area_m2 OR deals.floor IS DISTINCT FROM EXCLUDED.floor OR deals.year_built IS DISTINCT FROM EXCLUDED.year_built OR deals.price_rub IS DISTINCT FROM EXCLUDED.price_rub OR deals.price_per_m2 IS DISTINCT FROM EXCLUDED.price_per_m2 OR deals.deal_date IS DISTINCT FROM EXCLUDED.deal_date RETURNING (xmax = 0) AS was_inserted """), { "dedup_hash": row["dedup_hash"], "source_id": str(row["source_id_src"]), "address": row["address"], "region_code": row["region_code"], "city": row["city"], "rooms": row["rooms"], "area_m2": row["area_m2"], "floor_num": row["floor_num"], "year_built": row["year_built"], "price_rub": row["price_rub"], "price_per_m2": row["price_per_m2"], "deal_date": row["deal_date"], }, ).fetchone() if result is None: # DO UPDATE ... WHERE distinct не сработал → строка есть и # факты идентичны = легитимный dedup-skip (Fix C: только это # теперь считается skip, INSERT-ошибки — отдельно ниже). batch_skipped += 1 elif result[0]: batch_inserted += 1 else: batch_updated += 1 except Exception as exc: # Fix C: реальная per-row INSERT-ошибка — ОТДЕЛЬНЫЙ счётчик, не skip. logger.warning( "rosreestr_dkp_import run_id=%d: row id=%d INSERT failed: %s", run_id, row_id, exc, ) batch_errored += 1 last_id = batch_max_id db.commit() counters["rows_fetched"] += len(batch_rows) counters["rows_inserted"] += batch_inserted counters["rows_updated"] += batch_updated counters["rows_skipped"] += batch_skipped counters["rows_errored"] += batch_errored counters["batches_done"] = total_batches counters["last_id"] = last_id # type: ignore[assignment] # 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", run_id, total_batches, len(batch_rows), batch_inserted, batch_updated, batch_skipped, batch_errored, last_id, ) if len(batch_rows) < batch_size: break # Last partial batch — no more rows # Fix C: систематический per-row INSERT-сбой больше не рапортует success. # Отделив rows_errored от rows_skipped, проверяем долю ошибок: выше порога — # прогон FAILED (raise → внешний except → mark_failed), а не silent-green. fetched = counters["rows_fetched"] errored = counters["rows_errored"] if fetched > 0 and errored / fetched > DKP_IMPORT_ERROR_RATE_THRESHOLD: raise RuntimeError( f"rosreestr_dkp_import: per-row INSERT error rate " f"{errored}/{fetched} ({errored / fetched:.1%}) exceeds " f"{DKP_IMPORT_ERROR_RATE_THRESHOLD:.0%} threshold — marking run failed" ) runs_mod.mark_done(db, run_id, counters) logger.info( "rosreestr_dkp_import run_id=%d done: " "total_fetched=%d inserted=%d updated=%d skipped=%d errored=%d batches=%d", run_id, counters["rows_fetched"], counters["rows_inserted"], counters["rows_updated"], counters["rows_skipped"], counters["rows_errored"], total_batches, ) except Exception as exc: logger.exception("rosreestr_dkp_import run_id=%d failed at last_id=%d", run_id, last_id) runs_mod.mark_failed(db, run_id, str(exc)[:1000], counters) raise