feat: конвейер ingest→SQLite, API на БД, фронт сравнения цен

- etl/pipeline.py — extract→normalize→validate→SQLite, дедуп партнёров, версии цен, отчёт качества
- api/main.py — 7 эндпоинтов на SQLite; ядро /services/{id}/partners (сравнение)
- web/index.html — одностраничный фронт (палитра Nomad): поиск, сравнение, дашборд, unmatched
- загрузчик справочника: уникальный service_id

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
2026-06-26 17:30:10 +05:00
parent 1de66eaa85
commit 985faadf36
6 changed files with 438 additions and 24 deletions
+9 -1
View File
@@ -64,6 +64,7 @@ def load_dictionary(path: str | Path) -> list[Service]:
raise ValueError(f"В справочнике нет колонки Name_ru. Заголовки: {header}")
services: list[Service] = []
seen_ids: dict[str, int] = {}
for row in rows:
raw_name = row[name_idx]
if not raw_name or not str(raw_name).strip():
@@ -74,8 +75,15 @@ def load_dictionary(path: str | Path) -> list[Service]:
if tar and not TARIFICATOR_RE.match(tar):
tar = None # отбрасываем мусорные коды, оставляем только валидный формат
# «<ID>-<Code>» иногда повторяется в справочнике — гарантируем уникальность.
sid = f"{row[col['ID']]}-{row[col['Code']]}"
repeat = seen_ids.get(sid, 0)
seen_ids[sid] = repeat + 1
if repeat:
sid = f"{sid}#{repeat}"
service = Service(
service_id=f"{row[col['ID']]}-{row[col['Code']]}",
service_id=sid,
specialty=str(row[col["Специальность"]]).strip() if row[col["Специальность"]] else "",
name_ru=str(raw_name).strip(),
tarificator_code=tar,
+7 -3
View File
@@ -18,8 +18,12 @@ from etl.extractors import readers # noqa: E402
from etl.extractors.grid import extract_rows_from_grid # noqa: E402
def extract(path: str) -> ParsedDocument:
"""Разобрать один файл прайса в структурированный ParsedDocument."""
def extract(path: str, use_vision: bool = True) -> ParsedDocument:
"""Разобрать один файл прайса в структурированный ParsedDocument.
use_vision=False отключает добор через Gemini Vision (быстрый прогон без
обращений к API), даже если ключ задан.
"""
name = Path(path).name
fmt = readers.detect_format(path)
doc = ParsedDocument(file_name=name, file_format=fmt, partner_name=readers.partner_from_name(name))
@@ -55,7 +59,7 @@ def extract(path: str) -> ParsedDocument:
doc.parse_log.append(f"ошибка разбора таблицы: {type(exc).__name__}: {exc}")
# Трудные PDF (скан, битый слой, нулевой разбор) добираем через Gemini Vision.
needs_vision = fmt == "pdf" and (doc.file_format == "scan_pdf" or not doc.rows)
needs_vision = use_vision and fmt == "pdf" and (doc.file_format == "scan_pdf" or not doc.rows)
if needs_vision and os.environ.get("GEMINI_API_KEY"):
try:
from etl.extractors.vision import extract_with_vision
+137
View File
@@ -0,0 +1,137 @@
"""Конвейер обработки архива: извлечение → нормализация → валидация → SQLite.
Для MVP данные складываются в SQLite (разворачивается без Docker), а боевая схема
PostgreSQL + pgvector лежит в db/migrations. Здесь же — дедупликация партнёров,
версионирование цен (последний прайс активен, старые архивируются) и сбор отчёта о
качестве для дашборда и сдачи.
"""
from __future__ import annotations
import json
import sqlite3
import sys
import uuid
from collections import Counter
from datetime import date
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parents[1]))
from etl.dictionary import load_dictionary, normalize_name
from etl.extractors import extract
from etl.normalize import Matcher
from etl.validate import validate_item
SCHEMA = """
CREATE TABLE IF NOT EXISTS service (
service_id TEXT PRIMARY KEY, specialty TEXT, name_ru TEXT, name_norm TEXT, tarificator_code TEXT
);
CREATE TABLE IF NOT EXISTS partner (
partner_id TEXT PRIMARY KEY, name TEXT, name_norm TEXT UNIQUE, city TEXT
);
CREATE TABLE IF NOT EXISTS price_document (
doc_id TEXT PRIMARY KEY, partner_id TEXT, file_name TEXT, file_format TEXT,
effective_date TEXT, parse_status TEXT, rows_count INTEGER
);
CREATE TABLE IF NOT EXISTS price_item (
item_id INTEGER PRIMARY KEY AUTOINCREMENT, doc_id TEXT, partner_id TEXT,
service_name_raw TEXT, service_code_source TEXT, service_id TEXT,
prices TEXT, price_resident REAL, price_nonresident REAL, unit TEXT,
effective_date TEXT, map_method TEXT, map_confidence REAL, status TEXT, is_active INTEGER
);
CREATE INDEX IF NOT EXISTS price_item_service ON price_item(service_id);
CREATE INDEX IF NOT EXISTS price_item_partner ON price_item(partner_id);
"""
def _open(db_path: str) -> sqlite3.Connection:
conn = sqlite3.connect(db_path)
conn.executescript(SCHEMA)
for table in ("price_item", "price_document", "partner", "service"):
conn.execute(f"DELETE FROM {table}") # пересборка с нуля при каждом прогоне
return conn
def run(data_dir: str, db_path: str, dict_path: str, embedder=None, use_vision: bool = False,
today: date | None = None) -> dict:
"""Прогнать весь архив в SQLite и вернуть отчёт о качестве."""
today = today or date.today()
services = load_dictionary(dict_path)
matcher = Matcher(services, embedder=embedder)
conn = _open(db_path)
conn.executemany(
"INSERT INTO service VALUES (?,?,?,?,?)",
[(s.service_id, s.specialty, s.name_ru, s.name_norm, s.tarificator_code) for s in services],
)
partners: dict[str, str] = {}
staged: list[tuple] = [] # (doc_id, partner_id, raw, code, prices, unit, eff_iso)
report: Counter = Counter()
for path in sorted(Path(data_dir).glob("*")):
doc = extract(str(path), use_vision=use_vision)
partner_name = doc.partner_name or path.stem
partner_norm = normalize_name(partner_name)
partner_id = partners.get(partner_norm)
if partner_id is None:
partner_id = str(uuid.uuid4())
partners[partner_norm] = partner_id
conn.execute("INSERT INTO partner VALUES (?,?,?,?)", (partner_id, partner_name, partner_norm, None))
doc_id = str(uuid.uuid4())
eff_iso = doc.effective_date.isoformat() if doc.effective_date else None
conn.execute(
"INSERT INTO price_document VALUES (?,?,?,?,?,?,?)",
(doc_id, partner_id, doc.file_name, doc.file_format, eff_iso, "done", len(doc.rows)),
)
report["documents"] += 1
for row in doc.rows:
staged.append((doc_id, partner_id, row.service_name_raw, row.service_code_source,
row.prices, row.unit, eff_iso))
# Нормализация одной пачкой (эмбеддинги считаются разом — быстрее и дешевле).
matches = matcher.match_batch([s[2] for s in staged], [s[3] for s in staged])
for (doc_id, partner_id, raw, code, prices, unit, eff_iso), match in zip(staged, matches, strict=True):
status, _flags = validate_item(raw, prices, effective_date=None, today=today)
if match.service_id:
report["matched"] += 1
report[f"method_{match.method}"] += 1
conn.execute(
"""INSERT INTO price_item
(doc_id, partner_id, service_name_raw, service_code_source, service_id, prices,
price_resident, price_nonresident, unit, effective_date, map_method, map_confidence,
status, is_active)
VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,1)""",
(doc_id, partner_id, raw, code, match.service_id, json.dumps(prices, ensure_ascii=False),
prices.get("resident"), prices.get("nonresident"), unit, eff_iso, match.method,
round(match.confidence, 3), status),
)
report["items"] += 1
# Версионирование: если у партнёра по той же услуге есть более свежий прайс — старые архивируем.
conn.execute(
"""UPDATE price_item SET is_active = 0
WHERE effective_date IS NOT NULL AND EXISTS (
SELECT 1 FROM price_item b
WHERE b.partner_id = price_item.partner_id
AND b.service_name_raw = price_item.service_name_raw
AND b.effective_date > price_item.effective_date)"""
)
conn.commit()
total = report["items"]
summary = {
"documents": report["documents"],
"partners": len(partners),
"items": total,
"auto_matched": report["matched"],
"auto_matched_pct": round(100 * report["matched"] / total, 1) if total else 0.0,
"unmatched": total - report["matched"],
"by_method": {
method: report[f"method_{method}"] for method in ("code", "exact", "embedding", "fuzzy", None)
},
}
conn.close()
return summary