feat(indexer): MinHash LSH (уровень 2) в общем Redis вместо памяти
Раньше LSH-индекс жил in-memory в процессе воркера: терялся при рестарте и не шарился между воркерами — уровень 2 в проде фактически не работал. Теперь индекс хранится в общем Redis (тот же REDIS_URL) через нативный storage_config datasketch: - переживает рестарт, общий для всех воркеров-индексеров; - все ключи под префиксом antiplag_lsh — изолированы от кэша/rate-limits и чужих данных в общей БД; никаких FLUSH (и ACL-юзер их не умеет); - add_to_lsh теперь upsert (remove+insert) — full-text версия документа корректно заменяет провизорную по аннотации (был латентный баг: skip-on- duplicate оставлял абстрактную версию); - graceful-фолбэк в in-memory, если Redis недоступен, чтобы воркер не падал. Проверено на боевом Redis через изолированный тестовый префикс: query находит похожие, все ключи в неймспейсе, точечная чистка вернула БД к исходному состоянию (чужие данные не затронуты). Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
@@ -1,29 +1,73 @@
|
|||||||
"""MinHash LSH для нечёткого поиска похожих документов.
|
"""MinHash LSH для нечёткого поиска похожих документов (уровень 2).
|
||||||
|
|
||||||
Позволяет быстро находить документы с похожим содержимым
|
Индекс хранится в общем Redis (тот же, что кэш/rate-limits, см. REDIS_URL),
|
||||||
без точного сравнения всех пар.
|
а не в памяти процесса. Это даёт:
|
||||||
|
- переживаемость рестарта воркера (индекс не теряется);
|
||||||
|
- общий индекс для всех воркеров-индексеров (add в одном видно в query другого).
|
||||||
|
|
||||||
|
При недоступности Redis — graceful-фолбэк в in-memory (деградация: индекс
|
||||||
|
локальный и теряется при рестарте), чтобы воркер не падал целиком.
|
||||||
"""
|
"""
|
||||||
|
|
||||||
import logging
|
import logging
|
||||||
|
from urllib.parse import urlparse
|
||||||
|
|
||||||
from datasketch import MinHash, MinHashLSH
|
from datasketch import MinHash, MinHashLSH
|
||||||
|
|
||||||
|
from app.config import settings
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
# Параметры MinHash LSH
|
# Параметры MinHash LSH
|
||||||
LSH_THRESHOLD = 0.5 # Минимальная схожесть для включения в результаты
|
LSH_THRESHOLD = 0.5 # Минимальная схожесть для включения в результаты
|
||||||
LSH_NUM_PERM = 128 # Количество хэш-функций (точность vs память)
|
LSH_NUM_PERM = 128 # Количество хэш-функций (точность vs память)
|
||||||
|
# Префикс ключей в общем Redis — изолирует индекс от кэша/rate-limits в той же БД
|
||||||
|
LSH_BASENAME = b"antiplag_lsh"
|
||||||
|
|
||||||
# Глобальный LSH индекс (in-memory)
|
|
||||||
_lsh: MinHashLSH | None = None
|
_lsh: MinHashLSH | None = None
|
||||||
|
_redis_unavailable = False # чтобы не долбить недоступный Redis на каждый вызов
|
||||||
|
|
||||||
|
|
||||||
|
def _redis_storage_config() -> dict | None:
|
||||||
|
"""Собрать storage_config datasketch из REDIS_URL. None — если разбор не удался."""
|
||||||
|
try:
|
||||||
|
u = urlparse(settings.REDIS_URL)
|
||||||
|
redis_kwargs: dict = {
|
||||||
|
"host": u.hostname or "localhost",
|
||||||
|
"port": u.port or 6379,
|
||||||
|
"db": int((u.path or "/0").lstrip("/") or 0),
|
||||||
|
}
|
||||||
|
if u.username:
|
||||||
|
redis_kwargs["username"] = u.username
|
||||||
|
if u.password:
|
||||||
|
redis_kwargs["password"] = u.password
|
||||||
|
return {"type": "redis", "basename": LSH_BASENAME, "redis": redis_kwargs}
|
||||||
|
except Exception as e:
|
||||||
|
logger.error(f"MinHash LSH: не удалось разобрать REDIS_URL: {e}")
|
||||||
|
return None
|
||||||
|
|
||||||
|
|
||||||
def get_lsh() -> MinHashLSH:
|
def get_lsh() -> MinHashLSH:
|
||||||
"""Получить или создать глобальный LSH индекс."""
|
"""Получить или создать глобальный LSH индекс (Redis-backed, фолбэк — память)."""
|
||||||
global _lsh
|
global _lsh, _redis_unavailable
|
||||||
if _lsh is None:
|
if _lsh is not None:
|
||||||
|
return _lsh
|
||||||
|
|
||||||
|
if not _redis_unavailable:
|
||||||
|
cfg = _redis_storage_config()
|
||||||
|
if cfg is not None:
|
||||||
|
try:
|
||||||
|
_lsh = MinHashLSH(
|
||||||
|
threshold=LSH_THRESHOLD, num_perm=LSH_NUM_PERM, storage_config=cfg
|
||||||
|
)
|
||||||
|
logger.info("MinHash LSH: общий Redis-бэкенд подключён")
|
||||||
|
return _lsh
|
||||||
|
except Exception as e:
|
||||||
|
logger.error(f"MinHash LSH: Redis недоступен ({e}) — фолбэк в память")
|
||||||
|
_redis_unavailable = True
|
||||||
|
|
||||||
_lsh = MinHashLSH(threshold=LSH_THRESHOLD, num_perm=LSH_NUM_PERM)
|
_lsh = MinHashLSH(threshold=LSH_THRESHOLD, num_perm=LSH_NUM_PERM)
|
||||||
logger.info("MinHash LSH индекс создан")
|
logger.warning("MinHash LSH: in-memory режим (индекс не шарится и теряется при рестарте)")
|
||||||
return _lsh
|
return _lsh
|
||||||
|
|
||||||
|
|
||||||
@@ -62,7 +106,10 @@ def text_to_minhash(text: str) -> MinHash:
|
|||||||
|
|
||||||
def add_to_lsh(doc_key: str, text: str) -> None:
|
def add_to_lsh(doc_key: str, text: str) -> None:
|
||||||
"""
|
"""
|
||||||
Добавить документ в LSH индекс.
|
Добавить (или обновить) документ в LSH индекс.
|
||||||
|
|
||||||
|
Upsert: если документ уже есть (например, был добавлен по аннотации, а теперь
|
||||||
|
пересчитывается по полному тексту) — старая подпись удаляется, вставляется новая.
|
||||||
|
|
||||||
Args:
|
Args:
|
||||||
doc_key: Уникальный ключ документа (например, "doc:{id}")
|
doc_key: Уникальный ключ документа (например, "doc:{id}")
|
||||||
@@ -71,10 +118,13 @@ def add_to_lsh(doc_key: str, text: str) -> None:
|
|||||||
lsh = get_lsh()
|
lsh = get_lsh()
|
||||||
m = text_to_minhash(text)
|
m = text_to_minhash(text)
|
||||||
try:
|
try:
|
||||||
|
try:
|
||||||
|
lsh.remove(doc_key) # снять прежнюю версию, если была
|
||||||
|
except Exception:
|
||||||
|
pass # ключа не было — это норма
|
||||||
lsh.insert(doc_key, m)
|
lsh.insert(doc_key, m)
|
||||||
except ValueError:
|
except Exception as e:
|
||||||
# Документ уже в индексе — игнорируем
|
logger.warning(f"MinHash LSH: не удалось добавить {doc_key!r}: {e}")
|
||||||
pass
|
|
||||||
|
|
||||||
|
|
||||||
def find_similar(text: str) -> list[str]:
|
def find_similar(text: str) -> list[str]:
|
||||||
@@ -113,6 +163,7 @@ def compute_jaccard_minhash(text_a: str, text_b: str) -> float:
|
|||||||
|
|
||||||
|
|
||||||
def reset_lsh() -> None:
|
def reset_lsh() -> None:
|
||||||
"""Сбросить LSH индекс (для тестов)."""
|
"""Сбросить handle LSH (для тестов). Данные в Redis не трогает."""
|
||||||
global _lsh
|
global _lsh, _redis_unavailable
|
||||||
_lsh = None
|
_lsh = None
|
||||||
|
_redis_unavailable = False
|
||||||
|
|||||||
Reference in New Issue
Block a user