Files
anti-plagiarism/services/worker-gpu/app/qdrant_manager.py
jze9 012b17304f 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>
2026-08-11 20:34:43 +05:00

105 lines
4.3 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""Менеджер векторного индекса на 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