1ff5d93cd8
- очередь верификации: предложенный кандидат (топ эмбеддинг), подтвердить/другое/пропустить - POST /match пополняет синонимы (learned_synonym) — обучающаяся нормализация - сравнение: медиана + разброс цен + отклонение клиники (переплата/выгода) - эмбеддер: дедуп названий + ретрай при 429 (квота Gemini) - развёрнуто на med.secondbrain.tools Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
142 lines
6.6 KiB
Python
142 lines
6.6 KiB
Python
"""Конвейер обработки архива: извлечение → нормализация → валидация → 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,
|
|
suggested_service_id TEXT, suggested_score REAL
|
|
);
|
|
-- Синонимы, выученные при ручной верификации (для дообучения нормализации).
|
|
CREATE TABLE IF NOT EXISTS learned_synonym (name_norm TEXT PRIMARY KEY, service_id TEXT);
|
|
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:
|
|
# Полная пересборка: удаляем файл, чтобы схема всегда создавалась свежей.
|
|
Path(db_path).unlink(missing_ok=True)
|
|
conn = sqlite3.connect(db_path)
|
|
conn.executescript(SCHEMA)
|
|
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, suggested_service_id, suggested_score)
|
|
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, match.suggested_service_id,
|
|
round(match.suggested_score, 3) if match.suggested_score else None),
|
|
)
|
|
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
|