gendesign/tradein-mvp/backend/app/services/scheduler.py
bot-backend bf3214b9e4
All checks were successful
CI Trade-In / changes (pull_request) Successful in 8s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / changes (pull_request) Successful in 11s
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 / browser-tests (pull_request) Successful in 1m7s
CI Trade-In / backend-tests (pull_request) Successful in 4m57s
fix(tradein/scrapers): диагноз блока брался из текстовых маркеров чужой площадки, а не из HTTP-статуса (#3196)
Сайдкар вообще не читал код ответа page.goto: страница классифицировалась
только по маркерам, снятым с Авито. Домклик отдаёт статическую `403 | Домклик`
на 26 624 байта, где нет ни одного такого маркера (замер прода 28.08.2026) —
она уезжала наверх как валидный HTML, парсер не находил состояние, и прогон
получал блок неизвестной природы. За 14 дней все 14 прогонов домклика легли с
ban_kind='unknown'; у Яндекса счётчика blocked не было вовсе, поэтому ветка
перевода прогона в 'banned' была недостижима по построению — ноль банов.

- browser/server.py: статус целевой навигации сохраняется per-provider и
  доезжает в тело /fetch аддитивным ключом "status" (ключ "html" не тронут);
  403/429 с маркерами челленджа больше не ждут PoW — ждать нечего, статическая
  страница сама себя не перезагрузит. Наверх идёт BanPageDetectedError, а не
  заглушка: вернув её контентом, воскресили бы #3045.
- scraper_kit/browser_fetcher.py: BrowserFetcher.last_response_status +
  ban_kind_from_status (403/429 → platform, 5xx → infra, прочее → None).
  Поток управления не менялся: fetch() по-прежнему отдаёт str.
- domclick: DomClickBlockedError несёт .status — один тип исключения на
  маркер-детект и на сбой фетча разводится без размножения типов; прогон
  передаёт перепись диагнозов в mark_backfill_finished.
- yandex: появился счётчик blocked, оживляющий ветку бана. Серии блоков и
  промахов парсера считаются РАЗДЕЛЬНО: иначе четыре промаха плюс один 403
  пятым давали 'banned' с переписью {platform: 1}.
- cian: ban_kinds наполняется только диагностируемым статусом. HTTP 200 с
  пустым разбором — дрейф разметки на нашей стороне, а не отказ площадки;
  записав его блоком, мы бы штамповали фиктивные баны у здорового источника
  (13 done против 1 banned за 14 дней).

Инвариант: непустой ban_kinds ⟺ виден ответ 403/429/5xx. Значения остаются в
пределах CHECK scrape_runs.ban_kind.

Известный пробел: шов providers/domclick/detail.py `blocked.status = status`
тестами не покрыт — существующие домкликовые тесты подают исключение готовым
моком и боевой fetch_detail не исполняют.
2026-08-28 23:21:54 +03:00

605 lines
34 KiB
Python
Raw 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 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).
"""
from app.tasks.cian_history_backfill import CianBackfillResult, backfill_cian_history
batch_size = int(params.get("batch_size", 100))
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=True,
do_valuations=False,
on_progress=_heartbeat,
)
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