Some checks failed
CI Trade-In / changes (pull_request) Successful in 15s
CI / changes (pull_request) Failing after 14s
CI Trade-In / backend-tests (pull_request) Has been skipped
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
ВНИМАНИЕ: миграция 189 УДАЛЯЕТ строки на проде (см. «Что удаляется»). _UPSERT_NO_ACT_SQL заканчивается ON CONFLICT DO NOTHING, а единственный подходящий констрейнт — uq_land_reservation_cad_act UNIQUE (cad_num, act_number) с обычной NULL-семантикой. В Postgres NULL != NULL, поэтому у записей БЕЗ номера акта конфликт не наступает никогда: DO NOTHING не срабатывает, каждый недельный прогон вставляет копию. То есть ON CONFLICT здесь был декорацией. Замер прода 20.08.2026: строк всего 297 из них act_number IS NULL 297 (все) групп (cad_num, doc_url) с дублями 27 максимум копий одной записи 11 лишних строк 270 (91% таблицы) Что удаляется: копии сверх первой (минимальный id) в каждой группе. Это порождение бага, а не пользовательские данные: таблица — кэш OCR-разбора PDF с сайта, пересобираемый прогоном таски. Проверено, что ключ подходит: ни у одного cad_num нет более одного doc_url (max = 1), дубли внутри групп — точные копии. Прецеденты NULLS NOT DISTINCT в репо: м.110, м.125, м.140, м.158. Prod = PG16.4. Документация приведена к реальности. Docstring обещал python-дедуп по (cad_num, doc_url) и «двухшаговый UPSERT ниже» — ни того, ни другого в коде не было. Комментарий у варианта B был честнее, но его оценка «rare, data audit OK» не подтвердилась: 91% таблицы. Отложенный там вариант (уникальный индекс на NULL) и реализован этой миграцией. Тест репетирует миграцию на ВРЕМЕННОЙ копии, засеянной как прод (11 копий одной записи + соседний участок): исполняет РЕАЛЬНЫЕ выражения из файла, проверяет 12 строк → 2 и что повторная вставка стала no-op. Плюс фальсификация: со СТАРЫМ констрейнтом дубль обязан появиться — без неё зелёный тест неотличим от «оно и так работало». Плюс два контроля: записи С номером акта дедуплицировались и раньше, разные участки не схлопываются. Первая версия репетиции упала на ADD CONSTRAINT — я разбивал миграцию по «;» и молча выбрасывал куски, начинающиеся с комментария, вместе с DELETE. Это ровно то, что репетиция и должна ловить; разбор исправлен. Прогоны: tests/sql (живой Postgres) 6 passed rc=0; -k "izyatie or reservation" 76 passed rc=0. Шесть nodeid в skip_allowlist.txt — нужен Postgres, в CI идут. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
225 lines
9.9 KiB
Python
225 lines
9.9 KiB
Python
"""Celery task: OCR-пайплайн изъятия ЕКБ → land_reservation (#1062).
|
||
|
||
Загружает PDF-сканы «Сообщений о планируемом изъятии…» с раздела екатеринбург.рф,
|
||
выполняет OCR через Tesseract rus, извлекает кад-номера 66:41:NNNNNNN:NN,
|
||
UPSERT-ит в land_reservation (м.136). Reservation_lookup / analyze-wiring (#1118)
|
||
подхватывает данные автоматически — дополнительного wire-up не требуется.
|
||
|
||
Дедуп-ключ:
|
||
ON CONFLICT (cad_num, act_number) — унаследован из reservation_ingest.py.
|
||
Уникальность держит констрейнт uq_land_reservation_cad_act; с миграции 189 он
|
||
объявлен как UNIQUE NULLS NOT DISTINCT, поэтому записи без номера акта тоже
|
||
конфликтуют между собой и ON CONFLICT DO NOTHING реально их ловит.
|
||
|
||
До м.189 констрейнт был обычным UNIQUE, где NULL != NULL: у записей с
|
||
act_number IS NULL конфликт не наступал никогда, и каждый недельный прогон
|
||
вставлял копию. Замер прода 20.08.2026 до правки — 297 строк, все без номера
|
||
акта, 27 групп с дублями, до 11 копий, 270 лишних строк (91% таблицы).
|
||
|
||
Прежняя редакция этого docstring обещала python-дедуп по (cad_num, doc_url)
|
||
перед UPSERT и «двухшаговый UPSERT ниже». Ни того, ни другого в коде не было —
|
||
описание расходилось с реализацией и скрывало накопление дублей (#2464).
|
||
|
||
Beat: еженедельно (пятница 07:00 МСК) — изъятия выходят редко.
|
||
|
||
Conventions (зеркало reservation_ingest.py / backend.md):
|
||
• SessionLocal() + try/finally close.
|
||
• SAVEPOINT per-row (with db.begin_nested()).
|
||
• CAST(:x AS type) — НИКОГДА :x::type (psycopg v3).
|
||
• httpx, не requests.
|
||
• Сбой одного документа не валит прогон (try/except с logger.error).
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import logging
|
||
from typing import Any
|
||
|
||
from sqlalchemy import text
|
||
from sqlalchemy.orm import Session
|
||
|
||
from app.core.db import SessionLocal
|
||
from app.workers.celery_app import celery_app
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
# UPSERT — два варианта в зависимости от наличия act_number.
|
||
#
|
||
# Вариант A (act_number IS NOT NULL): ON CONFLICT (cad_num, act_number) → DO UPDATE.
|
||
# Stable key = (cad_num, act_number). Идемпотентно при повторном прогоне.
|
||
#
|
||
# Вариант B (act_number IS NULL): INSERT ... ON CONFLICT DO NOTHING.
|
||
# Работает с миграции 189: uq_land_reservation_cad_act объявлен как
|
||
# UNIQUE NULLS NOT DISTINCT, поэтому (cad_num, NULL) конфликтует с такой же
|
||
# строкой и повторный прогон становится no-op.
|
||
# Прежний комментарий здесь оценивал накопление дублей как «rare, data audit OK»
|
||
# и откладывал уникальный индекс. Оценка не подтвердилась: на 20.08.2026 дубли
|
||
# составляли 91% таблицы (270 лишних строк из 297), максимум 11 копий одной
|
||
# записи. Отложенный вариант и реализован м.189 (#2464).
|
||
|
||
_UPSERT_WITH_ACT_SQL = text(
|
||
"""
|
||
INSERT INTO land_reservation (
|
||
cad_num, reservation_kind, basis_act, act_number, act_date,
|
||
purpose, doc_url, source, is_active, raw_excerpt, fetched_at
|
||
) VALUES (
|
||
CAST(:cad_num AS text),
|
||
CAST(:reservation_kind AS text),
|
||
CAST(:basis_act AS text),
|
||
CAST(:act_number AS text),
|
||
CAST(:act_date AS date),
|
||
CAST(:purpose AS text),
|
||
CAST(:doc_url AS text),
|
||
CAST(:source AS text),
|
||
true,
|
||
CAST(:raw_excerpt AS text),
|
||
now()
|
||
)
|
||
ON CONFLICT (cad_num, act_number) DO UPDATE SET
|
||
reservation_kind = EXCLUDED.reservation_kind,
|
||
basis_act = EXCLUDED.basis_act,
|
||
act_date = EXCLUDED.act_date,
|
||
purpose = EXCLUDED.purpose,
|
||
doc_url = EXCLUDED.doc_url,
|
||
source = EXCLUDED.source,
|
||
raw_excerpt = COALESCE(EXCLUDED.raw_excerpt, land_reservation.raw_excerpt),
|
||
fetched_at = now()
|
||
"""
|
||
)
|
||
|
||
_UPSERT_NO_ACT_SQL = text(
|
||
"""
|
||
INSERT INTO land_reservation (
|
||
cad_num, reservation_kind, basis_act, act_number, act_date,
|
||
purpose, doc_url, source, is_active, raw_excerpt, fetched_at
|
||
) VALUES (
|
||
CAST(:cad_num AS text),
|
||
CAST(:reservation_kind AS text),
|
||
CAST(:basis_act AS text),
|
||
NULL,
|
||
CAST(:act_date AS date),
|
||
CAST(:purpose AS text),
|
||
CAST(:doc_url AS text),
|
||
CAST(:source AS text),
|
||
true,
|
||
CAST(:raw_excerpt AS text),
|
||
now()
|
||
)
|
||
ON CONFLICT DO NOTHING
|
||
"""
|
||
)
|
||
|
||
|
||
def _upsert_records(db: Session, records: list[dict[str, Any]]) -> int:
|
||
"""UPSERT записей в land_reservation. Возвращает число успешно обработанных строк."""
|
||
count = 0
|
||
for row in records:
|
||
upsert_sql = _UPSERT_WITH_ACT_SQL if row.get("act_number") else _UPSERT_NO_ACT_SQL
|
||
try:
|
||
with db.begin_nested(): # SAVEPOINT per-row
|
||
db.execute(upsert_sql, row)
|
||
count += 1
|
||
except Exception as exc:
|
||
logger.warning(
|
||
"_upsert_records: пропуск cad_num=%s doc_url=%s: %s",
|
||
row.get("cad_num"),
|
||
row.get("doc_url"),
|
||
exc,
|
||
)
|
||
return count
|
||
|
||
|
||
@celery_app.task(name="tasks.izyatie_ocr_ingest.ingest_izyatie_ocr")
|
||
def ingest_izyatie_ocr() -> dict[str, int]:
|
||
"""OCR-пайплайн изъятия ЕКБ: раздел сайта → PDF → Tesseract → land_reservation.
|
||
|
||
Returns:
|
||
{'docs': N, 'parcels': M} — обработано документов и вставлено/обновлено ЗУ.
|
||
"""
|
||
# Ленивые импорты — модули не нужны при старте celery_app.
|
||
from app.services.scrapers.izyatie_client import fetch_pdf, list_izyatie_documents
|
||
from app.services.scrapers.izyatie_ocr import extract_izyatie_records, ocr_pdf_text
|
||
|
||
docs = list_izyatie_documents()
|
||
if not docs:
|
||
logger.warning("ingest_izyatie_ocr: документов не найдено — прогон завершён без данных")
|
||
return {"docs": 0, "parcels": 0}
|
||
|
||
logger.info("ingest_izyatie_ocr: запуск пайплайна, %d документов", len(docs))
|
||
|
||
total_parcels = 0
|
||
processed_docs = 0
|
||
db: Session = SessionLocal()
|
||
|
||
try:
|
||
for doc in docs:
|
||
doc_title: str = doc.get("title", "")
|
||
doc_url: str = doc.get("url", "")
|
||
|
||
logger.info("ingest_izyatie_ocr: загружаем %r → %s", doc_title[:60], doc_url)
|
||
|
||
# Весь per-document body (fetch → OCR → extract → upsert → commit) — в
|
||
# одном try/except. Сбой ЛЮБОГО шага (битый PDF на fetch, Tesseract-crash
|
||
# на OCR, баг парсера на неожиданном OCR-тексте, ошибка upsert/commit)
|
||
# логируется и переходит к следующему документу вместо падения всей
|
||
# еженедельной пачки (issue #2445 D5). db.rollback() откатывает
|
||
# незакоммиченные изменения текущей итерации, чтобы они не отравили
|
||
# commit следующего документа.
|
||
try:
|
||
pdf_bytes = fetch_pdf(doc_url)
|
||
|
||
ocr_text = ocr_pdf_text(pdf_bytes)
|
||
if not ocr_text.strip():
|
||
logger.warning(
|
||
"ingest_izyatie_ocr: OCR вернул пустой текст для %s "
|
||
"(tesseract не установлен, битый PDF или чистый лист?)",
|
||
doc_url,
|
||
)
|
||
processed_docs += 1
|
||
continue
|
||
|
||
records = extract_izyatie_records(ocr_text, doc_title, doc_url)
|
||
logger.info(
|
||
"ingest_izyatie_ocr: %r → %d кад-номеров",
|
||
doc_title[:60],
|
||
len(records),
|
||
)
|
||
|
||
if not records:
|
||
processed_docs += 1
|
||
continue
|
||
|
||
count = _upsert_records(db, records)
|
||
db.commit()
|
||
total_parcels += count
|
||
processed_docs += 1
|
||
|
||
except Exception as exc:
|
||
logger.warning(
|
||
"ingest_izyatie_ocr: пропуск документа %r (%s): %s",
|
||
doc_title[:60],
|
||
doc_url,
|
||
exc,
|
||
)
|
||
try:
|
||
db.rollback()
|
||
except Exception as rollback_exc:
|
||
logger.warning(
|
||
"ingest_izyatie_ocr: rollback после сбоя документа %s тоже упал: %s",
|
||
doc_url,
|
||
rollback_exc,
|
||
)
|
||
continue
|
||
|
||
finally:
|
||
db.close()
|
||
|
||
logger.info(
|
||
"ingest_izyatie_ocr: завершено docs=%d parcels=%d",
|
||
processed_docs,
|
||
total_parcels,
|
||
)
|
||
return {"docs": processed_docs, "parcels": total_parcels}
|
||
|
||
|
||
__all__ = ["ingest_izyatie_ocr"]
|