feat(gpu): Qdrant как векторный бэкенд под флагом — снимает SPOF FAISS
FAISS-индекс — файловый синглтон в RAM одного воркера (save на каждую запись, без блокировок, без HA, не горизонтален). Добавлен альтернативный бэкенд Qdrant с тем же classmethod-интерфейсом, выбор через VECTOR_BACKEND — аддитивно и безопасно: дефолт остаётся faiss, ничего не ломается, пока не переключат. - app/qdrant_manager.py — search/add_vectors/save(no-op)/total_vectors поверх qdrant-client (коллекция Cosine, id точки = doc_id, upsert идемпотентен); - app/vector_store.py — get_backend() по настройке; задачи search/plagiarism переведены на него (больше не импортируют FAISSManager напрямую); - app/migrate_faiss_to_qdrant.py — перелив существующих векторов (reconstruct из IndexIDMap2 → upsert), идемпотентно; - docker-compose.prod.yml — сервис qdrant под профилем `qdrant` (по умолчанию не поднимается, ресурсов не ест) + volume; README — раздел про переключение. Тесты (9) гоняют QdrantManager против ВСТРОЕННОГО Qdrant (qdrant-client :memory:, не моки) + диспетчеризацию бэкенда. Всего тестов: 82 (indexer 24, gost 24, gpu 34). Плюсы Qdrant активируются только после явного переключения + миграции. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
@@ -48,6 +48,12 @@ class Settings(BaseSettings):
|
||||
FAISS_NLIST: int = 1024 # Количество кластеров для IVFFlat
|
||||
FAISS_NPROBE: int = 64 # Количество кластеров для поиска
|
||||
|
||||
# Векторный бэкенд: "faiss" (файловый синглтон, дефолт) или "qdrant" (сервис,
|
||||
# снимает SPOF и конкурентную запись). Переключается без изменения кода.
|
||||
VECTOR_BACKEND: str = "faiss"
|
||||
QDRANT_URL: str = "http://qdrant:6333"
|
||||
QDRANT_COLLECTION: str = "documents"
|
||||
|
||||
# App
|
||||
APP_URL: str = "https://academic.jze9.ru"
|
||||
ENVIRONMENT: str = "development"
|
||||
|
||||
47
services/worker-gpu/app/migrate_faiss_to_qdrant.py
Normal file
47
services/worker-gpu/app/migrate_faiss_to_qdrant.py
Normal file
@@ -0,0 +1,47 @@
|
||||
"""Одноразовая миграция векторов из файлового FAISS в Qdrant.
|
||||
|
||||
Запуск (после старта qdrant и выставления VECTOR_BACKEND=qdrant в .env):
|
||||
docker compose -f docker-compose.prod.yml --profile qdrant up -d qdrant
|
||||
docker compose -f docker-compose.prod.yml exec worker-gpu python -m app.migrate_faiss_to_qdrant
|
||||
|
||||
Идемпотентна: повторный запуск перезапишет те же точки по doc_id, дублей не будет.
|
||||
После миграции стоит сверить: QdrantManager.total_vectors() == count(documents).
|
||||
"""
|
||||
|
||||
import logging
|
||||
|
||||
import numpy as np
|
||||
|
||||
logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(message)s")
|
||||
logger = logging.getLogger("migrate_faiss_to_qdrant")
|
||||
|
||||
|
||||
def migrate(batch_size: int = 1000) -> int:
|
||||
"""Перелить все векторы из FAISS-индекса в Qdrant. Возвращает число векторов."""
|
||||
import faiss
|
||||
|
||||
from app.config import settings
|
||||
from app.qdrant_manager import QdrantManager
|
||||
|
||||
index = faiss.read_index(settings.FAISS_INDEX_PATH)
|
||||
if not hasattr(index, "id_map"):
|
||||
raise SystemExit("FAISS-индекс без id_map (несовместимый тип) — миграция невозможна")
|
||||
|
||||
ids = faiss.vector_to_array(index.id_map).astype(np.int64)
|
||||
total = int(len(ids))
|
||||
logger.info("FAISS: %d векторов к миграции в Qdrant", total)
|
||||
|
||||
migrated = 0
|
||||
for start in range(0, total, batch_size):
|
||||
chunk = ids[start : start + batch_size]
|
||||
vecs = np.vstack([index.reconstruct(int(i)) for i in chunk]).astype(np.float32)
|
||||
QdrantManager.add_vectors(vecs, [int(i) for i in chunk])
|
||||
migrated += len(chunk)
|
||||
logger.info("Мигрировано %d/%d", migrated, total)
|
||||
|
||||
logger.info("Готово. Точек в Qdrant: %d", QdrantManager.total_vectors())
|
||||
return migrated
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
migrate()
|
||||
104
services/worker-gpu/app/qdrant_manager.py
Normal file
104
services/worker-gpu/app/qdrant_manager.py
Normal file
@@ -0,0 +1,104 @@
|
||||
"""Менеджер векторного индекса на Qdrant — альтернатива FAISSManager.
|
||||
|
||||
Тот же classmethod-интерфейс (search / add_vectors / save / total_vectors), но
|
||||
вместо файлового синглтона в RAM одного воркера — сетевой сервис Qdrant. Это:
|
||||
- снимает единую точку отказа (индекс живёт отдельно от воркера);
|
||||
- допускает конкурентную запись из нескольких воркеров (upsert по doc_id);
|
||||
- масштабируется горизонтально и переживает рестарт воркера без перезагрузки.
|
||||
|
||||
Дистанция — Cosine; эмбеддинги нормированы, id точки = doc_id из PostgreSQL,
|
||||
поэтому поиск сразу возвращает doc_id. Включается через VECTOR_BACKEND=qdrant.
|
||||
"""
|
||||
|
||||
import logging
|
||||
|
||||
import numpy as np
|
||||
|
||||
from app.config import settings
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class QdrantManager:
|
||||
"""Синглтон-обёртка над коллекцией Qdrant (интерфейс как у FAISSManager)."""
|
||||
|
||||
COLLECTION = settings.QDRANT_COLLECTION
|
||||
|
||||
_client = None
|
||||
_ready = False
|
||||
|
||||
@classmethod
|
||||
def _get_client(cls):
|
||||
"""Ленивое подключение к Qdrant (или заранее внедрённый клиент в тестах)."""
|
||||
if cls._client is None:
|
||||
from qdrant_client import QdrantClient
|
||||
|
||||
cls._client = QdrantClient(url=settings.QDRANT_URL, timeout=30)
|
||||
return cls._client
|
||||
|
||||
@classmethod
|
||||
def _ensure(cls):
|
||||
"""Гарантировать существование коллекции нужной размерности/дистанции."""
|
||||
client = cls._get_client()
|
||||
if cls._ready:
|
||||
return client
|
||||
|
||||
from qdrant_client.models import Distance, VectorParams
|
||||
|
||||
if not client.collection_exists(cls.COLLECTION):
|
||||
client.create_collection(
|
||||
collection_name=cls.COLLECTION,
|
||||
vectors_config=VectorParams(
|
||||
size=settings.EMBED_DIM, distance=Distance.COSINE
|
||||
),
|
||||
)
|
||||
logger.info(f"Создана коллекция Qdrant {cls.COLLECTION!r} (dim={settings.EMBED_DIM})")
|
||||
cls._ready = True
|
||||
return client
|
||||
|
||||
@classmethod
|
||||
def search(cls, query_vector: np.ndarray, k: int = 20) -> list[tuple[int, float]]:
|
||||
"""Поиск k ближайших векторов. Возвращает [(doc_id, cosine_score), ...]."""
|
||||
client = cls._ensure()
|
||||
vec = np.asarray(query_vector, dtype=np.float32).reshape(-1).tolist()
|
||||
try:
|
||||
res = client.query_points(collection_name=cls.COLLECTION, query=vec, limit=k)
|
||||
return [(int(p.id), float(p.score)) for p in res.points]
|
||||
except Exception as e:
|
||||
logger.error(f"Ошибка поиска Qdrant: {e}")
|
||||
return []
|
||||
|
||||
@classmethod
|
||||
def add_vectors(cls, vectors: np.ndarray, doc_ids: list[int]) -> None:
|
||||
"""Upsert векторов по doc_id (идемпотентно: повторная запись перезаписывает)."""
|
||||
if len(doc_ids) == 0:
|
||||
return
|
||||
client = cls._ensure()
|
||||
|
||||
from qdrant_client.models import PointStruct
|
||||
|
||||
vectors = np.asarray(vectors, dtype=np.float32)
|
||||
points = [
|
||||
PointStruct(id=int(doc_id), vector=vectors[i].tolist())
|
||||
for i, doc_id in enumerate(doc_ids)
|
||||
]
|
||||
client.upsert(collection_name=cls.COLLECTION, points=points)
|
||||
logger.info(f"Upsert {len(points)} векторов в Qdrant. Всего: {cls.total_vectors()}")
|
||||
|
||||
@classmethod
|
||||
def save(cls) -> None:
|
||||
"""No-op: Qdrant персистит данные на своей стороне."""
|
||||
|
||||
@classmethod
|
||||
def total_vectors(cls) -> int:
|
||||
"""Количество точек в коллекции."""
|
||||
try:
|
||||
return cls._ensure().count(collection_name=cls.COLLECTION).count
|
||||
except Exception:
|
||||
return 0
|
||||
|
||||
@classmethod
|
||||
def _reset(cls) -> None:
|
||||
"""Сбросить клиент и флаг готовности (для тестов)."""
|
||||
cls._client = None
|
||||
cls._ready = False
|
||||
@@ -92,11 +92,12 @@ def check_plagiarism(
|
||||
session.commit()
|
||||
|
||||
try:
|
||||
from app.faiss_manager import FAISSManager
|
||||
from app.model_manager import ModelManager
|
||||
from app.ollama_client import OllamaClient
|
||||
from app.vector_store import get_backend
|
||||
|
||||
ollama = OllamaClient()
|
||||
vector_store = get_backend()
|
||||
semantic_matches: list[dict] = []
|
||||
|
||||
for i, fragment in enumerate(fragments):
|
||||
@@ -105,9 +106,9 @@ def check_plagiarism(
|
||||
# Пропустить слишком короткие фрагменты
|
||||
continue
|
||||
|
||||
# Уровень 3: Семантический поиск через FAISS
|
||||
# Уровень 3: Семантический поиск (FAISS или Qdrant по VECTOR_BACKEND)
|
||||
frag_vec = ModelManager.encode_single(frag_text)
|
||||
faiss_results = FAISSManager.search(frag_vec, k=10)
|
||||
faiss_results = vector_store.search(frag_vec, k=10)
|
||||
|
||||
for doc_id, score in faiss_results:
|
||||
if score < FAISS_SIMILARITY_THRESHOLD:
|
||||
@@ -201,9 +202,9 @@ def embed_documents(doc_ids: list[int]) -> dict[str, Any]:
|
||||
if not doc_ids:
|
||||
return {"status": "ok", "embedded": 0}
|
||||
|
||||
from app.faiss_manager import FAISSManager
|
||||
from app.model_manager import ModelManager
|
||||
from app.models import Document
|
||||
from app.vector_store import get_backend
|
||||
|
||||
logger.info(f"Построение эмбеддингов для {len(doc_ids)} документов...")
|
||||
|
||||
@@ -225,8 +226,9 @@ def embed_documents(doc_ids: list[int]) -> dict[str, Any]:
|
||||
|
||||
vectors = ModelManager.encode(texts)
|
||||
|
||||
FAISSManager.add_vectors(vectors, ids)
|
||||
FAISSManager.save()
|
||||
store = get_backend()
|
||||
store.add_vectors(vectors, ids)
|
||||
store.save()
|
||||
|
||||
# Обновить faiss_id в PostgreSQL (для IDMap2 faiss_id == doc_id)
|
||||
with db_session() as session:
|
||||
|
||||
@@ -222,10 +222,10 @@ def search_semantic(
|
||||
from app.model_manager import ModelManager
|
||||
query_vec = ModelManager.encode_single(query_normalized)
|
||||
|
||||
# FAISS семантический поиск
|
||||
from app.faiss_manager import FAISSManager
|
||||
faiss_results = FAISSManager.search(query_vec, k=50)
|
||||
logger.info(f"FAISS: найдено {len(faiss_results)} результатов")
|
||||
# Семантический поиск (FAISS или Qdrant по VECTOR_BACKEND)
|
||||
from app.vector_store import get_backend
|
||||
faiss_results = get_backend().search(query_vec, k=50)
|
||||
logger.info(f"Векторный поиск: найдено {len(faiss_results)} результатов")
|
||||
|
||||
# Elasticsearch BM25 поиск
|
||||
from app.es_client import search_fulltext
|
||||
|
||||
25
services/worker-gpu/app/vector_store.py
Normal file
25
services/worker-gpu/app/vector_store.py
Normal file
@@ -0,0 +1,25 @@
|
||||
"""Выбор бэкенда векторного индекса по настройке VECTOR_BACKEND.
|
||||
|
||||
Единая точка входа: задачи (search / plagiarism) зовут get_backend().search(...) /
|
||||
.add_vectors(...) / .save(...), не зная, какой бэкенд под капотом — FAISS (файловый
|
||||
синглтон в RAM воркера, дефолт) или Qdrant (сетевой сервис, снимает SPOF и допускает
|
||||
конкурентную запись). Оба класса имеют одинаковый classmethod-интерфейс.
|
||||
"""
|
||||
|
||||
import logging
|
||||
|
||||
from app.config import settings
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
def get_backend():
|
||||
"""Вернуть класс активного векторного бэкенда согласно settings.VECTOR_BACKEND."""
|
||||
if settings.VECTOR_BACKEND == "qdrant":
|
||||
from app.qdrant_manager import QdrantManager
|
||||
|
||||
return QdrantManager
|
||||
|
||||
from app.faiss_manager import FAISSManager
|
||||
|
||||
return FAISSManager
|
||||
Reference in New Issue
Block a user