Files
jze9 b0e1eadcbb 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>
2026-07-29 13:09:51 +05:00

170 lines
6.1 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.
"""MinHash LSH для нечёткого поиска похожих документов (уровень 2).
Индекс хранится в общем Redis (тот же, что кэш/rate-limits, см. REDIS_URL),
а не в памяти процесса. Это даёт:
- переживаемость рестарта воркера (индекс не теряется);
- общий индекс для всех воркеров-индексеров (add в одном видно в query другого).
При недоступности Redis — graceful-фолбэк в in-memory (деградация: индекс
локальный и теряется при рестарте), чтобы воркер не падал целиком.
"""
import logging
from urllib.parse import urlparse
from datasketch import MinHash, MinHashLSH
from app.config import settings
logger = logging.getLogger(__name__)
# Параметры MinHash LSH
LSH_THRESHOLD = 0.5 # Минимальная схожесть для включения в результаты
LSH_NUM_PERM = 128 # Количество хэш-функций (точность vs память)
# Префикс ключей в общем Redis — изолирует индекс от кэша/rate-limits в той же БД
LSH_BASENAME = b"antiplag_lsh"
_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:
"""Получить или создать глобальный LSH индекс (Redis-backed, фолбэк — память)."""
global _lsh, _redis_unavailable
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)
logger.warning("MinHash LSH: in-memory режим (индекс не шарится и теряется при рестарте)")
return _lsh
def get_shingles(text: str, k: int = 3) -> set[str]:
"""
Получить k-слоговые шинглы из текста.
Args:
text: Исходный текст
k: Размер шингла (количество слов)
Returns:
Множество шинглов
"""
words = text.lower().split()
if len(words) < k:
return {" ".join(words)} if words else set()
return {" ".join(words[i : i + k]) for i in range(len(words) - k + 1)}
def text_to_minhash(text: str) -> MinHash:
"""
Создать MinHash подпись для текста.
Args:
text: Исходный текст
Returns:
MinHash объект
"""
m = MinHash(num_perm=LSH_NUM_PERM)
for shingle in get_shingles(text):
m.update(shingle.encode("utf-8"))
return m
def add_to_lsh(doc_key: str, text: str) -> None:
"""
Добавить (или обновить) документ в LSH индекс.
Upsert: если документ уже есть (например, был добавлен по аннотации, а теперь
пересчитывается по полному тексту) — старая подпись удаляется, вставляется новая.
Args:
doc_key: Уникальный ключ документа (например, "doc:{id}")
text: Текст документа
"""
lsh = get_lsh()
m = text_to_minhash(text)
try:
try:
lsh.remove(doc_key) # снять прежнюю версию, если была
except Exception:
pass # ключа не было — это норма
lsh.insert(doc_key, m)
except Exception as e:
logger.warning(f"MinHash LSH: не удалось добавить {doc_key!r}: {e}")
def find_similar(text: str) -> list[str]:
"""
Найти похожие документы в LSH индексе.
Args:
text: Текст для поиска похожих
Returns:
Список ключей похожих документов
"""
lsh = get_lsh()
m = text_to_minhash(text)
try:
return lsh.query(m)
except Exception as e:
logger.warning(f"Ошибка запроса к MinHash LSH: {e}")
return []
def compute_jaccard_minhash(text_a: str, text_b: str) -> float:
"""
Оценить схожесть двух текстов через MinHash Jaccard.
Args:
text_a: Первый текст
text_b: Второй текст
Returns:
Оценка схожести Жаккара через MinHash
"""
m_a = text_to_minhash(text_a)
m_b = text_to_minhash(text_b)
return m_a.jaccard(m_b)
def reset_lsh() -> None:
"""Сбросить handle LSH (для тестов). Данные в Redis не трогает."""
global _lsh, _redis_unavailable
_lsh = None
_redis_unavailable = False