Files
medtech-hackathon/etl/pipeline.py
T

256 lines
11 KiB
Python

"""Конвейер обработки архива: извлечение → нормализация → валидация → SQLite.
Для MVP данные складываются в SQLite (разворачивается без Docker), а боевая схема
PostgreSQL + pgvector лежит в db/migrations. Здесь же — дедупликация партнёров,
версионирование цен (последний прайс активен, старые архивируются) и сбор отчёта о
качестве для дашборда и сдачи.
Два входа: `run()` собирает базу с нуля (полная пересборка), `ingest()` догружает
новые прайсы в уже собранную базу (append) — на нём держится приём ZIP через интерфейс.
"""
from __future__ import annotations
import json
import sqlite3
import sys
import uuid
from collections import Counter
from datetime import date
from pathlib import Path
import numpy as np
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);
"""
_INSERT_ITEM = """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,?,?)"""
def _emb_path(db_path: str) -> Path:
"""Путь к кэшу эмбеддингов справочника рядом с базой."""
return Path(db_path).parent / "dict_emb.npy"
def _open(db_path: str) -> sqlite3.Connection:
# Полная пересборка: удаляем файл, чтобы схема всегда создавалась свежей.
Path(db_path).unlink(missing_ok=True)
conn = sqlite3.connect(db_path)
conn.row_factory = sqlite3.Row
conn.executescript(SCHEMA)
return conn
def _load_services(conn, services) -> None:
"""Залить справочник услуг, если таблица ещё пуста."""
if conn.execute("SELECT 1 FROM service LIMIT 1").fetchone():
return
conn.executemany(
"INSERT INTO service VALUES (?,?,?,?,?)",
[(s.service_id, s.specialty, s.name_ru, s.name_norm, s.tarificator_code) for s in services],
)
def _process_documents(conn, data_dir, matcher, partners, today, use_vision) -> Counter:
"""Извлечь файлы из каталога, сопоставить со справочником и записать позиции.
Мутирует `conn` и словарь `partners` (норм. имя → partner_id, чтобы один и тот же
партнёр из разных файлов не задвоился). Возвращает счётчик для отчёта о качестве.
"""
report: Counter = Counter()
staged: list[tuple] = [] # (doc_id, partner_id, raw, code, prices, unit, eff_iso)
for path in sorted(Path(data_dir).glob("*")):
if path.is_dir():
continue
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),
)
report["new_partners"] += 1
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_ITEM,
(
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
return report
def _apply_versioning(conn) -> None:
"""У партнёра по той же услуге более свежий прайс вытесняет старые (is_active = 0)."""
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)"""
)
def _summary(report: Counter) -> dict:
total = report["items"]
return {
"documents": report["documents"],
"partners": report["new_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)
},
}
def run(
data_dir: str,
db_path: str,
dict_path: str,
embedder=None,
use_vision: bool = False,
today: date | None = None,
) -> dict:
"""Собрать базу с нуля по всему архиву и вернуть отчёт о качестве."""
today = today or date.today()
services = load_dictionary(dict_path)
matcher = Matcher(services, embedder=embedder)
# Кэшируем эмбеддинги справочника: при последующих догрузках их не пересчитываем.
if matcher.dict_emb is not None:
np.save(_emb_path(db_path), matcher.dict_emb)
conn = _open(db_path)
_load_services(conn, services)
partners: dict[str, str] = {}
report = _process_documents(conn, data_dir, matcher, partners, today, use_vision)
_apply_versioning(conn)
conn.commit()
summary = _summary(report)
conn.close()
return summary
def ingest(
data_dir: str,
db_path: str,
dict_path: str,
embedder=None,
use_vision: bool | None = None,
today: date | None = None,
) -> dict:
"""Догрузить новые прайсы из каталога в существующую базу (append, без пересборки).
На этом держится приём ZIP через интерфейс: справочник и его эмбеддинги уже готовы,
поэтому считается только новый файл. Партнёры дедуплицируются с уже загруженными,
версионирование вытесняет устаревшие прайсы той же клиники.
"""
today = today or date.today()
services = load_dictionary(dict_path)
emb_path = _emb_path(db_path)
dict_emb = np.load(emb_path) if (embedder is not None and emb_path.exists()) else None
matcher = Matcher(services, embedder=embedder, dict_emb=dict_emb)
if use_vision is None:
use_vision = embedder is not None # есть ключ Gemini — разрешаем распознавание сканов
conn = sqlite3.connect(db_path)
conn.row_factory = sqlite3.Row
conn.executescript(SCHEMA)
_load_services(conn, services)
# Уже загруженные партнёры — чтобы повторная выгрузка той же клиники не задвоила её.
partners = {
row["name_norm"]: row["partner_id"]
for row in conn.execute("SELECT partner_id, name_norm FROM partner")
}
report = _process_documents(conn, data_dir, matcher, partners, today, use_vision)
_apply_versioning(conn)
conn.commit()
summary = _summary(report)
conn.close()
return summary