gendesign/tradein-mvp/backend/app/services/scheduler.py
bot-backend fcf5887225
All checks were successful
CI Trade-In / changes (pull_request) Successful in 7s
CI / changes (pull_request) Successful in 9s
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 5m12s
merge(#3051): main (#3421) в ветку импорта по региону — московская дельта поверх region_code/doc_type
#3421 въехал в main параллельно с той же миграцией 288 (deals.doc_type,
параметры region_code/doc_types). Разрешение: 288 — целиком версия main;
наша дельта (FDW-колонки okato/quarter_cad_number/district, выключенный seed
rosreestr_dkp_import_77) переехала в 289. scheduler.py — doc_types из main +
canonical_city-маппинг/raw_payload/per-source чекпоинт. deploy-скрипт —
валидация REGION_CODE и DOC_TYPE (интерполируются в SQL текстом).
2026-09-08 23:19:47 +03:00

781 lines
48 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""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 json
import logging
from typing import Any
# kit_runs — ТОТ ЖЕ модуль, что и runs_mod ниже: с #3390 `app.services.scrape_runs`
# его алиас, реализация одна и counters везде МЕРЖАТСЯ (`counters || :counters`). До
# #3390 копии было две, и app-копия counters ЗАМЕНЯЛА — тогда чекпоинт курсора
# import_rosreestr_dkp (#3168) обязан был писаться именно kit-именем, иначе resume-вердикт
# со старта затирался первым же per-batch пульсом. Имя оставлено как есть: теперь это
# один объект, и переименование в runs_mod ничего не чинит и ничего не ломает.
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 scraper_kit.proxy_errors import caused_by_no_proxy
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
from app.services.regions import REGIONS
__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-м часу)."""
nonlocal counters
# #3384: снимок измеренного едет не только в БД, но и в `counters` — этот словарь
# уезжает в mark_failed из общего except ниже. Мерж (#3390) спасает лишь ключи,
# которых в payload нет; одноимённые он ПЕРЕЗАПИСЫВАЕТ, поэтому предынициализированные
# нули без этого присваивания легли бы поверх измеренного, и SQL-разбор простоя
# (#3288/#3367) прочитал бы «к площадке не ходили» про прогон, который ходил.
# Присваивание ДО записи в БД: сбой heartbeat'а не должен стирать сам факт замера.
counters = _counters(progress)
try:
runs_mod.update_heartbeat(db, run_id, counters)
except Exception:
logger.warning(
"scheduler: cian_history_backfill run_id=%d heartbeat failed (ignored)",
run_id,
exc_info=True,
)
# Стартовые нули: прогон виден в админке до первого прогресса. Держатся здесь ровно
# до первого `_heartbeat` — дальше в `counters` лежит измеренное (см. выше).
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.no_proxy_stop:
# #3197 (как #3288 у avito / #3283 у домклика): остановка из-за пустого пула —
# НЕ блок, поэтому и не mark_banned: иначе прогон уйдёт в 'banned' и запись
# будет утверждать про площадку то, чего не было. Это отказ нашей стороны.
counters["no_proxy_stop"] = 1
runs_mod.mark_failed(
db, run_id, "пул прокси пуст — к площадке не ходили (#3197)", counters
)
# INFO, как у соседей (avito_detail_backfill.py:1050, domclick:626): причина
# уже записана в mark_failed + counters.no_proxy_stop, а ERROR на финальной
# строке ставил в один разряд с падением задачи (logger.exception ниже).
logger.info(
"scheduler: cian_history_backfill run_id=%d СТОП (пул пуст) — "
"listings=%d/%d houses=%d/%d %.1fs",
run_id,
result.listings_succeeded,
result.listings_total,
result.houses_succeeded,
result.houses_total,
result.duration_sec,
)
return
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)
if caused_by_no_proxy(exc):
# #3384: пул был пуст ещё ДО первого объявления — lease берётся в
# BrowserFetcher.__aenter__, поэтому NoProxyAvailableError вылетает из
# самого `async with` (cian_history_backfill.py:218) мимо стоп-механики
# внутри цикла, которая и ставит no_proxy_stop. Без этого ключа прогон,
# который к площадке не ходил ВООБЩЕ, неотличим от любого другого падения:
# причина только в тексте, а разбор простоя идёт SQL'ём по
# counters.no_proxy_stop (#3288/#3367). Ключ тот же, что у ветки выше.
# Флаг ДОБАВЛЯЕТСЯ к последнему снимку `_heartbeat`, а не подменяет его:
# опустевший между стадиями пул — это отказ ПОСЛЕ реальной работы.
counters["no_proxy_stop"] = 1
runs_mod.mark_failed(db, run_id, str(exc)[:1000], counters)
raise
_DKP_SOURCE = "rosreestr_dkp_import"
def _dkp_source_for_region(region_code: int) -> str:
"""Имя scrape_runs.source для чекпоинта данного региона (#3051 п.3).
66 — байт-в-байт прежнее имя ('rosreestr_dkp_import'), под которым годами
писались scrape_runs. Остальные регионы получают суффикс кода — тот же
формат, что и у строки scrape_schedules ('rosreestr_dkp_import_77',
seed — миграция 289), которую резолвит wildcard 'rosreestr_dkp_import_*'
в product_handlers.py. Изоляция чекпоинтов между регионами держится именно
на разных source: _resume_dkp_cursor ищет ПРЕДЫДУЩИЙ прогон с ТЕМ ЖЕ source,
поэтому курсор региона 77 никогда не подхватит last_id региона 66 (и
наоборот) — они просто разные строки в scrape_runs.source.
"""
if region_code == 66:
return _DKP_SOURCE
return f"{_DKP_SOURCE}_{region_code}"
# Потолок возраста чекпоинта: старше — 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, source: str = _DKP_SOURCE
) -> tuple[int, dict[str, Any]]:
"""Продолжить last_id прошлого прогона или начать с 0 — решение + explain (issue #3168).
last_id раньше жил только в памяти процесса (init на 0 при каждом запуске): heartbeat
писал его в counters КАЖДЫЙ батч (checkpoint), но на старте никто это не читал
обратно. Обрыв (деплой/OOM/рестарт хоста) откатывал прогресс на 0 и заставлял
пере-сканировать источник с начала.
`source` (#3051 п.3) — per-region ключ чекпоинта (см. _dkp_source_for_region):
дефолт _DKP_SOURCE сохраняет прежнее поведение вызовов без явного аргумента
(регион 66). Кандидат ищется СТРОГО по этому source — прогон региона 77
(source='rosreestr_dkp_import_77') никогда не видит last_id региона 66
(source='rosreestr_dkp_import') и наоборот: разные регионы физически не
матчат друг друга в WHERE source = :source ниже.
Кандидат — ПОСЛЕДНИЙ прогон ЭТОГО 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": 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). #3051 п.3: параметризовано
по региону (params["region_code"], реестр — app.services.regions.REGIONS) —
было хардкод region_code=66.
Источник: foreign table gendesign_rosreestr_deals (создана в migration 072,
okato/quarter_cad_number/district добавлены миграцией 289).
SERVER gendesign_remote настроен в 060_postgres_fdw_extension.sql.
USER MAPPING создаётся при startup в core/fdw.py (tradein_fdw_reader).
Область покрытия задаётся параметром, а не литералом (#3051 п.6): region_code
приходит из params, дефолт 66 = вся Свердловская область (не только Екатеринбург —
прежний ILIKE-фильтр по подстроке города снят, unlocks +47183 сделок вне ЕКБ уже
сидящих в source foreign table). Неизвестный код региона (нет в REGIONS) — ValueError,
прогон падает явно, а не молча импортирует мусор с чужим region_code.
region_code=66 (регион БЕЗ canonical_city в реестре) — поведение байт-в-байт
прежнее: city/address строятся из city источника, обязателен фильтр
city IS NOT NULL AND trim(city) != ''.
Регион С canonical_city (77 — Москва): Росреестр отдаёт в city муниципальный
округ/поселение ("муниципальный округ Раменки", "поселение Сосенское"), НЕ
город — city/address подставляют region.canonical_city, а не city источника;
фильтр city IS NOT NULL НЕ применяется (иначе теряется ~10% строк с пустым
city источника). Исходные city/okato/quarter_cad_number/district уходят в
raw_payload (jsonb) — единственная ветка SQL решает это через bind-параметр
:canonical_city (CASE WHEN ... IS NOT NULL), а не отдельный Python if/else на
конкретный код региона.
Типы документов тоже параметр (#3051 п.3): doc_types, дефолт ['ДКП'] = прежнее
поведение (только вторичка — #549 / Fix_Rosreestr_Dkp_Filter_May24). Для Москвы
ДДУ идут по ценам котлована и медиану развалят, поэтому смешивать их с ДКП можно
только осознанно и с колонкой deals.doc_type (миграция 288), которая теперь
заполняется на импорте.
Фильтры (совпадают с import-rosreestr.sh + Fix_Rosreestr_Dkp_Filter_May24):
- region_code = :region_code (параметризовано, было хардкод 66)
- city IS NOT NULL AND trim(city) != '' — ТОЛЬКО если у региона нет canonical_city
- realestate_type_code = '002001003000' (квартира)
- area BETWEEN 18 AND 200
- deal_price BETWEEN 1000000 AND 100000000
- street IS NOT NULL AND trim(street) != ''
- doc_type = ANY(:doc_types) (param, default ['ДКП'])
- 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 (дедуп переустанавливаем).
Префикс ':dkp:' НАМЕРЕННО оставлен неизменным после параметризации doc_types: id
уникален в источнике сам по себе, независимо от типа документа, поэтому ключ и без
того не коллизирует; а вот смена формы ключа осиротила бы все уже загруженные строки
(их пришлось бы конвертировать ещё одной миграцией — ровно то, что делала 077).
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). Курсор — ПЕР
РЕГИОН (#3051 п.3): _resume_dkp_cursor вызывается с source=_dkp_source_for_region
(region_code), поэтому last_id региона 77 никогда не подхватывает last_id региона
66 — они разные scrape_runs.source ('rosreestr_dkp_import' vs
'rosreestr_dkp_import_77'), см. докстринг _dkp_source_for_region.
SAVEPOINT per row — один сбойный row не откатывает батч.
Координаты: NULL после импорта — геокодинг остаётся follow-up (geocode-deals).
TODO (follow-up): запустить geocode backfill после import.
Cleanup: удаляет legacy строки address='Екатеринбург, реальная сделка' (pre-#549,
region-agnostic — синтетические строки существовали только для ЕКБ).
"""
since: str = str(params.get("since", "2024-01-01"))
batch_size: int = int(params.get("batch_size", 2000))
# #3051 п.6: регион — параметр, дефолт 66 сохраняет текущее прод-поведение
# (расписание получает явный region_code в миграции 288).
region_code: int = int(params.get("region_code", 66))
# #3051 п.3: типы документов — параметр, дефолт ['ДКП'] = прежний литерал.
doc_types: list[str] = [str(t) for t in params.get("doc_types") or ["ДКП"]]
region = REGIONS.get(region_code)
if region is None:
raise ValueError(
f"rosreestr_dkp_import: region_code={region_code} не найден в "
f"app.services.regions.REGIONS (известны: {sorted(REGIONS)}) — "
"прогон остановлен, чтобы не импортировать сделки с неизвестным "
"региональным контекстом (city/address-правила для него не определены)"
)
dkp_source = _dkp_source_for_region(region_code)
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, source=dkp_source)
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,
-- #3051: регион с canonical_city (Москва) подставляет его вместо
-- city источника (округ/поселение, не город) — CASE на bind-параметре,
-- не Python if/else на код региона.
CASE
WHEN CAST(:canonical_city AS text) IS NOT NULL
THEN CAST(:canonical_city AS text) || ', ' || trim(street)
ELSE trim(city) || ', ' || trim(street)
END AS address,
region_code,
CASE
WHEN CAST(:canonical_city AS text) IS NOT NULL
THEN CAST(:canonical_city AS text)
ELSE trim(city)
END 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,
doc_type,
-- Исходный city/okato/quarter_cad_number/district — ТОЛЬКО когда
-- city перезаписан canonical_city выше (иначе NULL, регион 66
-- byte-for-byte прежний: raw_payload не заполнялся и не заполняется).
CASE
WHEN CAST(:canonical_city AS text) IS NOT NULL THEN
jsonb_build_object(
'src_city', city,
'okato', okato,
'quarter_cad_number', quarter_cad_number,
'district', district
)
ELSE NULL
END AS raw_payload
FROM gendesign_rosreestr_deals
WHERE region_code = CAST(:region_code AS int)
AND (
CAST(:canonical_city AS text) IS NOT NULL
OR (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 = ANY(CAST(:doc_types AS text[]))
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,
"region_code": region_code,
"canonical_city": region.canonical_city,
"doc_types": doc_types,
},
)
.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, doc_type, raw_payload
)
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),
CAST(:doc_type AS text),
CAST(:raw_payload AS jsonb)
)
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,
doc_type = EXCLUDED.doc_type,
raw_payload = EXCLUDED.raw_payload
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
OR deals.doc_type IS DISTINCT FROM EXCLUDED.doc_type
OR deals.raw_payload IS DISTINCT FROM EXCLUDED.raw_payload
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"],
"doc_type": row["doc_type"],
"raw_payload": (
json.dumps(row["raw_payload"], ensure_ascii=False)
if row["raw_payload"] is not None
else None
),
},
).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.
# Пульс МЕРЖИТ counters (`counters || :counters`), поэтому resume_verdict,
# записанный _resume_dkp_cursor'ом перед циклом, переживает per-batch запись
# (issue #3168; с #3390 мерж — единственная семантика, см. scrape_runs).
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