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.

  1. Tokenize
    • Split text into tokens. (Default: character‑level, but you can plug in any callable.)
  2. n-grams
    • Group tokens into n‑grams (n_grams=5) to capture context.
  3. MinHash
    • 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.
    • This module uses rensa.RMinHash which is a fast MinHash implementation by Rust language.
  4. Banding and compression
    • Split the signature into b bands, each containing r integers (num_perm ≈ b×r).
    • Treat the r integers as raw bytes, hash them with xxhash‑128, and format as v2:+.
  5. Output
    • Store all band keys in document.extras['dedup_lsh'] as a list of strings.

InlineDeduplicator, RedisDeduplicator, and RedisBloomDeduplicator filters use the generated LSH keys to mark documents as duplicates.

  • InlineDeduplicator stores LSH keys in a local set, so it works only in a single process.
  • RedisDeduplicator stores LSH keys in Redis, so it works in a distributed environment.
  • RedisBloomDeduplicator stores 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
def char_level_splitter(text: str) -> list[str]:
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.

def non_alpha_num_splitter(text: str) -> list[str]:
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.

def japanese_word_splitter(text: str) -> list[str]:
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.

class GenerateDedupLSH(hojichar.core.filter_interface.Filter):
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.

GenerateDedupLSH( num_perm: int = 500, threshold: float = 0.8, tokenizer: Callable[[str], Iterable[str]] = <function char_level_splitter>, n_grams: int = 5, seed: int = 42, *, num_bands: Optional[int] = None, band_size: Optional[int] = None, **kwargs: Any)
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.

def calculate_minhash_signature( self, text: str) -> numpy.ndarray[tuple[typing.Any, ...], numpy.dtype[numpy.uint32]]:
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.

def signature_to_lsh_digest( self, signature: numpy.ndarray[tuple[typing.Any, ...], numpy.dtype[numpy.uint32]], band_size: int, band_idx: int) -> int:
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.

def apply( self, document: hojichar.core.models.Document) -> hojichar.core.models.Document:
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:+<128-bit-digest-hex>'. Keys are stored in document.extras['dedup_lsh'].

Args: document: Document object with 'text' attribute.

Returns: The same Document object with 'dedup_lsh' added in extras.

class InlineDeduplicator(hojichar.core.filter_interface.Filter):
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.

InlineDeduplicator(**kwargs: Any)
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.

def apply( self, document: hojichar.core.models.Document) -> hojichar.core.models.Document:
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.

class RedisDeduplicator(hojichar.core.filter_interface.Filter):
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.

RedisDeduplicator( *, host: str = 'localhost', port: int = 6379, db: int = 0, key_prefix: str = 'dedup', **kwargs: Any)
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.

def apply( self, document: hojichar.core.models.Document) -> hojichar.core.models.Document:
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

class RedisBloomDeduplicator(hojichar.core.filter_interface.Filter):
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>
RedisBloomDeduplicator( *, expected_docs: int, host: str = 'localhost', port: int = 6379, db: int = 0, key_prefix: str = 'bloomdedup', error_rate: float = 1e-07, expansion: int = 2, num_bands: int | None = None, **kwargs: Any)
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.

def apply( self, document: hojichar.core.models.Document) -> hojichar.core.models.Document:
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

class InlineDuplicateAnalyzer(hojichar.core.filter_interface.Filter):
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.")
InlineDuplicateAnalyzer(**kwargs: Any)
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.

def apply( self, document: hojichar.core.models.Document) -> hojichar.core.models.Document:
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.