feat(ingest+поиск): приём ZIP через интерфейс, регистронезависимый поиск, чистка очереди

- Загрузка архива (§4.1): POST /ingest распаковывает ZIP, дораспознаёт и
  догружает прайсы в базу (append). Виджет в админке. Кэш эмбеддингов справочника
  (dict_emb.npy) — догрузка считает только новый файл, эмбеддинги best-effort с
  откатом на fuzzy при сбое Gemini.
- Поиск по name_norm (lowercase): «УЗИ» и «узи» дают одинаковый результат —
  SQLite LIKE не сворачивает регистр для кириллицы.
- Фильтр названий: код-мнемоники («МОЗО.15») больше не попадают в очередь.
- /stats считается из базы — цифры на дашборде обновляются после загрузки.
- Деплой через Dockerfile (зависимости в образе, код монтируется томом).
- ruff: ядро и весь репозиторий проходят проверку.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
2026-06-27 11:35:12 +05:00
parent 89928511ba
commit 69a9884db5
25 changed files with 2217 additions and 182 deletions
+4 -3
View File
@@ -10,6 +10,7 @@
отвечает только за загрузку и подготовку индексов; само сопоставление — в
`etl/normalize`.
"""
from __future__ import annotations
import re
@@ -40,9 +41,9 @@ def normalize_name(text: str | None) -> str:
class Service:
"""Одна запись эталонного справочника услуг."""
service_id: str # стабильный идентификатор «<ID>-<Code>»
specialty: str # специальность — служит и категорией, и сужением поиска
name_ru: str # официальное название
service_id: str # стабильный идентификатор «<ID>-<Code>»
specialty: str # специальность — служит и категорией, и сужением поиска
name_ru: str # официальное название
tarificator_code: str | None # код тарификатора (если задан)
name_norm: str = field(default="") # нормализованная форма name_ru
+4 -1
View File
@@ -4,6 +4,7 @@
PDF с битым текстовым слоем помечается как scan_pdf — его добирает путь через Vision.
Любой сбой чтения логируется, но не роняет обработку остального архива.
"""
from __future__ import annotations
import os
@@ -26,7 +27,9 @@ def extract(path: str, use_vision: bool = True) -> ParsedDocument:
"""
name = Path(path).name
fmt = readers.detect_format(path)
doc = ParsedDocument(file_name=name, file_format=fmt, partner_name=readers.partner_from_name(name))
doc = ParsedDocument(
file_name=name, file_format=fmt, partner_name=readers.partner_from_name(name)
)
year = readers.year_from_name(name)
if year:
+7 -1
View File
@@ -6,6 +6,7 @@
в числах, многотарифные колонки цен, строки-заголовки секций и шапку не в первой
строке таблицы.
"""
from __future__ import annotations
import re
@@ -85,7 +86,12 @@ def is_real_name(text) -> bool:
"""
s = clean(text)
cyrillic = sum(1 for ch in s if "а" <= ch.lower() <= "я" or ch.lower() == "ё")
return cyrillic >= 3 and not looks_like_code(s)
if cyrillic < 3 or looks_like_code(s):
return False
# Код-мнемоника: цифры + только заглавные буквы + без пробела (напр. «МОЗО.15») — не услуга.
if any(c.isdigit() for c in s) and " " not in s and not any(c.islower() for c in s):
return False
return True
def classify_price_tier(header_text) -> str | None:
+22 -5
View File
@@ -9,6 +9,7 @@
- настоящая цена медуслуги — сотни и тысячи тенге, а не мелкое число;
- название услуги — самая «кириллическая» колонка, код — цифро-точечная.
"""
from __future__ import annotations
import statistics
@@ -97,7 +98,11 @@ def _assign_columns(grid: Grid, header_idx: int):
# --- колонка названия: явный заголовок, иначе самая кириллическая ---
name_col = next(
(c for c in range(ncol) if c not in price_cols and any(k in header[c].lower() for k in _NAME_KEYS)),
(
c
for c in range(ncol)
if c not in price_cols and any(k in header[c].lower() for k in _NAME_KEYS)
),
None,
)
if name_col is None:
@@ -109,7 +114,9 @@ def _assign_columns(grid: Grid, header_idx: int):
(
c
for c in range(ncol)
if c != name_col and c not in price_cols and any(k in header[c].lower() for k in _CODE_HEADER_KEYS)
if c != name_col
and c not in price_cols
and any(k in header[c].lower() for k in _CODE_HEADER_KEYS)
),
None,
)
@@ -123,7 +130,11 @@ def _assign_columns(grid: Grid, header_idx: int):
break
unit_col = next(
(c for c in range(ncol) if c not in price_cols and any(k in header[c].lower() for k in _UNIT_KEYS)),
(
c
for c in range(ncol)
if c not in price_cols and any(k in header[c].lower() for k in _UNIT_KEYS)
),
None,
)
return name_col, code_col, unit_col, price_cols
@@ -165,10 +176,16 @@ def extract_rows_from_grid(grid: Grid) -> list[RawRow]:
RawRow(
service_name_raw=name,
service_code_source=(
r[code_col] if code_col is not None and code_col < len(r) and r[code_col] else None
r[code_col]
if code_col is not None and code_col < len(r) and r[code_col]
else None
),
prices=prices,
unit=(r[unit_col] if unit_col is not None and unit_col < len(r) and r[unit_col] else None),
unit=(
r[unit_col]
if unit_col is not None and unit_col < len(r) and r[unit_col]
else None
),
section=section,
)
)
+8 -2
View File
@@ -4,6 +4,7 @@
многолистовых книг и многотабличных PDF. Сырой текст нужен для аудита и для
детектора битого текстового слоя.
"""
from __future__ import annotations
import re
@@ -29,14 +30,19 @@ def partner_from_name(name: str) -> str:
def detect_format(path: str) -> str:
ext = Path(path).suffix.lower()
return {".xlsx": "xlsx", ".xls": "xls", ".docx": "docx", ".pdf": "pdf"}.get(ext, ext.lstrip("."))
return {".xlsx": "xlsx", ".xls": "xls", ".docx": "docx", ".pdf": "pdf"}.get(
ext, ext.lstrip(".")
)
def docx_grids(path: str) -> tuple[list[Grid], str]:
import docx
document = docx.Document(path)
grids = [[[clean(cell.text) for cell in row.cells] for row in table.rows] for table in document.tables]
grids = [
[[clean(cell.text) for cell in row.cells] for row in table.rows]
for table in document.tables
]
text = "\n".join(p.text for p in document.paragraphs)
return grids, text
+1
View File
@@ -5,6 +5,7 @@
вернуть позиции прайса строго по JSON-схеме. Требует переменную окружения
GEMINI_API_KEY; модель задаётся через GEMINI_MODEL (по умолчанию gemini-2.0-flash).
"""
from __future__ import annotations
import json
+1
View File
@@ -1,4 +1,5 @@
"""Нормализация извлечённых позиций к справочнику услуг."""
from etl.normalize.matcher import Matcher
__all__ = ["Matcher"]
+4 -1
View File
@@ -4,6 +4,7 @@
раздувается и не держит модель в RAM. Размерность 768 выбрана, чтобы вектор лез в
HNSW-индекс pgvector. Требуется переменная окружения GEMINI_API_KEY.
"""
from __future__ import annotations
import os
@@ -54,7 +55,9 @@ class GeminiEmbedder:
chunk = items[start : start + self.batch]
for attempt in range(8):
try:
response = self._client.models.embed_content(model=self.model, contents=chunk, config=config)
response = self._client.models.embed_content(
model=self.model, contents=chunk, config=config
)
vectors.extend(embedding.values for embedding in response.embeddings)
break
except Exception as exc: # пауза и повтор при превышении квоты (429)
+15 -5
View File
@@ -11,6 +11,7 @@
клиник встречаются с лишним хвостом («A02.020.000.2»), поэтому сравниваем по
канонической части кода тарификатора.
"""
from __future__ import annotations
import re
@@ -37,6 +38,7 @@ class Matcher:
embedder=None,
emb_threshold: float = 0.70,
fuzzy_threshold: int = 88,
dict_emb=None,
):
self.services = services
self.by_code = {s.tarificator_code: s for s in services if s.tarificator_code}
@@ -47,8 +49,10 @@ class Matcher:
self.emb_threshold = emb_threshold
self.fuzzy_threshold = fuzzy_threshold
self.embedder = embedder
self.dict_emb = None
if embedder is not None:
# Готовые эмбеддинги справочника (dict_emb) переиспользуются между прогонами:
# при догрузке нового прайса не нужно заново векторизовать тысячи услуг.
self.dict_emb = dict_emb
if embedder is not None and self.dict_emb is None:
self.dict_emb = embedder.encode(
[s.name_ru for s in services],
normalize_embeddings=True,
@@ -65,7 +69,9 @@ class Matcher:
best = process.extractOne(name_norm, self.names_norm, scorer=fuzz.token_set_ratio)
if best and best[1] >= self.fuzzy_threshold:
return MatchResult(
service_id=self.services[best[2]].service_id, method="fuzzy", confidence=best[1] / 100
service_id=self.services[best[2]].service_id,
method="fuzzy",
confidence=best[1] / 100,
)
return MatchResult()
@@ -78,12 +84,16 @@ class Matcher:
for i, (name, code) in enumerate(zip(names, codes, strict=True)):
service = self._match_by_code(code)
if service:
results[i] = MatchResult(service_id=service.service_id, method="code", confidence=1.0)
results[i] = MatchResult(
service_id=service.service_id, method="code", confidence=1.0
)
continue
name_norm = normalize_name(name)
exact = self.by_norm.get(name_norm)
if exact:
results[i] = MatchResult(service_id=exact.service_id, method="exact", confidence=1.0)
results[i] = MatchResult(
service_id=exact.service_id, method="exact", confidence=1.0
)
continue
pending_idx.append(i)
pending_norm.append(name_norm)
+143 -29
View File
@@ -4,7 +4,11 @@
PostgreSQL + pgvector лежит в db/migrations. Здесь же — дедупликация партнёров,
версионирование цен (последний прайс активен, старые архивируются) и сбор отчёта о
качестве для дашборда и сдачи.
Два входа: `run()` собирает базу с нуля (полная пересборка), `ingest()` догружает
новые прайсы в уже собранную базу (append) — на нём держится приём ZIP через интерфейс.
"""
from __future__ import annotations
import json
@@ -15,6 +19,8 @@ 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
@@ -46,33 +52,49 @@ 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 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)
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],
)
partners: dict[str, str] = {}
staged: list[tuple] = [] # (doc_id, partner_id, raw, code, prices, unit, eff_iso)
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)
@@ -80,7 +102,11 @@ def run(data_dir: str, db_path: str, dict_path: str, embedder=None, use_vision:
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))
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
@@ -90,31 +116,53 @@ def run(data_dir: str, db_path: str, dict_path: str, embedder=None, use_vision:
)
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))
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):
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),
_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 (
@@ -123,19 +171,85 @@ def run(data_dir: str, db_path: str, dict_path: str, embedder=None, use_vision:
AND b.service_name_raw = price_item.service_name_raw
AND b.effective_date > price_item.effective_date)"""
)
conn.commit()
def _summary(report: Counter) -> dict:
total = report["items"]
summary = {
return {
"documents": report["documents"],
"partners": len(partners),
"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)
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
+1
View File
@@ -1,4 +1,5 @@
"""Валидация позиций прайса по правилам ТЗ §4.4."""
from etl.validate.rules import Flag, to_kzt, validate_item
__all__ = ["Flag", "to_kzt", "validate_item"]
+9 -2
View File
@@ -4,6 +4,7 @@
ошибка → error, предупреждение → needs_review, иначе → done. Конвертация валют и
сравнение с прошлой версией (детектор аномалий) тоже здесь.
"""
from __future__ import annotations
from dataclasses import dataclass
@@ -34,7 +35,9 @@ def validate_prices(prices: dict[str, float]) -> list[Flag]:
flags: list[Flag] = []
for tier, value in prices.items():
if value is None or value <= 0:
flags.append(Flag("warning", "price_nonpositive", f"тариф «{tier}»: цена не положительна"))
flags.append(
Flag("warning", "price_nonpositive", f"тариф «{tier}»: цена не положительна")
)
resident, nonresident = prices.get("resident"), prices.get("nonresident")
if resident and nonresident and nonresident < resident:
flags.append(
@@ -50,7 +53,11 @@ def validate_date(effective_date: date | None, today: date | None = None) -> lis
def validate_anomaly(new_price: float | None, previous_price: float | None) -> list[Flag]:
if previous_price and new_price and abs(new_price - previous_price) / previous_price > ANOMALY_RATIO:
if (
previous_price
and new_price
and abs(new_price - previous_price) / previous_price > ANOMALY_RATIO
):
return [
Flag(
"warning",