From 012b17304f1b49976174aebd5f07651bb1142c07 Mon Sep 17 00:00:00 2001 From: jze9 Date: Tue, 11 Aug 2026 20:34:43 +0500 Subject: [PATCH] =?UTF-8?q?feat(gpu):=20Qdrant=20=D0=BA=D0=B0=D0=BA=20?= =?UTF-8?q?=D0=B2=D0=B5=D0=BA=D1=82=D0=BE=D1=80=D0=BD=D1=8B=D0=B9=20=D0=B1?= =?UTF-8?q?=D1=8D=D0=BA=D0=B5=D0=BD=D0=B4=20=D0=BF=D0=BE=D0=B4=20=D1=84?= =?UTF-8?q?=D0=BB=D0=B0=D0=B3=D0=BE=D0=BC=20=E2=80=94=20=D1=81=D0=BD=D0=B8?= =?UTF-8?q?=D0=BC=D0=B0=D0=B5=D1=82=20SPOF=20FAISS?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 --- README.md | 25 +++++ docker-compose.prod.yml | 16 +++ services/worker-gpu/app/config.py | 6 + .../worker-gpu/app/migrate_faiss_to_qdrant.py | 47 ++++++++ services/worker-gpu/app/qdrant_manager.py | 104 ++++++++++++++++++ services/worker-gpu/app/tasks/plagiarism.py | 14 ++- services/worker-gpu/app/tasks/search.py | 8 +- services/worker-gpu/app/vector_store.py | 25 +++++ services/worker-gpu/requirements-test.txt | 1 + services/worker-gpu/requirements.txt | 1 + .../worker-gpu/tests/test_qdrant_manager.py | 83 ++++++++++++++ .../worker-gpu/tests/test_vector_store.py | 25 +++++ 12 files changed, 345 insertions(+), 10 deletions(-) create mode 100644 services/worker-gpu/app/migrate_faiss_to_qdrant.py create mode 100644 services/worker-gpu/app/qdrant_manager.py create mode 100644 services/worker-gpu/app/vector_store.py create mode 100644 services/worker-gpu/tests/test_qdrant_manager.py create mode 100644 services/worker-gpu/tests/test_vector_store.py diff --git a/README.md b/README.md index 1fbc799..2d3276a 100644 --- a/README.md +++ b/README.md @@ -159,6 +159,31 @@ make migrate Nginx работает внутри контейнера (`docker-compose.prod.yml`, сервис `nginx`, Dockerfile в `infra/nginx/Dockerfile`) — собирает `services/frontend` в статику и отдаёт её вместе с проксированием `/api/`, `/ws/` на `api:8000` по конфигу `infra/nginx/nginx.conf`. Контейнер монтирует `/etc/letsencrypt` с хоста как read-only — сертификат обновляется на хосте (`certbot renew`), контейнер просто читает его. +## Векторный бэкенд (FAISS / Qdrant) + +Семантический индекс (уровень 3) спрятан за `app.vector_store.get_backend()` и +переключается настройкой `VECTOR_BACKEND` — без изменения кода: + +- **`faiss`** (по умолчанию) — файловый `IndexIDMap2(IndexFlatIP)` в RAM воркера. + Просто, но это единая точка отказа и без конкурентной записи. +- **`qdrant`** — сетевой сервис: снимает SPOF, допускает конкурентный upsert из + нескольких воркеров, переживает рестарт, масштабируется горизонтально. + +Переключение на Qdrant (аддитивно, ничего не ломает до шага 2): + +```bash +# 1. Поднять Qdrant (профиль qdrant в docker-compose.prod.yml) +docker compose -f docker-compose.prod.yml --profile qdrant up -d qdrant + +# 2. В .env выставить VECTOR_BACKEND=qdrant (QDRANT_URL по умолчанию http://qdrant:6333) + +# 3. Перелить существующие векторы FAISS → Qdrant (идемпотентно) +docker compose -f docker-compose.prod.yml exec worker-gpu python -m app.migrate_faiss_to_qdrant + +# 4. Перезапустить GPU-воркер +docker compose -f docker-compose.prod.yml up -d worker-gpu +``` + ## Тестирование и качество кода Перед деплоем CI (`.gitea/workflows/deploy.yml`, job `test`) прогоняет два гейта, diff --git a/docker-compose.prod.yml b/docker-compose.prod.yml index 9bb00d6..dd4e792 100644 --- a/docker-compose.prod.yml +++ b/docker-compose.prod.yml @@ -13,6 +13,7 @@ networks: volumes: es_prod_data: faiss_index_prod: + qdrant_prod_data: x-app-env: &app-env env_file: .env @@ -49,6 +50,21 @@ services: retries: 10 start_period: 60s + # Qdrant — альтернативный векторный бэкенд (снимает SPOF файлового FAISS). + # Опционален: поднимается только с профилем и активируется VECTOR_BACKEND=qdrant. + # 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 + # Тег сервера при необходимости поднять до версии клиента (qdrant-client 1.19). + qdrant: + image: qdrant/qdrant:v1.12.4 + container_name: antiplagiator-qdrant + profiles: ["qdrant"] + volumes: + - qdrant_prod_data:/qdrant/storage + networks: + - antiplagiator + restart: unless-stopped + # ─── Приложение ──────────────────────────────────────────────────────────── # RabbitMQ вынесен на отдельный сервер 192.168.20.82 (см. RABBITMQ_URL в .env) api: diff --git a/services/worker-gpu/app/config.py b/services/worker-gpu/app/config.py index a5ee271..9a0076d 100644 --- a/services/worker-gpu/app/config.py +++ b/services/worker-gpu/app/config.py @@ -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" diff --git a/services/worker-gpu/app/migrate_faiss_to_qdrant.py b/services/worker-gpu/app/migrate_faiss_to_qdrant.py new file mode 100644 index 0000000..b3d5eed --- /dev/null +++ b/services/worker-gpu/app/migrate_faiss_to_qdrant.py @@ -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() diff --git a/services/worker-gpu/app/qdrant_manager.py b/services/worker-gpu/app/qdrant_manager.py new file mode 100644 index 0000000..9375c42 --- /dev/null +++ b/services/worker-gpu/app/qdrant_manager.py @@ -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 diff --git a/services/worker-gpu/app/tasks/plagiarism.py b/services/worker-gpu/app/tasks/plagiarism.py index 8cddab4..b23c2ab 100644 --- a/services/worker-gpu/app/tasks/plagiarism.py +++ b/services/worker-gpu/app/tasks/plagiarism.py @@ -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: diff --git a/services/worker-gpu/app/tasks/search.py b/services/worker-gpu/app/tasks/search.py index 8f8a40b..746a2c4 100644 --- a/services/worker-gpu/app/tasks/search.py +++ b/services/worker-gpu/app/tasks/search.py @@ -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 diff --git a/services/worker-gpu/app/vector_store.py b/services/worker-gpu/app/vector_store.py new file mode 100644 index 0000000..82b980b --- /dev/null +++ b/services/worker-gpu/app/vector_store.py @@ -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 diff --git a/services/worker-gpu/requirements-test.txt b/services/worker-gpu/requirements-test.txt index 7f96039..eee1795 100644 --- a/services/worker-gpu/requirements-test.txt +++ b/services/worker-gpu/requirements-test.txt @@ -6,3 +6,4 @@ faiss-cpu==1.8.0 numpy==1.26.4 pydantic-settings==2.2.1 httpx==0.27.0 +qdrant-client==1.19.0 # local :memory:-режим → тесты бэкенда без сервера diff --git a/services/worker-gpu/requirements.txt b/services/worker-gpu/requirements.txt index 37e8d70..86bef1c 100644 --- a/services/worker-gpu/requirements.txt +++ b/services/worker-gpu/requirements.txt @@ -4,6 +4,7 @@ sqlalchemy==2.0.30 psycopg2-binary==2.9.9 sentence-transformers==3.0.0 faiss-cpu==1.8.0 # faiss-gpu нет в pip для py3.11; индексы в коде CPU-типа, GPU занят эмбеддингами (torch) и LLM (Ollama) +qdrant-client==1.19.0 # альтернативный векторный бэкенд (VECTOR_BACKEND=qdrant), снимает SPOF FAISS torch==2.3.0 numpy==1.26.4 httpx==0.27.0 diff --git a/services/worker-gpu/tests/test_qdrant_manager.py b/services/worker-gpu/tests/test_qdrant_manager.py new file mode 100644 index 0000000..992cd0d --- /dev/null +++ b/services/worker-gpu/tests/test_qdrant_manager.py @@ -0,0 +1,83 @@ +"""Юнит-тесты QdrantManager против встроенного Qdrant (:memory:) — реальный бэкенд. + +Не моки: qdrant-client в local-режиме поднимает in-process Qdrant, поэтому +проверяется настоящее поведение коллекции (косинус, upsert, count). +""" + +import numpy as np +import pytest + +pytest.importorskip("qdrant_client") + +from app.config import settings # noqa: E402 +from app.qdrant_manager import QdrantManager # noqa: E402 +from qdrant_client import QdrantClient # noqa: E402 + +DIM = settings.EMBED_DIM + + +def _unit(vecs: np.ndarray) -> np.ndarray: + vecs = np.asarray(vecs, dtype=np.float32) + norms = np.linalg.norm(vecs, axis=1, keepdims=True) + norms[norms == 0] = 1.0 + return vecs / norms + + +@pytest.fixture(autouse=True) +def in_memory_backend(): + """Свежий встроенный Qdrant на каждый тест.""" + QdrantManager._reset() + QdrantManager._client = QdrantClient(location=":memory:") + yield + QdrantManager._reset() + + +def test_empty_search_returns_nothing(): + assert QdrantManager.search(np.zeros(DIM, dtype=np.float32)) == [] + assert QdrantManager.total_vectors() == 0 + + +def test_add_empty_is_noop(): + QdrantManager.add_vectors(np.zeros((0, DIM), dtype=np.float32), []) + assert QdrantManager.total_vectors() == 0 + + +def test_add_and_self_search_returns_postgres_doc_id(): + rng = np.random.default_rng(42) + vecs = _unit(rng.standard_normal((3, DIM))) + QdrantManager.add_vectors(vecs, [101, 202, 303]) + assert QdrantManager.total_vectors() == 3 + + res = QdrantManager.search(vecs[1], k=1) + assert res, "поиск ничего не вернул" + top_id, score = res[0] + assert top_id == 202 # id точки == doc_id из PostgreSQL + assert score == pytest.approx(1.0, abs=1e-3) # self-match: косинус ≈ 1 + + +def test_upsert_is_idempotent_by_doc_id(): + rng = np.random.default_rng(0) + QdrantManager.add_vectors(_unit(rng.standard_normal((2, DIM))), [1, 2]) + + updated = _unit(rng.standard_normal((2, DIM))) + QdrantManager.add_vectors(updated, [1, 2]) # те же id → перезапись, не дубли + assert QdrantManager.total_vectors() == 2 + + top_id, score = QdrantManager.search(updated[0], k=1)[0] + assert top_id == 1 + assert score == pytest.approx(1.0, abs=1e-3) + + +def test_search_ranks_nearest_first(): + rng = np.random.default_rng(7) + v = _unit(rng.standard_normal((1, DIM)))[0] + near = _unit((v + 0.01 * rng.standard_normal(DIM))[None, :])[0] + far = _unit(rng.standard_normal((1, DIM)))[0] + QdrantManager.add_vectors(np.stack([near, far]), [10, 20]) + + res = QdrantManager.search(v, k=2) + assert res[0][0] == 10 # ближайший — near + + +def test_save_is_noop(): + QdrantManager.save() # Qdrant персистит сам — вызов не должен падать diff --git a/services/worker-gpu/tests/test_vector_store.py b/services/worker-gpu/tests/test_vector_store.py new file mode 100644 index 0000000..ee000fd --- /dev/null +++ b/services/worker-gpu/tests/test_vector_store.py @@ -0,0 +1,25 @@ +"""Тест выбора векторного бэкенда по настройке VECTOR_BACKEND.""" + +from app import vector_store +from app.config import settings + + +def test_default_backend_is_faiss(monkeypatch): + monkeypatch.setattr(settings, "VECTOR_BACKEND", "faiss") + from app.faiss_manager import FAISSManager + + assert vector_store.get_backend() is FAISSManager + + +def test_qdrant_backend_selected(monkeypatch): + monkeypatch.setattr(settings, "VECTOR_BACKEND", "qdrant") + from app.qdrant_manager import QdrantManager + + assert vector_store.get_backend() is QdrantManager + + +def test_unknown_backend_falls_back_to_faiss(monkeypatch): + monkeypatch.setattr(settings, "VECTOR_BACKEND", "что-то не то") + from app.faiss_manager import FAISSManager + + assert vector_store.get_backend() is FAISSManager