Files
trueconf_bot/search_bot/smart_indexer.py

650 lines
20 KiB
Python

import os
import sys
import re
import time
import uuid
import json
import hashlib
import requests
from glob import glob
from datetime import datetime
# Импорт SentenceTransformer УДАЛЕН для разгрузки ОЗУ сервера
from qdrant_client import QdrantClient
from qdrant_client.models import (
Distance,
VectorParams,
PointStruct,
PayloadSchemaType,
Filter,
FieldCondition,
MatchValue
)
# =========================================================
# ИМПОРТ ЦЕНТРАЛЬНОГО КОНФИГА БОТА
# =========================================================
# Добавляем корневую папку бота в sys.path, чтобы скрипт увидел модуль config,
# даже если скрипт запускается из подпапки.
sys.path.append(os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
import config.config as config
# =========================================================
# CONFIG (Берем данные из центрального config.py)
# =========================================================
DOCS_DIR = config.SEARCH_DOCS_PATH
QDRANT_HOST = config.QDRANT_HOST
QDRANT_PORT = config.QDRANT_PORT
COLLECTION_NAME = config.QDRANT_COLLECTION
# Формируем путь к реестру внутри папки логов, указанной в конфиге
REGISTRY_PATH = os.path.join(config.LOGS_PATH, "indexing_registry.json")
# =========================================================
# ЛОКАЛЬНЫЕ НАСТРОЙКИ ИНДЕКСАТОРА
# =========================================================
UNSTRUCTURED_API_URL = "http://172.16.80.250:9000/general/v0/general"
BATCH_SIZE = 32
RECREATE_COLLECTION = False
MIN_TEXT_LENGTH = 25
# =========================================================
# INIT
# =========================================================
# Локальная инициализация нейросети убрана. Работаем через API.
print("🔄 Подключение к Qdrant...")
client = QdrantClient(
host=QDRANT_HOST,
port=QDRANT_PORT
)
# =========================================================
# REGISTRY & HASHING
# =========================================================
def get_file_hash(filepath):
"""Вычисляет MD5 хэш физического файла"""
hasher = hashlib.md5()
with open(filepath, 'rb') as f:
buf = f.read()
hasher.update(buf)
return hasher.hexdigest()
def load_registry():
"""Загружает текущее состояние проиндексированных файлов"""
if os.path.exists(REGISTRY_PATH):
with open(REGISTRY_PATH, 'r', encoding='utf-8') as f:
return json.load(f)
return {}
def save_registry(registry):
"""Сохраняет реестр на диск"""
os.makedirs(os.path.dirname(REGISTRY_PATH), exist_ok=True)
with open(REGISTRY_PATH, 'w', encoding='utf-8') as f:
json.dump(registry, f, ensure_ascii=False, indent=4)
# =========================================================
# HELPERS
# =========================================================
def sha256_text(text):
return hashlib.sha256(
text.encode("utf-8")
).hexdigest()
def clean_text(text):
if not text:
return ""
text = re.sub(
r"\s+",
" ",
text
)
return text.strip()
def is_page_number(text):
return bool(
re.match(
r"^Стр\.\s*\d+\s*из\s*\d+$",
text
)
)
def is_header_noise(text):
patterns = [
"АО «ХК «Сибцем»",
"Тип документа:",
"Ведущее подразделение:",
"Дата утверждения:",
"Редакция 1",
"Оглавление"
]
return any(
p in text
for p in patterns
)
def is_junk(item, text):
if not text:
return True
if len(text) < MIN_TEXT_LENGTH:
return True
if item["type"] in [
"Header",
"Footer",
"Image"
]:
return True
if is_page_number(text):
return True
if is_header_noise(text):
return True
return False
# =========================================================
# TABLE PARSER
# =========================================================
def html_table_to_text(html):
html = re.sub(
r"\n+",
" ",
html
)
rows = re.findall(
r"<tr.*?>(.*?)</tr>",
html,
re.I | re.S
)
parsed_rows = []
for row in rows:
cells = re.findall(
r"<t[hd].*?>(.*?)</t[hd]>",
row,
re.I | re.S
)
clean_cells = []
for cell in cells:
box = re.sub(
r"<[^>]+>",
" ",
cell
)
box = clean_text(box)
clean_cells.append(box)
if clean_cells:
parsed_rows.append(
clean_cells
)
if not parsed_rows:
return ""
headers = parsed_rows[0]
result_lines = []
for row_idx, row in enumerate(parsed_rows[1:], 1):
result_lines.append(
f"Запись {row_idx}:"
)
for col_idx, cell in enumerate(row):
if not cell:
continue
if col_idx < len(headers):
header = headers[col_idx]
else:
header = f"Колонка {col_idx+1}"
result_lines.append(
f"{header}: {cell}"
)
result_lines.append("")
return "\n".join(result_lines).strip()
# =========================================================
# MERGE FRAGMENTS
# =========================================================
def merge_fragments(elements):
merged = []
buffer = []
buffer_page = None
buffer_section = None
def flush_buffer():
nonlocal buffer
if not buffer:
return
text = clean_text(
" ".join(buffer)
)
if text:
merged.append({
"type": "MergedText",
"text": text,
"page_number": buffer_page,
"section_title": buffer_section
})
buffer = []
for el in elements:
text = clean_text(
el["text"]
)
if not text:
continue
el_type = el["type"]
page = el["metadata"].get("page_number", 0)
section = el.get("section_title", "")
# =========================================
# TABLE
# =========================================
if el_type == "Table":
flush_buffer()
merged.append({
"type": "Table",
"text": text,
"html": el["metadata"].get("text_as_html", ""),
"page_number": page,
"section_title": section
})
continue
# =========================================
# SHORT TEXT
# =========================================
if el_type in [
"NarrativeText",
"ListItem",
"UncategorizedText"
]:
if len(text) < 80:
buffer.append(text)
buffer_page = page
buffer_section = section
else:
flush_buffer()
merged.append({
"type": "Text",
"text": text,
"page_number": page,
"section_title": section
})
else:
flush_buffer()
flush_buffer()
return merged
# =========================================================
# QDRANT INIT & DELETE
# =========================================================
def init_collection():
collections = [
c.name
for c in client.get_collections().collections
]
if (
RECREATE_COLLECTION
and COLLECTION_NAME in collections
):
print("🗑 Удаление коллекции...")
client.delete_collection(COLLECTION_NAME)
collections = [
c.name
for c in client.get_collections().collections
]
if COLLECTION_NAME not in collections:
print("⚙️ Создание коллекции...")
client.create_collection(
collection_name=COLLECTION_NAME,
vectors_config=VectorParams(
size=1024,
distance=Distance.COSINE
)
)
indexes = [
("metadata.filename", PayloadSchemaType.KEYWORD),
("metadata.section_title", PayloadSchemaType.KEYWORD),
("metadata.page_number", PayloadSchemaType.INTEGER),
("metadata.content_type", PayloadSchemaType.KEYWORD)
]
for field_name, field_schema in indexes:
client.create_payload_index(
collection_name=COLLECTION_NAME,
field_name=field_name,
field_schema=field_schema
)
def delete_file_from_qdrant(filename):
"""Удаляет все векторы, привязанные к конкретному файлу"""
client.delete(
collection_name=COLLECTION_NAME,
points_selector=Filter(
must=[
FieldCondition(
key="metadata.filename",
match=MatchValue(value=filename)
)
]
)
)
print(f" 🗑 Удалены старые данные файла {filename} из базы.")
# =========================================================
# PARSE DOCUMENT
# =========================================================
def parse_document(file_path):
filename = os.path.basename(
file_path
)
print(f"\n📄 Парсинг: {filename}")
payload_data = {
"strategy": "hi_res",
"languages": ["rus"],
"pdf_infer_table_structure": "true"
}
ext = filename.lower().split(".")[-1]
if ext == "pdf":
mime_type = "application/pdf"
else:
mime_type = (
"application/vnd.openxmlformats-officedocument."
"wordprocessingml.document"
)
with open(file_path, "rb") as f:
files = {
"files": (
filename,
f,
mime_type
)
}
response = requests.post(
UNSTRUCTURED_API_URL,
files=files,
data=payload_data,
timeout=900
)
response.raise_for_status()
return response.json()
# =========================================================
# INDEX DOCUMENT
# =========================================================
def index_document(file_path, data):
filename = os.path.basename(file_path)
cleaned = []
current_title = "Общая информация"
# =====================================================
# CLEAN
# =====================================================
for item in data:
raw_text = item.get("text", "")
text = clean_text(raw_text)
if is_junk(item, text):
continue
if item["type"] == "Title":
current_title = text
continue
item["section_title"] = current_title
cleaned.append(item)
# =====================================================
# MERGE
# =====================================================
merged = merge_fragments(cleaned)
points = []
# =====================================================
# BUILD CHUNKS
# =====================================================
for idx, chunk in enumerate(merged):
chunk_type = chunk["type"]
section_title = chunk["section_title"]
page_number = chunk["page_number"]
# =================================================
# TABLE / TEXT
# =================================================
if chunk_type == "Table":
table_text = html_table_to_text(chunk["html"])
content = (
f"Раздел: {section_title}\n"
f"Тип: Таблица\n\n"
f"{table_text}"
)
else:
content = (
f"Раздел: {section_title}\n\n"
f"{chunk['text']}"
)
content = clean_text(content)
if len(content) < 40:
continue
# =================================================
# VECTOR (Переведено на центральный Embedding API)
# =================================================
try:
payload_api = {
"model": config.EMBEDDING_MODEL_NAME, # Отправляем "bge-m3"
"input_type": "passage", # Для сохранения документов используем тип passage
"texts": [content]
}
headers_api = {"Content-Type": "application/json"}
if hasattr(config, 'EMBEDDING_API_TOKEN') and config.EMBEDDING_API_TOKEN:
headers_api["X-Embedding-Token"] = config.EMBEDDING_API_TOKEN
api_resp = requests.post(config.EMBEDDING_API_URL, json=payload_api, headers=headers_api)
api_resp.raise_for_status()
vector = api_resp.json()["embeddings"][0]
except Exception as e:
print(f"❌ Ошибка векторизации чанка {idx} через API: {e}")
continue # Пропускаем чанк, если центральный сервис недоступен
chunk_hash = sha256_text(
filename + str(idx) + content
)
point_id = str(
uuid.uuid5(
uuid.NAMESPACE_DNS,
chunk_hash
)
)
payload = {
"page_content": content,
"metadata": {
"filename": filename,
"source": file_path,
"chunk_index": idx,
"section_title": section_title,
"page_number": page_number,
"content_type": chunk_type,
"chunk_hash": chunk_hash
}
}
points.append(
PointStruct(
id=point_id,
vector=vector,
payload=payload
)
)
# =================================================
# UPSERT BATCH
# =================================================
if len(points) >= BATCH_SIZE:
client.upsert(
collection_name=COLLECTION_NAME,
points=points
)
print(f" ⬆️ batch={len(points)}")
points = []
# =====================================================
# FINAL UPSERT
# =====================================================
if points:
client.upsert(
collection_name=COLLECTION_NAME,
points=points
)
print(f" ✅ Загружено chunks: {len(merged)}")
return len(merged)
# =========================================================
# MAIN
# =========================================================
def main():
init_collection()
target_files = []
for ext in ["*.pdf", "*.docx"]:
target_files.extend(
glob(
os.path.join(DOCS_DIR, "**", ext),
recursive=True
)
)
# ФИЛЬТРАЦИЯ СКРЫТЫХ ВРЕМЕННЫХ ФАЙЛОВ WORD (~$)
target_files = [f for f in target_files if not os.path.basename(f).startswith("~$")]
target_files.sort()
if not target_files:
print("❌ PDF/DOCX файлы не найдены")
return
print(f"📂 Найдено файлов для индексации: {len(target_files)}")
print(f"📡 Использование удаленной векторизации: {config.EMBEDDING_API_URL} ({config.EMBEDDING_MODEL_NAME})")
# Загружаем реестр
registry = load_registry()
# =====================================================
# PROCESS
# =====================================================
for file_path in target_files:
filename = os.path.basename(file_path)
current_hash = get_file_hash(file_path)
file_info = registry.get(filename)
# Проверка: Файл уже есть, хэш совпадает, статус был успешным
if file_info and file_info.get("hash") == current_hash and file_info.get("status") == "SUCCESS":
print(f"\n⏭️ Пропуск (не изменился): {filename}")
continue
print(f"\n🔄 Обработка: {filename}")
# Если файл есть в реестре со статусом SUCCESS, но хэш другой (файл обновили)
if file_info and file_info.get("status") == "SUCCESS":
print(f" ⚠️ Обнаружена новая версия файла! Обновляем индекс...")
delete_file_from_qdrant(filename)
try:
data = parse_document(file_path)
chunks_count = index_document(file_path, data)
# Запись успешного результата
registry[filename] = {
"hash": current_hash,
"status": "SUCCESS",
"chunks": chunks_count,
"timestamp": datetime.now().strftime("%Y-%m-%d %H:%M:%S")
}
save_registry(registry)
# ПАУЗА ПРИ НЕДОСТУПНОСТИ СЕРВЕРА ПАРСИНГА
except requests.exceptions.ConnectionError:
print(f" ⚠️ Сервер парсинга документов недоступен. Ждем 60 секунд...")
time.sleep(60)
registry[filename] = {
"hash": current_hash,
"status": "ERROR",
"error_details": "Connection refused",
"timestamp": datetime.now().strftime("%Y-%m-%d %H:%M:%S")
}
save_registry(registry)
except Exception as e:
print(f" ❌ Ошибка обработки файла {filename}: {e}")
# Запись ошибки в реестр
registry[filename] = {
"hash": current_hash,
"status": "ERROR",
"error_details": str(e),
"timestamp": datetime.now().strftime("%Y-%m-%d %H:%M:%S")
}
save_registry(registry)
print("\n🎉 ГОТОВО. База Qdrant полностью синхронизирована через API.")
# =========================================================
if __name__ == "__main__":
main()