650 lines
20 KiB
Python
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() |