diff --git a/services/worker-gpu/app/config.py b/services/worker-gpu/app/config.py index 1406834..91737f8 100644 --- a/services/worker-gpu/app/config.py +++ b/services/worker-gpu/app/config.py @@ -58,7 +58,10 @@ class Settings(BaseSettings): EMBED_BACKEND: str = "ollama" EMBED_MODEL: str = "bge-m3" EMBED_DEVICE: str = "cuda" # используется только при EMBED_BACKEND=sentence_transformers - EMBED_BATCH_SIZE: int = 64 + # Ollama обрабатывает тексты в пачке последовательно (один "slot" в llama.cpp, + # не параллельно) — 64 реальных документа в один HTTP-запрос регулярно не + # укладывались в таймаут. Меньше пачка — короче и надёжнее каждый запрос. + EMBED_BATCH_SIZE: int = 16 EMBED_DIM: int = 1024 # bge-m3; было 768 у paraphrase-multilingual-mpnet-base-v2 # DeepInfra — OpenAI-совместимый инференс BAAI/bge-m3 в облаке, та же модель, # что и на VM109, дешёвый пей-пер-токен. Используется при EMBED_BACKEND=cloud. diff --git a/services/worker-gpu/app/model_manager.py b/services/worker-gpu/app/model_manager.py index 2e4f091..826eb51 100644 --- a/services/worker-gpu/app/model_manager.py +++ b/services/worker-gpu/app/model_manager.py @@ -91,7 +91,10 @@ class ModelManager: response = httpx.post( f"{settings.OLLAMA_URL}/api/embed", json={"model": settings.EMBED_MODEL, "input": batch}, - timeout=120.0, + timeout=300.0, # Ollama обрабатывает тексты в пачке последовательно + # (виден один "slot" в логах llama.cpp), не параллельно — на реальных + # (не тестовых) текстах 64 шт. в 120с не укладывались, отсюда ReadTimeout + # и потеря батча целиком (без ретрая). ) response.raise_for_status() all_vectors.extend(response.json()["embeddings"]) diff --git a/services/worker-gpu/app/tasks/plagiarism.py b/services/worker-gpu/app/tasks/plagiarism.py index 14fc74f..0f38c75 100644 --- a/services/worker-gpu/app/tasks/plagiarism.py +++ b/services/worker-gpu/app/tasks/plagiarism.py @@ -219,8 +219,13 @@ def check_plagiarism( raise self.retry(exc=exc, countdown=120) from exc -@celery_app.task(name="gpu.embed_documents") -def embed_documents(doc_ids: list[int]) -> dict[str, Any]: +@celery_app.task( + name="gpu.embed_documents", + bind=True, + max_retries=3, + default_retry_delay=30, +) +def embed_documents(self, doc_ids: list[int]) -> dict[str, Any]: """ Построить эмбеддинги для документов и добавить их в FAISS индекс. @@ -254,7 +259,14 @@ def embed_documents(doc_ids: list[int]) -> dict[str, Any]: ] ids = [d.id for d in docs] - vectors = ModelManager.encode(texts) + try: + vectors = ModelManager.encode(texts) + except Exception as exc: + # Таймаут/недоступность бэкенда эмбеддингов — не терять батч молча, + # а повторить (раньше падало без ретрая, документы просто выпадали + # из переиндексации). + logger.warning(f"Не удалось построить эмбеддинги для {ids}: {exc}") + raise self.retry(exc=exc) from exc store = get_backend() store.add_vectors(vectors, ids)