gendesign/backend/app/workers/tasks/izyatie_ocr_ingest.py
bot-backend 0792e34172
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
fix(ptica): land_reservation перестаёт копить дубли — 91% таблицы были копиями (#2464)
ВНИМАНИЕ: миграция 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>
2026-08-20 14:45:56 +05:00

225 lines
9.9 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.

"""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"]