hojichar.filters.deduplication
This module provides filters for deduplication using MinHash and Locality-Sensitive Hashing (LSH).
If you want to use this module, install hojichar via pip install 'hojichar[dedup]'
What is MinHash LSH?
A gentle introduction for first‑time leaner for MinHash LSH
1 What problem does this solve?
When you have millions of documents it is too expensive to compare every pair directly. MinHash + Locality‑Sensitive Hashing (LSH) lets you
- estimate Jaccard similarity extremely fast, and memory‑efficiently
- retrieve near‑duplicates in sub‑linear time.
In practice you can keep a single set or Redis index of the generated LSH keys and ask:
"Does any existing document share at least one LSH key with mine?"
If the answer is yes the two documents are almost certainly similar; if no they are very likely different.
How the pipeline works
GenerateDedupLSH filter generates LSH keys for each document.
- Tokenize
- Split text into tokens. (Default: character‑level, but you can plug in any callable.)
- n-grams
- Group tokens into n‑grams (n_grams=5) to capture context.
- MinHash
- Hash each n‑gram with
num_permindependent permutations and keep only the minimum value. The resulting signature is a vector of num_perm 32‑bit integers. - This module uses
rensa.RMinHashwhich is a fast MinHash implementation by Rust language.
- Hash each n‑gram with
- Banding and compression
- Split the signature into b bands, each containing
rintegers (num_perm ≈ b×r). - Treat the
rintegers as raw bytes, hash them with xxhash‑128, and format as v2:+ .
- Split the signature into b bands, each containing
- Output
- Store all band keys in
document.extras['dedup_lsh']as a list of strings.
- Store all band keys in
InlineDeduplicator, RedisDeduplicator, and RedisBloomDeduplicator filters use the generated LSH keys to mark documents as duplicates.
InlineDeduplicatorstores LSH keys in a local set, so it works only in a single process.RedisDeduplicatorstores LSH keys in Redis, so it works in a distributed environment.RedisBloomDeduplicatorstores LSH keys in RedisBloom, which is a scalable Bloom filter. It uses less memory than Redis keys but may return false positives.
1""" 2This module provides filters for deduplication using MinHash and Locality-Sensitive Hashing (LSH). 3If you want to use this module, install hojichar via `pip install 'hojichar[dedup]'` 4 5## What is MinHash LSH? 6*A gentle introduction for first‑time leaner for MinHash LSH* 7 8--- 9 10### 1 What problem does this solve? 11 12When you have **millions of documents** it is too expensive to compare every pair directly. 13**MinHash + Locality‑Sensitive Hashing (LSH)** lets you 14 15- estimate Jaccard similarity extremely fast, and memory‑efficiently 16- retrieve **near‑duplicates** in sub‑linear time. 17 18In practice you can keep a single `set` or Redis index of the generated **LSH keys** and ask: 19"Does any existing document share *at least one* LSH key with mine?" 20 21If the answer is yes the two documents are almost certainly similar; if no they are very likely different. 22 23### How the pipeline works 24 25`GenerateDedupLSH` filter generates LSH keys for each document. 26 271. Tokenize 28 - Split text into tokens. (Default: character‑level, but you can plug in any callable.) 292. n-grams 30 - Group tokens into n‑grams (n_grams=5) to capture context. 313. MinHash 32 - Hash each n‑gram with `num_perm` independent permutations and keep only the minimum value. The resulting **signature** is a vector of num_perm 32‑bit integers. 33 - This module uses `rensa.RMinHash` which is a fast MinHash implementation by Rust language. 344. Banding and compression 35 - Split the signature into b bands, each containing `r` integers (num_perm ≈ b×r). 36 - Treat the `r` integers as raw bytes, hash them with xxhash‑128, and format as v2:<band_idx>+<digest32hex>. 375. Output 38 - Store all band keys in `document.extras['dedup_lsh']` as a list of strings. 39 40`InlineDeduplicator`, `RedisDeduplicator`, and `RedisBloomDeduplicator` filters use the generated LSH keys to mark documents as duplicates. 41 42- `InlineDeduplicator` stores LSH keys in a local set, so it works only in a single process. 43- `RedisDeduplicator` stores LSH keys in Redis, so it works in a distributed environment. 44- `RedisBloomDeduplicator` stores LSH keys in RedisBloom, which is a scalable Bloom filter. It uses less memory than Redis keys but may return false positives. 45 46""" 47 48from __future__ import annotations 49 50import importlib 51import re 52import struct 53import sys 54from collections import deque 55from itertools import islice 56from typing import Any, Callable, Final, Iterable, Optional, cast 57 58import numpy as np 59from numpy.typing import NDArray 60 61try: 62 import redis 63 import xxhash 64 from datasketch.lsh import _optimal_param # type: ignore 65 from rensa import RMinHash # type: ignore 66 67 is_loaded_dedup = True 68except ImportError: 69 is_loaded_dedup = False 70 71 72from hojichar import Document, Filter 73 74_japanese_tagger: Optional["fugashi.Tagger"] = None # type: ignore[name-defined] # noqa: F821 75NON_ALPHA = re.compile("[^A-Za-z_0-9]") 76IS_LOADED_DEDUP_ERROR_MSG = ( 77 "Failed to import redis, xxhash, rensa, or datasketch. " 78 "Please install the extra dependencies with `pip install 'hojichar[dedup]'`" 79) 80 81 82def char_level_splitter(text: str) -> list[str]: 83 """ 84 Split the text into characters. 85 This is a simple implementation that splits the text into individual characters. 86 """ 87 return list(text) 88 89 90def non_alpha_num_splitter(text: str) -> list[str]: 91 """ 92 Split the text into alphanumeric tokens. 93 This is a simple implementation that splits on non-alphanumeric characters. 94 """ 95 return [token for token in NON_ALPHA.split(text) if token] 96 97 98def _ngrams(tokens: Iterable[str], n: int) -> Iterable[tuple[str, ...]]: 99 """Yield sliding windows of *n* tokens.""" 100 iterator = iter(tokens) 101 window = deque(islice(iterator, n), maxlen=n) 102 if len(window) == n: 103 yield tuple(window) 104 for token in iterator: 105 window.append(token) 106 yield tuple(window) 107 108 109def japanese_word_splitter(text: str) -> list[str]: 110 """ 111 Split the text into Japanese words using fugashi. 112 This will import fugashi and instantiate Tagger on first use. 113 """ 114 global _japanese_tagger 115 if _japanese_tagger is None: 116 fugashi = importlib.import_module("fugashi") 117 _japanese_tagger = fugashi.Tagger() 118 return [token.surface for token in _japanese_tagger(text)] 119 120 121class GenerateDedupLSH(Filter): 122 """ 123 Filter that uses MinHash + Locality-Sensitive Hashing (LSH) to assign 124 deduplication keys to documents, allowing fast near-duplicate detection. 125 126 Attributes: 127 num_perm (int): Number of permutations (hash functions) for MinHash. 128 threshold (float): Similarity threshold for tuning LSH parameters. 129 tokenizer (Callable[[str], Iterable[str]]): Function to tokenize text. 130 n_grams (int): n-gram size for token grouping. 131 seed (int): Random seed for MinHash. 132 num_bands (int): Number of LSH bands, specified explicitly or computed automatically. 133 band_size (int): Number of hashes per band. 134 135 Notes 136 ----- 137 When band parameters are omitted, `_optimal_param` searches for the optimal 138 number of **bands** (`b`) and 139 **rows per band** (`r`) that minimise a weighted sum of false positives / 140 false negatives at the specified *threshold*. 141 142 HojiChar 0.18.0 introduces ``v2:`` keys using Rensa 0.5 and an optimized 143 hashing path. Hash generation was over 4x faster than HojiChar 0.17.3 144 in our 2,000-character English benchmark. Hashes differ from 145 HojiChar 0.17.x. Please rebuild your deduplication fingerprints, 146 or keep using HojiChar 0.17.x if you want to use existing LSH pool. 147 """ 148 149 _BYTES_PER_U32: Final[int] = 4 150 _LSH_KEY_VERSION: Final[str] = "v2" 151 152 def __init__( 153 self, 154 num_perm: int = 500, 155 threshold: float = 0.8, 156 tokenizer: Callable[[str], Iterable[str]] = char_level_splitter, 157 n_grams: int = 5, 158 seed: int = 42, 159 *, 160 num_bands: Optional[int] = None, 161 band_size: Optional[int] = None, 162 **kwargs: Any, 163 ) -> None: 164 """ 165 Initialize the deduplication filter with MinHash and LSH settings. 166 167 Args: 168 num_perm: Number of hash permutations for MinHash signature length. 169 Ignored when num_bands and band_size are both provided. 170 threshold: Similarity threshold to decide optimal LSH parameters. 171 Ignored when num_bands and band_size are both provided. 172 tokenizer: Function to split text into tokens. 173 n_grams: Number of tokens per n-gram for MinHash update. 174 seed: Seed for hash permutation consistency. 175 num_bands: Explicit number of LSH bands. Must be a positive integer 176 and provided together with band_size. 177 band_size: Explicit number of hashes per band. Must be a positive integer 178 and provided together with num_bands. When both are provided, 179 automatic selection is bypassed and the MinHash signature length 180 (self.num_perm) is set to num_bands * band_size. 181 **kwargs: Additional keyword arguments for parent Filter. 182 """ 183 super().__init__(**kwargs) 184 if not is_loaded_dedup: 185 raise ImportError(IS_LOADED_DEDUP_ERROR_MSG) 186 if n_grams <= 0: 187 raise ValueError("n_grams must be positive") 188 self.num_perm = num_perm 189 self.threshold = threshold 190 self.tokenizer = tokenizer 191 self.n_grams = n_grams 192 self.seed = seed 193 194 if num_bands is None and band_size is None: 195 self.num_bands, self.band_size = _optimal_param( 196 threshold=self.threshold, 197 num_perm=self.num_perm, 198 false_negative_weight=0.5, 199 false_positive_weight=0.5, 200 ) 201 else: 202 if num_bands is None or band_size is None: 203 raise ValueError("num_bands and band_size must be provided together") 204 if num_bands <= 0 or band_size <= 0: 205 raise ValueError("num_bands and band_size must be positive") 206 self.num_bands = num_bands 207 self.band_size = band_size 208 self.num_perm = num_bands * band_size 209 210 def _calculate_minhash_digest(self, text: str) -> list[int]: 211 """Compute the raw digest shared by the array API and the LSH fast path.""" 212 if self.tokenizer is char_level_splitter: 213 n = self.n_grams 214 tokens = [text[i : i + n] for i in range(len(text) - n + 1)] 215 else: 216 tokens = [" ".join(grams) for grams in _ngrams(self.tokenizer(text), self.n_grams)] 217 minhash = RMinHash(num_perm=self.num_perm, seed=self.seed) 218 minhash.update(tokens) 219 return cast(list[int], minhash.digest()) 220 221 def calculate_minhash_signature(self, text: str) -> NDArray[np.uint32]: 222 """ 223 Compute MinHash signature of input text as an array of uint32. 224 225 Steps: 226 1. Tokenize text using the provided tokenizer. 227 2. Generate n-gram tokens. 228 3. Update MinHash with n-gram tokens. 229 230 Args: 231 text: Input document text to be hashed. 232 233 Returns: 234 A 1D numpy array of shape (num_perm,) with dtype uint32. 235 """ 236 return np.asarray(self._calculate_minhash_digest(text), dtype=np.uint32) 237 238 def _sig_bytes_le(self, sig: NDArray[np.uint32]) -> memoryview: 239 """ 240 Return the signature as *little-endian* byte view. 241 Platform‑independent way to get bytes from a numpy array. 242 """ 243 if sys.byteorder == "little": 244 # amd64 / arm64 (little) 245 return memoryview(cast(Any, sig)).cast("B") # zero-copy 246 else: 247 # big-endian CPU 248 return memoryview(cast(Any, sig.byteswap())).cast("B") # 1 copy 249 250 def signature_to_lsh_digest( 251 self, signature: NDArray[np.uint32], band_size: int, band_idx: int 252 ) -> int: 253 """ 254 Convert a slice of the MinHash signature into an LSH digest with less memory overhead. 255 256 This method is optimized for speed by avoiding copies: 257 - We view the uint32 array as raw bytes (uint8 view). 258 - We create a memoryview of the byte slice for the specified band. 259 - We compute a 128-bit hash directly on the slice. 260 261 Args: 262 signature: 1D numpy array of uint32 representing MinHash signature. 263 band_size: Number of hashes per LSH band. 264 band_idx: Index of the band to hash (0-based). 265 266 Returns: 267 An integer representing the 128-bit hash digest of the band. 268 269 Raises: 270 AssertionError: If signature shape/dtype or band index is invalid. 271 """ 272 assert signature.dtype == np.uint32 and signature.ndim == 1, ( 273 "signature must be a 1D numpy array of uint32" 274 ) 275 assert 0 <= band_idx < self.num_bands, ( 276 f"band_idx {band_idx} out of range [0, {self.num_bands})" 277 ) 278 assert len(signature) >= band_size * self.num_bands, ( 279 "signature length is too short for given band_size and num_bands" 280 ) 281 282 # Compute byte offsets for the selected band 283 start = band_idx * band_size * self._BYTES_PER_U32 284 stop = start + band_size * self._BYTES_PER_U32 285 286 # View signature as raw bytes without copy. memoryview avoids creating new bytes. 287 mv = self._sig_bytes_le(signature)[start:stop] # slice view 288 289 return xxhash.xxh128_intdigest(mv) 290 291 def _format_lsh_key(self, band_idx: int, digest: int) -> str: 292 """ 293 Format the LSH key with the HojiChar scheme version, band index, and digest. 294 """ 295 return f"{self._LSH_KEY_VERSION}:{band_idx}+{digest:032x}" 296 297 def apply(self, document: Document) -> Document: 298 """ 299 Decorate the document with LSH deduplication keys. 300 301 For each band, compute the digest and format as a hex string: 302 'v2:<band_idx>+<128-bit-digest-hex>'. 303 Keys are stored in document.extras['dedup_lsh']. 304 305 Args: 306 document: Document object with 'text' attribute. 307 308 Returns: 309 The same Document object with 'dedup_lsh' added in extras. 310 """ 311 digest = self._calculate_minhash_digest(document.text) 312 # Pack once in little-endian order without an intermediate NumPy array. 313 signature_bytes = memoryview(struct.pack(f"<{self.num_perm}I", *digest)) 314 band_bytes = self.band_size * self._BYTES_PER_U32 315 lsh_keys = [ 316 f"{self._LSH_KEY_VERSION}:{band_idx}+" 317 + xxhash.xxh128_hexdigest( 318 signature_bytes[band_idx * band_bytes : (band_idx + 1) * band_bytes] 319 ) 320 for band_idx in range(self.num_bands) 321 ] 322 323 document.extras["dedup_lsh"] = lsh_keys 324 return document 325 326 327class InlineDeduplicator(Filter): 328 """ 329 Simple in‑memory deduplicator. 330 331 Stores every LSH key in a local :pyclass:`set`. If any key of the incoming 332 document is already present, the document is marked as duplicate via 333 `document.is_rejected = True`. 334 335 **Limitations** 336 ------------- 337 *State is per‑process only.* Running multiple workers or machines will *not* 338 share the key set – use :class:`RedisDeduplicator` or 339 :class:`RedisBloomDeduplicator` for distributed setups. 340 """ 341 342 def __init__(self, **kwargs: Any): 343 super().__init__(**kwargs) 344 self.hash_pool: set[str] = set() 345 346 def apply(self, document: Document) -> Document: 347 """ 348 Inline deduplication based on the LSH keys in the document. 349 This filter cannot use in the distributed environment because it uses a local hash pool. 350 """ 351 lsh_keys = document.extras.get("dedup_lsh") 352 if lsh_keys is None: 353 raise ValueError( 354 "Document does not contain LSH keys for deduplication. Please apply GenerateDedupLSH first." 355 ) 356 357 for lsh in lsh_keys: 358 if lsh in self.hash_pool: 359 document.is_rejected = True 360 else: 361 self.hash_pool.add(lsh) 362 return document 363 364 365class RedisDeduplicator(Filter): 366 """ 367 Distributed deduplicator using **plain Redis keys**. 368 You have to run a Redis server and pass its connection parameters. 369 """ 370 371 def __init__( 372 self, 373 *, 374 host: str = "localhost", 375 port: int = 6379, 376 db: int = 0, 377 key_prefix: str = "dedup", 378 **kwargs: Any, 379 ) -> None: 380 """ 381 Initialize the Redis deduplicator. 382 Args: 383 host (str): Redis server hostname. 384 port (int): Redis server port. 385 db (int): Redis database number. 386 key_prefix (str): Prefix for Redis keys to avoid collisions. You should use a unique prefix for each deduplication task. 387 **kwargs: Additional keyword arguments for parent Filter. 388 """ 389 if not is_loaded_dedup: 390 raise ImportError(IS_LOADED_DEDUP_ERROR_MSG) 391 super().__init__(**kwargs) 392 self.rds = redis.Redis(host=host, port=port, db=db, decode_responses=False) 393 self.key_prefix = key_prefix.encode() 394 395 try: 396 self.rds.ping() 397 except redis.exceptions.RedisError as exc: 398 raise RuntimeError(f"Cannot connect to Redis server {host}:{port}/{db}") from exc 399 400 def apply(self, document: Document) -> Document: 401 lsh_keys = document.extras.get("dedup_lsh") 402 if lsh_keys is None: 403 raise ValueError("Apply GenerateDedupLSH first") 404 405 pipe = self.rds.pipeline(transaction=False) 406 for k in lsh_keys: 407 pipe.set(self.key_prefix + b":" + k.encode(), b"1", nx=True) 408 results: list[bool | None] = pipe.execute() # If instance already exists, it returns None 409 410 if any(r is None for r in results): 411 document.is_rejected = True 412 return document 413 414 415class RedisBloomDeduplicator(Filter): 416 """ 417 Distributed deduplicator backed by **RedisBloom scalable Bloom filters**. 418 You can use this filter to store-LSHs with less memory than Redis keys with the risk of false positives. 419 420 Each *band* gets its own scalable Bloom filter on the Redis side: 421 422 ```text 423 BF.RESERVE <prefix>:<band_idx> <error> <capacity> EXPANSION <n> 424 ``` 425 426 """ 427 428 def __init__( 429 self, 430 *, 431 expected_docs: int, 432 host: str = "localhost", 433 port: int = 6379, 434 db: int = 0, 435 key_prefix: str = "bloomdedup", 436 error_rate: float = 1e-7, 437 expansion: int = 2, 438 num_bands: int | None = None, 439 **kwargs: Any, 440 ): 441 """ 442 Initialize the RedisBloom deduplicator. 443 Args: 444 expected_docs (int): Expected number of documents to deduplicate. This is used to set the initial capacity of the Bloom filter. 445 host (str): Redis server hostname. 446 port (int): Redis server port. 447 db (int): Redis database number. 448 key_prefix (str): Prefix for Redis keys to avoid collisions. You should use a unique prefix for each deduplication task. 449 error_rate (float): Desired error rate for the Bloom filter. 450 expansion (int): Expansion factor for the Bloom filter. This is used to increase the capacity of the filter dynamically. 451 num_bands (int | None): Number of bands to use for LSH to calculate the capacity of BloomFilter. If None, it will be set to 32. 452 **kwargs: Additional keyword arguments for parent Filter. 453 """ 454 if not is_loaded_dedup: 455 raise ImportError(IS_LOADED_DEDUP_ERROR_MSG) 456 super().__init__(**kwargs) 457 self.rds = redis.Redis(host=host, port=port, db=db) 458 self.key_prefix = key_prefix.encode() 459 460 _num_bands = num_bands if num_bands is not None else 32 461 462 try: 463 self.rds.execute_command( 464 "BF.RESERVE", 465 self.key_prefix, 466 error_rate, 467 expected_docs * _num_bands, 468 "EXPANSION", 469 expansion, 470 ) 471 except redis.ResponseError as e: 472 if "exists" not in str(e): 473 raise 474 475 def apply(self, document: Document) -> Document: 476 lsh_keys: list[str] | None = document.extras.get("dedup_lsh") 477 if lsh_keys is None: 478 raise ValueError( 479 "Document does not contain LSH keys for deduplication. Please apply GenerateDedupLSH first." 480 ) 481 482 key_bytes = [k.encode() for k in lsh_keys] 483 484 # Return value of BF.MADD is [1,0,1,...] (0 = already exists, 1 = insertion successful) 485 flags: Iterable[int] = self.rds.execute_command("BF.MADD", self.key_prefix, *key_bytes) 486 if 0 in flags: 487 document.is_rejected = True 488 return document 489 490 491class InlineDuplicateAnalyzer(Filter): 492 def __init__(self, **kwargs: Any): 493 super().__init__(**kwargs) 494 self.docs: dict[int, Document] = dict() # doc_id -> Document mapping 495 self.hash_pool: dict[str, int] = dict() # LSH key -> doc_id mapping 496 497 self._current_doc_id = 0 498 499 def apply(self, document: Document) -> Document: 500 """ 501 Analyze duplicates inline based on the LSH keys in the document. 502 This filter cannot use in the distributed environment because it uses a local hash pool. 503 """ 504 lsh_keys = document.extras.get("dedup_lsh") 505 if lsh_keys is None: 506 raise ValueError( 507 "Document does not contain LSH keys for deduplication. Please apply GenerateDedupLSH first." 508 ) 509 510 for lsh in lsh_keys: 511 if lsh in self.hash_pool: 512 document.is_rejected = True 513 document.extras["similar_doc"] = self.docs[self.hash_pool[lsh]].text 514 else: 515 self.hash_pool[lsh] = self._current_doc_id 516 self.docs[self._current_doc_id] = document 517 518 self._current_doc_id += 1 519 return document
83def char_level_splitter(text: str) -> list[str]: 84 """ 85 Split the text into characters. 86 This is a simple implementation that splits the text into individual characters. 87 """ 88 return list(text)
Split the text into characters. This is a simple implementation that splits the text into individual characters.
91def non_alpha_num_splitter(text: str) -> list[str]: 92 """ 93 Split the text into alphanumeric tokens. 94 This is a simple implementation that splits on non-alphanumeric characters. 95 """ 96 return [token for token in NON_ALPHA.split(text) if token]
Split the text into alphanumeric tokens. This is a simple implementation that splits on non-alphanumeric characters.
110def japanese_word_splitter(text: str) -> list[str]: 111 """ 112 Split the text into Japanese words using fugashi. 113 This will import fugashi and instantiate Tagger on first use. 114 """ 115 global _japanese_tagger 116 if _japanese_tagger is None: 117 fugashi = importlib.import_module("fugashi") 118 _japanese_tagger = fugashi.Tagger() 119 return [token.surface for token in _japanese_tagger(text)]
Split the text into Japanese words using fugashi. This will import fugashi and instantiate Tagger on first use.
122class GenerateDedupLSH(Filter): 123 """ 124 Filter that uses MinHash + Locality-Sensitive Hashing (LSH) to assign 125 deduplication keys to documents, allowing fast near-duplicate detection. 126 127 Attributes: 128 num_perm (int): Number of permutations (hash functions) for MinHash. 129 threshold (float): Similarity threshold for tuning LSH parameters. 130 tokenizer (Callable[[str], Iterable[str]]): Function to tokenize text. 131 n_grams (int): n-gram size for token grouping. 132 seed (int): Random seed for MinHash. 133 num_bands (int): Number of LSH bands, specified explicitly or computed automatically. 134 band_size (int): Number of hashes per band. 135 136 Notes 137 ----- 138 When band parameters are omitted, `_optimal_param` searches for the optimal 139 number of **bands** (`b`) and 140 **rows per band** (`r`) that minimise a weighted sum of false positives / 141 false negatives at the specified *threshold*. 142 143 HojiChar 0.18.0 introduces ``v2:`` keys using Rensa 0.5 and an optimized 144 hashing path. Hash generation was over 4x faster than HojiChar 0.17.3 145 in our 2,000-character English benchmark. Hashes differ from 146 HojiChar 0.17.x. Please rebuild your deduplication fingerprints, 147 or keep using HojiChar 0.17.x if you want to use existing LSH pool. 148 """ 149 150 _BYTES_PER_U32: Final[int] = 4 151 _LSH_KEY_VERSION: Final[str] = "v2" 152 153 def __init__( 154 self, 155 num_perm: int = 500, 156 threshold: float = 0.8, 157 tokenizer: Callable[[str], Iterable[str]] = char_level_splitter, 158 n_grams: int = 5, 159 seed: int = 42, 160 *, 161 num_bands: Optional[int] = None, 162 band_size: Optional[int] = None, 163 **kwargs: Any, 164 ) -> None: 165 """ 166 Initialize the deduplication filter with MinHash and LSH settings. 167 168 Args: 169 num_perm: Number of hash permutations for MinHash signature length. 170 Ignored when num_bands and band_size are both provided. 171 threshold: Similarity threshold to decide optimal LSH parameters. 172 Ignored when num_bands and band_size are both provided. 173 tokenizer: Function to split text into tokens. 174 n_grams: Number of tokens per n-gram for MinHash update. 175 seed: Seed for hash permutation consistency. 176 num_bands: Explicit number of LSH bands. Must be a positive integer 177 and provided together with band_size. 178 band_size: Explicit number of hashes per band. Must be a positive integer 179 and provided together with num_bands. When both are provided, 180 automatic selection is bypassed and the MinHash signature length 181 (self.num_perm) is set to num_bands * band_size. 182 **kwargs: Additional keyword arguments for parent Filter. 183 """ 184 super().__init__(**kwargs) 185 if not is_loaded_dedup: 186 raise ImportError(IS_LOADED_DEDUP_ERROR_MSG) 187 if n_grams <= 0: 188 raise ValueError("n_grams must be positive") 189 self.num_perm = num_perm 190 self.threshold = threshold 191 self.tokenizer = tokenizer 192 self.n_grams = n_grams 193 self.seed = seed 194 195 if num_bands is None and band_size is None: 196 self.num_bands, self.band_size = _optimal_param( 197 threshold=self.threshold, 198 num_perm=self.num_perm, 199 false_negative_weight=0.5, 200 false_positive_weight=0.5, 201 ) 202 else: 203 if num_bands is None or band_size is None: 204 raise ValueError("num_bands and band_size must be provided together") 205 if num_bands <= 0 or band_size <= 0: 206 raise ValueError("num_bands and band_size must be positive") 207 self.num_bands = num_bands 208 self.band_size = band_size 209 self.num_perm = num_bands * band_size 210 211 def _calculate_minhash_digest(self, text: str) -> list[int]: 212 """Compute the raw digest shared by the array API and the LSH fast path.""" 213 if self.tokenizer is char_level_splitter: 214 n = self.n_grams 215 tokens = [text[i : i + n] for i in range(len(text) - n + 1)] 216 else: 217 tokens = [" ".join(grams) for grams in _ngrams(self.tokenizer(text), self.n_grams)] 218 minhash = RMinHash(num_perm=self.num_perm, seed=self.seed) 219 minhash.update(tokens) 220 return cast(list[int], minhash.digest()) 221 222 def calculate_minhash_signature(self, text: str) -> NDArray[np.uint32]: 223 """ 224 Compute MinHash signature of input text as an array of uint32. 225 226 Steps: 227 1. Tokenize text using the provided tokenizer. 228 2. Generate n-gram tokens. 229 3. Update MinHash with n-gram tokens. 230 231 Args: 232 text: Input document text to be hashed. 233 234 Returns: 235 A 1D numpy array of shape (num_perm,) with dtype uint32. 236 """ 237 return np.asarray(self._calculate_minhash_digest(text), dtype=np.uint32) 238 239 def _sig_bytes_le(self, sig: NDArray[np.uint32]) -> memoryview: 240 """ 241 Return the signature as *little-endian* byte view. 242 Platform‑independent way to get bytes from a numpy array. 243 """ 244 if sys.byteorder == "little": 245 # amd64 / arm64 (little) 246 return memoryview(cast(Any, sig)).cast("B") # zero-copy 247 else: 248 # big-endian CPU 249 return memoryview(cast(Any, sig.byteswap())).cast("B") # 1 copy 250 251 def signature_to_lsh_digest( 252 self, signature: NDArray[np.uint32], band_size: int, band_idx: int 253 ) -> int: 254 """ 255 Convert a slice of the MinHash signature into an LSH digest with less memory overhead. 256 257 This method is optimized for speed by avoiding copies: 258 - We view the uint32 array as raw bytes (uint8 view). 259 - We create a memoryview of the byte slice for the specified band. 260 - We compute a 128-bit hash directly on the slice. 261 262 Args: 263 signature: 1D numpy array of uint32 representing MinHash signature. 264 band_size: Number of hashes per LSH band. 265 band_idx: Index of the band to hash (0-based). 266 267 Returns: 268 An integer representing the 128-bit hash digest of the band. 269 270 Raises: 271 AssertionError: If signature shape/dtype or band index is invalid. 272 """ 273 assert signature.dtype == np.uint32 and signature.ndim == 1, ( 274 "signature must be a 1D numpy array of uint32" 275 ) 276 assert 0 <= band_idx < self.num_bands, ( 277 f"band_idx {band_idx} out of range [0, {self.num_bands})" 278 ) 279 assert len(signature) >= band_size * self.num_bands, ( 280 "signature length is too short for given band_size and num_bands" 281 ) 282 283 # Compute byte offsets for the selected band 284 start = band_idx * band_size * self._BYTES_PER_U32 285 stop = start + band_size * self._BYTES_PER_U32 286 287 # View signature as raw bytes without copy. memoryview avoids creating new bytes. 288 mv = self._sig_bytes_le(signature)[start:stop] # slice view 289 290 return xxhash.xxh128_intdigest(mv) 291 292 def _format_lsh_key(self, band_idx: int, digest: int) -> str: 293 """ 294 Format the LSH key with the HojiChar scheme version, band index, and digest. 295 """ 296 return f"{self._LSH_KEY_VERSION}:{band_idx}+{digest:032x}" 297 298 def apply(self, document: Document) -> Document: 299 """ 300 Decorate the document with LSH deduplication keys. 301 302 For each band, compute the digest and format as a hex string: 303 'v2:<band_idx>+<128-bit-digest-hex>'. 304 Keys are stored in document.extras['dedup_lsh']. 305 306 Args: 307 document: Document object with 'text' attribute. 308 309 Returns: 310 The same Document object with 'dedup_lsh' added in extras. 311 """ 312 digest = self._calculate_minhash_digest(document.text) 313 # Pack once in little-endian order without an intermediate NumPy array. 314 signature_bytes = memoryview(struct.pack(f"<{self.num_perm}I", *digest)) 315 band_bytes = self.band_size * self._BYTES_PER_U32 316 lsh_keys = [ 317 f"{self._LSH_KEY_VERSION}:{band_idx}+" 318 + xxhash.xxh128_hexdigest( 319 signature_bytes[band_idx * band_bytes : (band_idx + 1) * band_bytes] 320 ) 321 for band_idx in range(self.num_bands) 322 ] 323 324 document.extras["dedup_lsh"] = lsh_keys 325 return document
Filter that uses MinHash + Locality-Sensitive Hashing (LSH) to assign deduplication keys to documents, allowing fast near-duplicate detection.
Attributes: num_perm (int): Number of permutations (hash functions) for MinHash. threshold (float): Similarity threshold for tuning LSH parameters. tokenizer (Callable[[str], Iterable[str]]): Function to tokenize text. n_grams (int): n-gram size for token grouping. seed (int): Random seed for MinHash. num_bands (int): Number of LSH bands, specified explicitly or computed automatically. band_size (int): Number of hashes per band.
Notes
When band parameters are omitted, _optimal_param searches for the optimal
number of bands (b) and
rows per band (r) that minimise a weighted sum of false positives /
false negatives at the specified threshold.
HojiChar 0.18.0 introduces v2: keys using Rensa 0.5 and an optimized
hashing path. Hash generation was over 4x faster than HojiChar 0.17.3
in our 2,000-character English benchmark. Hashes differ from
HojiChar 0.17.x. Please rebuild your deduplication fingerprints,
or keep using HojiChar 0.17.x if you want to use existing LSH pool.
153 def __init__( 154 self, 155 num_perm: int = 500, 156 threshold: float = 0.8, 157 tokenizer: Callable[[str], Iterable[str]] = char_level_splitter, 158 n_grams: int = 5, 159 seed: int = 42, 160 *, 161 num_bands: Optional[int] = None, 162 band_size: Optional[int] = None, 163 **kwargs: Any, 164 ) -> None: 165 """ 166 Initialize the deduplication filter with MinHash and LSH settings. 167 168 Args: 169 num_perm: Number of hash permutations for MinHash signature length. 170 Ignored when num_bands and band_size are both provided. 171 threshold: Similarity threshold to decide optimal LSH parameters. 172 Ignored when num_bands and band_size are both provided. 173 tokenizer: Function to split text into tokens. 174 n_grams: Number of tokens per n-gram for MinHash update. 175 seed: Seed for hash permutation consistency. 176 num_bands: Explicit number of LSH bands. Must be a positive integer 177 and provided together with band_size. 178 band_size: Explicit number of hashes per band. Must be a positive integer 179 and provided together with num_bands. When both are provided, 180 automatic selection is bypassed and the MinHash signature length 181 (self.num_perm) is set to num_bands * band_size. 182 **kwargs: Additional keyword arguments for parent Filter. 183 """ 184 super().__init__(**kwargs) 185 if not is_loaded_dedup: 186 raise ImportError(IS_LOADED_DEDUP_ERROR_MSG) 187 if n_grams <= 0: 188 raise ValueError("n_grams must be positive") 189 self.num_perm = num_perm 190 self.threshold = threshold 191 self.tokenizer = tokenizer 192 self.n_grams = n_grams 193 self.seed = seed 194 195 if num_bands is None and band_size is None: 196 self.num_bands, self.band_size = _optimal_param( 197 threshold=self.threshold, 198 num_perm=self.num_perm, 199 false_negative_weight=0.5, 200 false_positive_weight=0.5, 201 ) 202 else: 203 if num_bands is None or band_size is None: 204 raise ValueError("num_bands and band_size must be provided together") 205 if num_bands <= 0 or band_size <= 0: 206 raise ValueError("num_bands and band_size must be positive") 207 self.num_bands = num_bands 208 self.band_size = band_size 209 self.num_perm = num_bands * band_size
Initialize the deduplication filter with MinHash and LSH settings.
Args: num_perm: Number of hash permutations for MinHash signature length. Ignored when num_bands and band_size are both provided. threshold: Similarity threshold to decide optimal LSH parameters. Ignored when num_bands and band_size are both provided. tokenizer: Function to split text into tokens. n_grams: Number of tokens per n-gram for MinHash update. seed: Seed for hash permutation consistency. num_bands: Explicit number of LSH bands. Must be a positive integer and provided together with band_size. band_size: Explicit number of hashes per band. Must be a positive integer and provided together with num_bands. When both are provided, automatic selection is bypassed and the MinHash signature length (self.num_perm) is set to num_bands * band_size. **kwargs: Additional keyword arguments for parent Filter.
222 def calculate_minhash_signature(self, text: str) -> NDArray[np.uint32]: 223 """ 224 Compute MinHash signature of input text as an array of uint32. 225 226 Steps: 227 1. Tokenize text using the provided tokenizer. 228 2. Generate n-gram tokens. 229 3. Update MinHash with n-gram tokens. 230 231 Args: 232 text: Input document text to be hashed. 233 234 Returns: 235 A 1D numpy array of shape (num_perm,) with dtype uint32. 236 """ 237 return np.asarray(self._calculate_minhash_digest(text), dtype=np.uint32)
Compute MinHash signature of input text as an array of uint32.
Steps: 1. Tokenize text using the provided tokenizer. 2. Generate n-gram tokens. 3. Update MinHash with n-gram tokens.
Args: text: Input document text to be hashed.
Returns: A 1D numpy array of shape (num_perm,) with dtype uint32.
251 def signature_to_lsh_digest( 252 self, signature: NDArray[np.uint32], band_size: int, band_idx: int 253 ) -> int: 254 """ 255 Convert a slice of the MinHash signature into an LSH digest with less memory overhead. 256 257 This method is optimized for speed by avoiding copies: 258 - We view the uint32 array as raw bytes (uint8 view). 259 - We create a memoryview of the byte slice for the specified band. 260 - We compute a 128-bit hash directly on the slice. 261 262 Args: 263 signature: 1D numpy array of uint32 representing MinHash signature. 264 band_size: Number of hashes per LSH band. 265 band_idx: Index of the band to hash (0-based). 266 267 Returns: 268 An integer representing the 128-bit hash digest of the band. 269 270 Raises: 271 AssertionError: If signature shape/dtype or band index is invalid. 272 """ 273 assert signature.dtype == np.uint32 and signature.ndim == 1, ( 274 "signature must be a 1D numpy array of uint32" 275 ) 276 assert 0 <= band_idx < self.num_bands, ( 277 f"band_idx {band_idx} out of range [0, {self.num_bands})" 278 ) 279 assert len(signature) >= band_size * self.num_bands, ( 280 "signature length is too short for given band_size and num_bands" 281 ) 282 283 # Compute byte offsets for the selected band 284 start = band_idx * band_size * self._BYTES_PER_U32 285 stop = start + band_size * self._BYTES_PER_U32 286 287 # View signature as raw bytes without copy. memoryview avoids creating new bytes. 288 mv = self._sig_bytes_le(signature)[start:stop] # slice view 289 290 return xxhash.xxh128_intdigest(mv)
Convert a slice of the MinHash signature into an LSH digest with less memory overhead.
This method is optimized for speed by avoiding copies:
- We view the uint32 array as raw bytes (uint8 view).
- We create a memoryview of the byte slice for the specified band.
- We compute a 128-bit hash directly on the slice.
Args: signature: 1D numpy array of uint32 representing MinHash signature. band_size: Number of hashes per LSH band. band_idx: Index of the band to hash (0-based).
Returns: An integer representing the 128-bit hash digest of the band.
Raises: AssertionError: If signature shape/dtype or band index is invalid.
298 def apply(self, document: Document) -> Document: 299 """ 300 Decorate the document with LSH deduplication keys. 301 302 For each band, compute the digest and format as a hex string: 303 'v2:<band_idx>+<128-bit-digest-hex>'. 304 Keys are stored in document.extras['dedup_lsh']. 305 306 Args: 307 document: Document object with 'text' attribute. 308 309 Returns: 310 The same Document object with 'dedup_lsh' added in extras. 311 """ 312 digest = self._calculate_minhash_digest(document.text) 313 # Pack once in little-endian order without an intermediate NumPy array. 314 signature_bytes = memoryview(struct.pack(f"<{self.num_perm}I", *digest)) 315 band_bytes = self.band_size * self._BYTES_PER_U32 316 lsh_keys = [ 317 f"{self._LSH_KEY_VERSION}:{band_idx}+" 318 + xxhash.xxh128_hexdigest( 319 signature_bytes[band_idx * band_bytes : (band_idx + 1) * band_bytes] 320 ) 321 for band_idx in range(self.num_bands) 322 ] 323 324 document.extras["dedup_lsh"] = lsh_keys 325 return document
Decorate the document with LSH deduplication keys.
For each band, compute the digest and format as a hex string:
'v2:
Args: document: Document object with 'text' attribute.
Returns: The same Document object with 'dedup_lsh' added in extras.
328class InlineDeduplicator(Filter): 329 """ 330 Simple in‑memory deduplicator. 331 332 Stores every LSH key in a local :pyclass:`set`. If any key of the incoming 333 document is already present, the document is marked as duplicate via 334 `document.is_rejected = True`. 335 336 **Limitations** 337 ------------- 338 *State is per‑process only.* Running multiple workers or machines will *not* 339 share the key set – use :class:`RedisDeduplicator` or 340 :class:`RedisBloomDeduplicator` for distributed setups. 341 """ 342 343 def __init__(self, **kwargs: Any): 344 super().__init__(**kwargs) 345 self.hash_pool: set[str] = set() 346 347 def apply(self, document: Document) -> Document: 348 """ 349 Inline deduplication based on the LSH keys in the document. 350 This filter cannot use in the distributed environment because it uses a local hash pool. 351 """ 352 lsh_keys = document.extras.get("dedup_lsh") 353 if lsh_keys is None: 354 raise ValueError( 355 "Document does not contain LSH keys for deduplication. Please apply GenerateDedupLSH first." 356 ) 357 358 for lsh in lsh_keys: 359 if lsh in self.hash_pool: 360 document.is_rejected = True 361 else: 362 self.hash_pool.add(lsh) 363 return document
Simple in‑memory deduplicator.
Stores every LSH key in a local :pyclass:set. If any key of the incoming
document is already present, the document is marked as duplicate via
document.is_rejected = True.
Limitations
State is per‑process only. Running multiple workers or machines will not
share the key set – use RedisDeduplicator or
RedisBloomDeduplicator for distributed setups.
343 def __init__(self, **kwargs: Any): 344 super().__init__(**kwargs) 345 self.hash_pool: set[str] = set()
Initialize the filter.
Parameters
p : float
The probability of applying the filter.
If p is 1, the filter will always be applied.
skip_rejected : bool
If True, the filter will skip documents that are already rejected.
If you want to apply the filter to all documents (e.g., postprocess), set this to False.
random_state : Optional[Union[int, np.random.Generator]]
Seed for the random number generator.
If None, a new random number generator will be created.
If None, and use in the Compose class, the random state is shared with the Compose object.
use_batch : bool
If True, the filter will process documents in batches in the apply_stream method.
batch_size : int
The size of the batch to process documents in the apply_stream method.
kwargs : Any
Additional keyword arguments to pass to the filter.
347 def apply(self, document: Document) -> Document: 348 """ 349 Inline deduplication based on the LSH keys in the document. 350 This filter cannot use in the distributed environment because it uses a local hash pool. 351 """ 352 lsh_keys = document.extras.get("dedup_lsh") 353 if lsh_keys is None: 354 raise ValueError( 355 "Document does not contain LSH keys for deduplication. Please apply GenerateDedupLSH first." 356 ) 357 358 for lsh in lsh_keys: 359 if lsh in self.hash_pool: 360 document.is_rejected = True 361 else: 362 self.hash_pool.add(lsh) 363 return document
Inline deduplication based on the LSH keys in the document. This filter cannot use in the distributed environment because it uses a local hash pool.
366class RedisDeduplicator(Filter): 367 """ 368 Distributed deduplicator using **plain Redis keys**. 369 You have to run a Redis server and pass its connection parameters. 370 """ 371 372 def __init__( 373 self, 374 *, 375 host: str = "localhost", 376 port: int = 6379, 377 db: int = 0, 378 key_prefix: str = "dedup", 379 **kwargs: Any, 380 ) -> None: 381 """ 382 Initialize the Redis deduplicator. 383 Args: 384 host (str): Redis server hostname. 385 port (int): Redis server port. 386 db (int): Redis database number. 387 key_prefix (str): Prefix for Redis keys to avoid collisions. You should use a unique prefix for each deduplication task. 388 **kwargs: Additional keyword arguments for parent Filter. 389 """ 390 if not is_loaded_dedup: 391 raise ImportError(IS_LOADED_DEDUP_ERROR_MSG) 392 super().__init__(**kwargs) 393 self.rds = redis.Redis(host=host, port=port, db=db, decode_responses=False) 394 self.key_prefix = key_prefix.encode() 395 396 try: 397 self.rds.ping() 398 except redis.exceptions.RedisError as exc: 399 raise RuntimeError(f"Cannot connect to Redis server {host}:{port}/{db}") from exc 400 401 def apply(self, document: Document) -> Document: 402 lsh_keys = document.extras.get("dedup_lsh") 403 if lsh_keys is None: 404 raise ValueError("Apply GenerateDedupLSH first") 405 406 pipe = self.rds.pipeline(transaction=False) 407 for k in lsh_keys: 408 pipe.set(self.key_prefix + b":" + k.encode(), b"1", nx=True) 409 results: list[bool | None] = pipe.execute() # If instance already exists, it returns None 410 411 if any(r is None for r in results): 412 document.is_rejected = True 413 return document
Distributed deduplicator using plain Redis keys. You have to run a Redis server and pass its connection parameters.
372 def __init__( 373 self, 374 *, 375 host: str = "localhost", 376 port: int = 6379, 377 db: int = 0, 378 key_prefix: str = "dedup", 379 **kwargs: Any, 380 ) -> None: 381 """ 382 Initialize the Redis deduplicator. 383 Args: 384 host (str): Redis server hostname. 385 port (int): Redis server port. 386 db (int): Redis database number. 387 key_prefix (str): Prefix for Redis keys to avoid collisions. You should use a unique prefix for each deduplication task. 388 **kwargs: Additional keyword arguments for parent Filter. 389 """ 390 if not is_loaded_dedup: 391 raise ImportError(IS_LOADED_DEDUP_ERROR_MSG) 392 super().__init__(**kwargs) 393 self.rds = redis.Redis(host=host, port=port, db=db, decode_responses=False) 394 self.key_prefix = key_prefix.encode() 395 396 try: 397 self.rds.ping() 398 except redis.exceptions.RedisError as exc: 399 raise RuntimeError(f"Cannot connect to Redis server {host}:{port}/{db}") from exc
Initialize the Redis deduplicator. Args: host (str): Redis server hostname. port (int): Redis server port. db (int): Redis database number. key_prefix (str): Prefix for Redis keys to avoid collisions. You should use a unique prefix for each deduplication task. **kwargs: Additional keyword arguments for parent Filter.
401 def apply(self, document: Document) -> Document: 402 lsh_keys = document.extras.get("dedup_lsh") 403 if lsh_keys is None: 404 raise ValueError("Apply GenerateDedupLSH first") 405 406 pipe = self.rds.pipeline(transaction=False) 407 for k in lsh_keys: 408 pipe.set(self.key_prefix + b":" + k.encode(), b"1", nx=True) 409 results: list[bool | None] = pipe.execute() # If instance already exists, it returns None 410 411 if any(r is None for r in results): 412 document.is_rejected = True 413 return document
Definition of filter behavior.
The document must have a protocol TextContent,
and mostly used hojichar.Document class.
In this method, the filter will modify document.text or
document.extras and set document.is_rejected = True to discard the document.
Parameters
document : Document Input document
Returns
Document Processed Document
416class RedisBloomDeduplicator(Filter): 417 """ 418 Distributed deduplicator backed by **RedisBloom scalable Bloom filters**. 419 You can use this filter to store-LSHs with less memory than Redis keys with the risk of false positives. 420 421 Each *band* gets its own scalable Bloom filter on the Redis side: 422 423 ```text 424 BF.RESERVE <prefix>:<band_idx> <error> <capacity> EXPANSION <n> 425 ``` 426 427 """ 428 429 def __init__( 430 self, 431 *, 432 expected_docs: int, 433 host: str = "localhost", 434 port: int = 6379, 435 db: int = 0, 436 key_prefix: str = "bloomdedup", 437 error_rate: float = 1e-7, 438 expansion: int = 2, 439 num_bands: int | None = None, 440 **kwargs: Any, 441 ): 442 """ 443 Initialize the RedisBloom deduplicator. 444 Args: 445 expected_docs (int): Expected number of documents to deduplicate. This is used to set the initial capacity of the Bloom filter. 446 host (str): Redis server hostname. 447 port (int): Redis server port. 448 db (int): Redis database number. 449 key_prefix (str): Prefix for Redis keys to avoid collisions. You should use a unique prefix for each deduplication task. 450 error_rate (float): Desired error rate for the Bloom filter. 451 expansion (int): Expansion factor for the Bloom filter. This is used to increase the capacity of the filter dynamically. 452 num_bands (int | None): Number of bands to use for LSH to calculate the capacity of BloomFilter. If None, it will be set to 32. 453 **kwargs: Additional keyword arguments for parent Filter. 454 """ 455 if not is_loaded_dedup: 456 raise ImportError(IS_LOADED_DEDUP_ERROR_MSG) 457 super().__init__(**kwargs) 458 self.rds = redis.Redis(host=host, port=port, db=db) 459 self.key_prefix = key_prefix.encode() 460 461 _num_bands = num_bands if num_bands is not None else 32 462 463 try: 464 self.rds.execute_command( 465 "BF.RESERVE", 466 self.key_prefix, 467 error_rate, 468 expected_docs * _num_bands, 469 "EXPANSION", 470 expansion, 471 ) 472 except redis.ResponseError as e: 473 if "exists" not in str(e): 474 raise 475 476 def apply(self, document: Document) -> Document: 477 lsh_keys: list[str] | None = document.extras.get("dedup_lsh") 478 if lsh_keys is None: 479 raise ValueError( 480 "Document does not contain LSH keys for deduplication. Please apply GenerateDedupLSH first." 481 ) 482 483 key_bytes = [k.encode() for k in lsh_keys] 484 485 # Return value of BF.MADD is [1,0,1,...] (0 = already exists, 1 = insertion successful) 486 flags: Iterable[int] = self.rds.execute_command("BF.MADD", self.key_prefix, *key_bytes) 487 if 0 in flags: 488 document.is_rejected = True 489 return document
Distributed deduplicator backed by RedisBloom scalable Bloom filters. You can use this filter to store-LSHs with less memory than Redis keys with the risk of false positives.
Each band gets its own scalable Bloom filter on the Redis side:
BF.RESERVE <prefix>:<band_idx> <error> <capacity> EXPANSION <n>
429 def __init__( 430 self, 431 *, 432 expected_docs: int, 433 host: str = "localhost", 434 port: int = 6379, 435 db: int = 0, 436 key_prefix: str = "bloomdedup", 437 error_rate: float = 1e-7, 438 expansion: int = 2, 439 num_bands: int | None = None, 440 **kwargs: Any, 441 ): 442 """ 443 Initialize the RedisBloom deduplicator. 444 Args: 445 expected_docs (int): Expected number of documents to deduplicate. This is used to set the initial capacity of the Bloom filter. 446 host (str): Redis server hostname. 447 port (int): Redis server port. 448 db (int): Redis database number. 449 key_prefix (str): Prefix for Redis keys to avoid collisions. You should use a unique prefix for each deduplication task. 450 error_rate (float): Desired error rate for the Bloom filter. 451 expansion (int): Expansion factor for the Bloom filter. This is used to increase the capacity of the filter dynamically. 452 num_bands (int | None): Number of bands to use for LSH to calculate the capacity of BloomFilter. If None, it will be set to 32. 453 **kwargs: Additional keyword arguments for parent Filter. 454 """ 455 if not is_loaded_dedup: 456 raise ImportError(IS_LOADED_DEDUP_ERROR_MSG) 457 super().__init__(**kwargs) 458 self.rds = redis.Redis(host=host, port=port, db=db) 459 self.key_prefix = key_prefix.encode() 460 461 _num_bands = num_bands if num_bands is not None else 32 462 463 try: 464 self.rds.execute_command( 465 "BF.RESERVE", 466 self.key_prefix, 467 error_rate, 468 expected_docs * _num_bands, 469 "EXPANSION", 470 expansion, 471 ) 472 except redis.ResponseError as e: 473 if "exists" not in str(e): 474 raise
Initialize the RedisBloom deduplicator. Args: expected_docs (int): Expected number of documents to deduplicate. This is used to set the initial capacity of the Bloom filter. host (str): Redis server hostname. port (int): Redis server port. db (int): Redis database number. key_prefix (str): Prefix for Redis keys to avoid collisions. You should use a unique prefix for each deduplication task. error_rate (float): Desired error rate for the Bloom filter. expansion (int): Expansion factor for the Bloom filter. This is used to increase the capacity of the filter dynamically. num_bands (int | None): Number of bands to use for LSH to calculate the capacity of BloomFilter. If None, it will be set to 32. **kwargs: Additional keyword arguments for parent Filter.
476 def apply(self, document: Document) -> Document: 477 lsh_keys: list[str] | None = document.extras.get("dedup_lsh") 478 if lsh_keys is None: 479 raise ValueError( 480 "Document does not contain LSH keys for deduplication. Please apply GenerateDedupLSH first." 481 ) 482 483 key_bytes = [k.encode() for k in lsh_keys] 484 485 # Return value of BF.MADD is [1,0,1,...] (0 = already exists, 1 = insertion successful) 486 flags: Iterable[int] = self.rds.execute_command("BF.MADD", self.key_prefix, *key_bytes) 487 if 0 in flags: 488 document.is_rejected = True 489 return document
Definition of filter behavior.
The document must have a protocol TextContent,
and mostly used hojichar.Document class.
In this method, the filter will modify document.text or
document.extras and set document.is_rejected = True to discard the document.
Parameters
document : Document Input document
Returns
Document Processed Document
492class InlineDuplicateAnalyzer(Filter): 493 def __init__(self, **kwargs: Any): 494 super().__init__(**kwargs) 495 self.docs: dict[int, Document] = dict() # doc_id -> Document mapping 496 self.hash_pool: dict[str, int] = dict() # LSH key -> doc_id mapping 497 498 self._current_doc_id = 0 499 500 def apply(self, document: Document) -> Document: 501 """ 502 Analyze duplicates inline based on the LSH keys in the document. 503 This filter cannot use in the distributed environment because it uses a local hash pool. 504 """ 505 lsh_keys = document.extras.get("dedup_lsh") 506 if lsh_keys is None: 507 raise ValueError( 508 "Document does not contain LSH keys for deduplication. Please apply GenerateDedupLSH first." 509 ) 510 511 for lsh in lsh_keys: 512 if lsh in self.hash_pool: 513 document.is_rejected = True 514 document.extras["similar_doc"] = self.docs[self.hash_pool[lsh]].text 515 else: 516 self.hash_pool[lsh] = self._current_doc_id 517 self.docs[self._current_doc_id] = document 518 519 self._current_doc_id += 1 520 return document
Base class for all filters. Document-level filters must inherit from this class.
The definition of text processing is in apply method.
If you define a new filter, override the method.
When this class is called, apply the filter from string to string.
With context manager, you can use the filter as follows:
with YourFilter(p=0.5) as filt:
text = filt("This is a sample text.")
493 def __init__(self, **kwargs: Any): 494 super().__init__(**kwargs) 495 self.docs: dict[int, Document] = dict() # doc_id -> Document mapping 496 self.hash_pool: dict[str, int] = dict() # LSH key -> doc_id mapping 497 498 self._current_doc_id = 0
Initialize the filter.
Parameters
p : float
The probability of applying the filter.
If p is 1, the filter will always be applied.
skip_rejected : bool
If True, the filter will skip documents that are already rejected.
If you want to apply the filter to all documents (e.g., postprocess), set this to False.
random_state : Optional[Union[int, np.random.Generator]]
Seed for the random number generator.
If None, a new random number generator will be created.
If None, and use in the Compose class, the random state is shared with the Compose object.
use_batch : bool
If True, the filter will process documents in batches in the apply_stream method.
batch_size : int
The size of the batch to process documents in the apply_stream method.
kwargs : Any
Additional keyword arguments to pass to the filter.
500 def apply(self, document: Document) -> Document: 501 """ 502 Analyze duplicates inline based on the LSH keys in the document. 503 This filter cannot use in the distributed environment because it uses a local hash pool. 504 """ 505 lsh_keys = document.extras.get("dedup_lsh") 506 if lsh_keys is None: 507 raise ValueError( 508 "Document does not contain LSH keys for deduplication. Please apply GenerateDedupLSH first." 509 ) 510 511 for lsh in lsh_keys: 512 if lsh in self.hash_pool: 513 document.is_rejected = True 514 document.extras["similar_doc"] = self.docs[self.hash_pool[lsh]].text 515 else: 516 self.hash_pool[lsh] = self._current_doc_id 517 self.docs[self._current_doc_id] = document 518 519 self._current_doc_id += 1 520 return document
Analyze duplicates inline based on the LSH keys in the document. This filter cannot use in the distributed environment because it uses a local hash pool.