Curate TextProcess DataDeduplication

Semantic Deduplication

View as Markdown

Detect and remove semantically redundant data from your large text datasets using NeMo Curator.

Unlike exact or fuzzy deduplication, which focus on textual similarity, semantic deduplication leverages the meaning of content to identify duplicates. This approach can significantly reduce dataset size while maintaining or even improving model performance.

The technique uses embeddings to identify “semantic duplicates” - content pairs that convey similar meaning despite using different words.

GPU Acceleration: Semantic deduplication requires GPU acceleration for both embedding generation and clustering operations. This method uses cuDF for GPU-accelerated dataframe operations and PyTorch models on GPU for optimal performance.

How It Works

Semantic deduplication identifies meaning-based duplicates using embeddings:

  1. Generates embeddings for each document using transformer models
  2. Clusters embeddings using K-means
  3. Computes pairwise cosine similarities within clusters
  4. Identifies semantic duplicates based on similarity threshold
  5. Removes duplicates, keeping one representative per group

Before You Start

Prerequisites:

  • GPU acceleration (required for embedding generation and clustering)
  • Stable document identifiers for removal (either existing IDs or IDs managed by the workflow and removal stages)
  • Access to the gated default model: accept the EmbeddingGemma usage license, then authenticate with a Hugging Face read token through HF_TOKEN or the hf_token parameter

Running in Docker: When running semantic deduplication inside the NeMo Curator container, ensure the container is started with --gpus all so that CUDA GPUs are available. Without this flag, you will see RuntimeError: No CUDA GPUs are available. Also activate the virtual environment with source /opt/venv/env.sh after entering the container.

Quick Start

Get started with semantic deduplication using the following example of identifying duplicates, then remove them in one step:

from nemo_curator.stages.text.deduplication.semantic import TextSemanticDeduplicationWorkflow
# Default: uses vLLM with google/embeddinggemma-300m
workflow = TextSemanticDeduplicationWorkflow(
input_path="input_data/",
output_path="./results",
cache_path="./sem_cache",
n_clusters=100,
kmeans_fit_data_fraction=0.1, # Fit centroids on about 10% of generated embedding files
eps=0.07, # Similarity threshold
id_field="doc_id",
perform_removal=True, # Complete deduplication
)
results = workflow.run()
# Clean dataset saved to ./results/deduplicated/

Configuration

Configure semantic deduplication using these key parameters:

For fine-grained control, break semantic deduplication into separate stages:

from nemo_curator.stages.deduplication.id_generator import create_id_generator_actor
from nemo_curator.stages.text.embedders.vllm import VLLMEmbeddingModelStage
from nemo_curator.stages.deduplication.semantic import SemanticDeduplicationWorkflow
from nemo_curator.stages.text.deduplication.removal_workflow import TextDuplicatesRemovalWorkflow
from nemo_curator.pipeline import Pipeline
from nemo_curator.stages.text.io.reader import ParquetReader
from nemo_curator.stages.text.io.writer import ParquetWriter
# Step 1: Create ID generator
create_id_generator_actor()
# Step 2: Generate embeddings separately (using vLLM)
embedding_pipeline = Pipeline(
name="embedding_pipeline",
stages=[
ParquetReader(file_paths=input_path, files_per_partition=1, fields=["text"], _generate_ids=True),
# VLLMEmbeddingModelStage uses shorter parameter names than the workflow wrapper:
# pretokenize (not embedding_pretokenize), vllm_init_kwargs (not embedding_vllm_init_kwargs),
# max_chars (not embedding_max_chars), cache_dir (not model_cache_dir)
VLLMEmbeddingModelStage(
model_identifier="google/embeddinggemma-300m",
text_field="text",
),
ParquetWriter(path=embedding_output_path, fields=["_curator_dedup_id", "embeddings"]),
],
)
embedding_out = embedding_pipeline.run()
# Step 3: Run clustering and pairwise similarity (without duplicate identification)
semantic_workflow = SemanticDeduplicationWorkflow(
input_path=embedding_output_path,
output_path=semantic_workflow_path,
n_clusters=100,
fit_data_fraction=0.1,
id_field="_curator_dedup_id",
embedding_field="embeddings",
eps=None, # Skip duplicate identification for analysis
)
result = semantic_workflow.run()
# result.metadata contains: total_time, num_duplicates, kmeans_time, pairwise_time
# Step 4: Analyze similarity distribution to choose eps
# Step 5: Identify duplicates with chosen eps
# Step 6: Remove duplicates from original dataset

This approach enables analysis of intermediate results and parameter tuning.

Disable Weighted K-Means Task Assignment

By default, the Ray Actor Pool executor uses each input partition’s file size to balance K-means work across RAFT actors. To restore equal task-count assignment, override the decomposed KMeansStage with USE_TASK_WEIGHTS=False:

from nemo_curator.backends.ray_actor_pool import RayActorPoolExecutor
from nemo_curator.backends.utils import RayStageSpecKeys
from nemo_curator.pipeline import Pipeline
from nemo_curator.stages.deduplication.semantic.kmeans import KMeansStage
kmeans_stage = KMeansStage(
input_path="embeddings/",
output_path="clustered_embeddings/",
n_clusters=1000,
id_field="_curator_dedup_id",
embedding_field="embeddings",
).with_(
{
"KMeansStage": {
"ray_stage_spec": {
RayStageSpecKeys.USE_TASK_WEIGHTS: False,
}
}
}
)
Pipeline(name="kmeans", stages=[kmeans_stage]).run(
executor=RayActorPoolExecutor(),
)

This override is available when using KMeansStage directly. The high-level semantic deduplication workflows construct their K-means stage internally and keep weighted assignment enabled. Disabling weights is mainly useful for comparisons or troubleshooting; it can leave actors with uneven amounts of input data.

Comparison with Other Deduplication Methods

Compare semantic deduplication with other methods:

MethodReturn Value Optionsperform_removal ParameterWorkflow
ExactDuplicatesDuplicates (ID list only)❌ Not supported (must remain False; use TextDuplicatesRemovalWorkflow)Two-step (identification + removal workflow)
FuzzyDuplicatesDuplicates (ID list only)❌ Not supported (must remain False; use TextDuplicatesRemovalWorkflow)Two-step (identification + removal workflow)
TextSemanticDeduplicationWorkflowDuplicates or Clean Dataset✅ AvailableOne-step or two-step

Key Parameters

ParameterTypeDefaultDescription
input_pathstr | list[str]RequiredInput file, directory, glob, or list of paths. A single directory is scanned recursively; listed directories are scanned only at their top level.
input_filetypestr"parquet"Original text reader and removal format. Supported values are "parquet" and "jsonl"; generated embeddings are always Parquet.
input_file_extensionslist[str] | NoneNoneExtensions to discover. None or [] uses [".parquet"] for Parquet and [".jsonl", ".json"] for JSONL. A non-empty list replaces those defaults.
model_identifierstr"google/embeddinggemma-300m"Pre-trained model for embedding generation (vLLM backend)
hf_tokenstr | NoneNoneHugging Face read token for private or gated models
embedding_pretokenizeboolFalseWhether to pre-tokenize input before passing to vLLM
embedding_vllm_init_kwargsdictNoneAdditional keyword arguments passed to the vLLM LLM initializer
embedding_max_charsintNoneMaximum number of characters for text truncation
model_cache_dirstrNoneDirectory to cache model weights
n_clustersint100Number of clusters for k-means clustering
kmeans_max_iterint300Maximum iterations for clustering
kmeans_fit_data_fractionfloat | NoneNoneFraction of generated Parquet embedding files used to fit K-means. None automatically selects complete embedding files within the live GPU-memory budget.
epsfloat0.01Threshold for deduplication (higher = more aggressive)
which_to_keepstr"hard"Strategy for keeping duplicates (“hard”, “easy”, or “random”)
kmeans_embedding_output_dtypestr"float16"Precision used to store KMeans embeddings for Pairwise ("float16" or "float32")
pairwise_compute_dtypestr"float16"Pairwise multiplication precision: "auto", "float16", or "float32"
pairwise_batch_sizeint1024Positive batch size for the bounded similarity workspace
distance_metricstr"cosine"Distance metric for similarity (“cosine” or “l2”)
perform_removalboolTrueWhether to perform duplicate removal
text_fieldstr"text"Name of the text field in input data
id_fieldstr"_curator_dedup_id"Name of the ID field in the data
use_id_generatorboolFalseWhether to assign IDs with the shared ID generator. Set this to True when the input does not already contain id_field. The generated ID is removed from deduplicated output when output_fields is not set.
output_fieldslist[str] | NoneNoneFields to write for deduplicated output. When set, it takes precedence over automatic ID removal.
read_kwargsdict{}Additional reader keyword arguments. Include storage_options for input files on remote storage.
cache_kwargsdict{}Additional reader and writer arguments, including storage options, for intermediate files and duplicate-ID artifacts, including those under output_path.
write_kwargsdict{}Additional reader and writer arguments, including storage options, for final deduplicated output and reading duplicate-ID artifacts.

For a remote output_path, configure storage_options in both cache_kwargs and write_kwargs, even when cache_path is local. The options in cache_kwargs must also allow access to cache_path if it is remote.

The same discovery rules apply to TextSemanticDeduplicationWorkflow, the pre-computed-embedding SemanticDeduplicationWorkflow, and TextDuplicatesRemovalWorkflow. The extension filter does not infer the reader format: keep input_filetype="parquet" when reading a custom Parquet suffix such as .pq. See Input File Discovery for examples and recursion details.

TextSemanticDeduplicationWorkflow always writes Parquet embeddings before clustering, regardless of the original input_filetype. The lower-level SemanticDeduplicationWorkflow and KMeansStage accept precomputed Parquet or JSONL embeddings and use the parameter name fit_data_fraction.

Control Pairwise memory and precision

Pairwise comparison keeps a reusable N x B similarity workspace, where N is the cluster size and B is pairwise_batch_size. The default is 1024; a smaller positive value lowers peak memory at the cost of more matrix multiplications. A larger value also computes more self/later-row products before masking them, so the largest batch that fits is not necessarily the fastest. Only earlier-ranked rows are valid neighbors, and equal similarity scores deterministically choose the earliest-ranked row.

How Pairwise chooses duplicates

Input order is the preservation preference: A is preferred over B, B over C, and so on. Consider these normalized embeddings:

A [1.00, 0.00, 0.00, 0.00]
B [0.96, 0.28, 0.00, 0.00]
C [0.00, 1.00, 0.00, 0.00]
D [0.00, 0.00, 1.00, 0.00]
E [0.00, 0.96, 0.00, 0.28]

Because they have unit length, X @ X.T is their cosine-similarity matrix:

Query \ candidateABCDE
A1.00000.96000.00000.00000.0000
B0.96001.00000.28000.00000.2688
C0.00000.28001.00000.00000.9600
D0.00000.00000.00001.00000.0000
E0.00000.26880.96000.00001.0000

Only candidates to the left of the query’s diagonal are eligible. Pairwise therefore selects B -> A (0.96), C -> B (0.28), D -> A (0.00), and E -> C (0.96). With eps=0.1, the duplicate threshold is 1 - eps = 0.9, so the next stage removes query rows B and E while preserving A, C, and D. An earlier-ranked candidate remains eligible even if that candidate is itself later marked for removal.

KMeans fit, prediction, and distance calculations remain FP32. By default, KMeans casts normalized embeddings to FP16 and stores their bit patterns in uint16 list leaves because cuDF does not support numeric FP16 columns. Set kmeans_embedding_output_dtype="float32" to retain FP32 embeddings. The uint16 representation does not guarantee half-sized Parquet files because Parquet stores uint16 values using an INT32 physical type.

Pairwise infers storage precision from the embedding list leaves: uint16 leaves are decoded as FP16 bit patterns, FP32 leaves keep their precision, and legacy FP64 leaves are normalized to FP32. With pairwise_compute_dtype="auto", FP16 embeddings use FP16 multiplication and FP32 embeddings use FP32. FP32 embeddings can instead use "float16" multiplication; upcasting stored FP16 embeddings is rejected because it cannot restore the discarded precision.

FP16 scores remain FP16 through multiplication and reduction, then are promoted to FP32 because cuDF does not support an FP16 numeric output column. Pairwise records cumulative footer scan, read, rank, precision conversion, compute, and write time.

Use FP32 storage and compute when exact or near-exact similarity thresholds make FP16 rounding material, or when the selected nearest-neighbor ID is part of the required output contract.

Control K-means fitting memory

For the lower-level workflow and stage, the behavior of fit_data_fraction=None depends on the embedding input format:

  • Parquet: Selects as many complete files as fit the actor’s live GPU-memory budget. Only one contiguous float32 embedding array remains resident during fitting; each input frame is released after its embeddings are copied into that array.
  • JSONL: Fits all input files in one pass because row and embedding counts are not available without reading the data.

Parquet is recommended for large semantic-deduplication workloads. Its footer metadata sizes the fit array, while cuDF’s chunked reader bounds fit-pass memory and output column sizes. After fitting, Parquet rereads all files in bounded frames to predict and write every row, including when fit_data_fraction=1.0 or automatic sizing selects every file. The fit array is released before prediction and writing, so it does not compete with their temporary allocations. Prediction groups respect cuDF’s embedding child-column limit, but this is not a GPU-memory limit: a large group or metadata-heavy input can still exhaust memory during reading, prediction, or writing. An explicit JSONL fit fraction also rereads the input, but loads all files together for prediction.

Set a fraction in (0, 1] to choose the number of sampled files explicitly. Both formats sample round(fit_data_fraction * number_of_actor_files) complete files per actor, with a minimum of one, but their I/O differs:

  • Parquet: Exact footer counts size one contiguous float32 fit array. cuDF reads the sampled files in chunks that share the remaining GPU memory between output and decompression. The sampled frames are released during the first pass; the second pass rereads all files to predict and write every row.
  • JSONL: The first pass reads the sampled files and fits K-means. The second pass reads all files, including the sampled files, to predict and write every row.

Because sampling is file-based rather than row-based, the realized row fraction can differ from the configured value when file sizes vary. Every actor contributes at least one file, so very small fractions can also sample more data than expected on highly parallel runs. Choose a fraction that leaves enough representative rows for at least n_clusters centroids, and use random_state or kmeans_random_state to make file selection repeatable.

All input rows receive cluster assignments; the fraction affects only centroid fitting. For Parquet input, both fitting and prediction use exact footer element counts to keep each read below cuDF’s column-size limit.

from nemo_curator.stages.text.deduplication.semantic import TextSemanticDeduplicationWorkflow
workflow = TextSemanticDeduplicationWorkflow(
input_path="input_data/",
output_path="results/",
cache_path="semdedup_cache/",
n_clusters=1000,
kmeans_fit_data_fraction=0.1,
)
workflow.run()

KMeansStage.cache_path controls only centroid persistence. When it is set, actor 0 writes kmeans_centroids.npy after fitting; when it is None, centroids are not saved. This differs from SemanticDeduplicationWorkflow.cache_path, which stores the workflow’s K-means and pairwise intermediate results. The workflow intentionally does not forward its cache path to KMeansStage.cache_path, so use KMeansStage directly when you need the centroid array.

Similarity Threshold

Control deduplication aggressiveness with eps:

  • Lower values (such as 0.001): More strict, less deduplication, higher confidence
  • Higher values (such as 0.1): Less strict, more aggressive deduplication

Experiment with different values to balance data reduction and dataset diversity.

Embedding generation uses vLLM as the inference backend. The default model is google/embeddinggemma-300m.

Default (vLLM):

workflow = TextSemanticDeduplicationWorkflow(
# Uses google/embeddinggemma-300m by default
input_path="input_data/",
output_path="./results",
cache_path="./sem_cache",
)

Custom model with vLLM options:

workflow = TextSemanticDeduplicationWorkflow(
model_identifier="google/embeddinggemma-300m",
embedding_pretokenize=True,
embedding_vllm_init_kwargs={"enforce_eager": True, "max_model_len": 2048},
# ... other parameters
)

vLLM Embedder (recommended for large models):

For large embedding models, you can generate embeddings separately using VLLMEmbeddingModelStage before running the deduplication workflow. This provides better GPU utilization and throughput for models with 500M+ parameters. See vLLM Embedder for details.

Generate embeddings with VLLMEmbeddingModelStage using the vLLM Embedder pipeline, then pass the output to SemanticDeduplicationWorkflow:

from nemo_curator.stages.deduplication.semantic import SemanticDeduplicationWorkflow
# After generating embeddings to embedding_output_path using VLLMEmbeddingModelStage
semantic_workflow = SemanticDeduplicationWorkflow(
input_path=embedding_output_path,
output_path=output_path,
n_clusters=100,
eps=0.07,
id_field="_curator_dedup_id",
embedding_field="embeddings",
)
semantic_workflow.run()
# Step 3: Filter original text dataset using the IDs to remove
# See TextDuplicatesRemovalWorkflow for the removal step

When choosing a model:

  • Use models that support vLLM pooling (embedding) mode
  • Choose models appropriate for your language or domain
  • Prefer models trained for sentence embeddings (for example, EmbeddingGemma, E5, BGE, or SBERT)
  • Use embedding_pretokenize=True for models that benefit from explicit tokenization control
  • Pass additional vLLM configuration through embedding_vllm_init_kwargs
  • For more control over the embedding process, consider using VLLMEmbeddingModelStage separately
workflow = TextSemanticDeduplicationWorkflow(
# I/O
input_path="input_data/",
output_path="results/",
cache_path="semdedup_cache",
input_filetype="jsonl", # Discovers .jsonl and .json by default
# Embedding generation (vLLM backend)
text_field="text",
model_identifier="google/embeddinggemma-300m",
embedding_pretokenize=False,
embedding_max_chars=None,
model_cache_dir=None,
# Deduplication
n_clusters=100,
eps=0.01, # Similarity threshold
distance_metric="cosine",
which_to_keep="hard",
# K-means
kmeans_max_iter=300,
kmeans_tol=1e-4,
kmeans_fit_data_fraction=0.1,
pairwise_batch_size=1024,
perform_removal=True,
)

GPU write tuning:

NeMo Curator defaults KVIKIO_AUTO_DIRECT_IO_WRITE=0 because KvikIO direct writes can reduce throughput on distributed filesystems such as Lustre. If the workflow writes to a local SSD, setting the variable to 1 may improve throughput:

export KVIKIO_AUTO_DIRECT_IO_WRITE=1

Set the variable before importing NeMo Curator or starting the job, ensure every distributed worker inherits it, and benchmark both values against the target filesystem. The setting affects the workflow’s GPU-backed cuDF intermediate and output writes, not general CPU-based writes. See Tune GPU Writes for the Target Filesystem for details.

Output Format

The semantic deduplication process produces the following directory structure in your configured cache_path:

cache_path/
├── embeddings/ # Embedding outputs
│ └── *.parquet # Parquet files containing document embeddings
├── semantic_dedup/ # Semantic deduplication cache
│ ├── kmeans_results/ # K-means clustering outputs
│ │ └── embs_by_nearest_center/ # Embeddings organized by cluster
│ │ └── nearest_cent={0..n-1}/ # Subdirectories for each cluster
│ │ └── *.parquet # Cluster member embeddings
│ └── pairwise_results/ # Pairwise similarity results
│ └── *.parquet # Similarity scores by cluster
└── output_path/
├── duplicates/ # Duplicate identification results
│ └── *.parquet # Document IDs to remove
└── deduplicated/ # Final clean dataset (if perform_removal=True)
└── *.parquet # Deduplicated documents

File Formats

The workflow produces these output files:

  1. Document Embeddings (embeddings/*.parquet):

    • Contains document IDs and their vector embeddings
    • Format: Parquet files with columns: [id_column, embedding_column]
  2. Cluster Assignments (semantic_dedup/kmeans_results/):

    • embs_by_nearest_center/: Parquet files containing cluster members
    • Format: Parquet files with columns: [id_column, embedding_column, cluster_id]

    A direct KMeansStage(cache_path="kmeans_cache/") additionally writes the fitted cluster centers to kmeans_cache/kmeans_centroids.npy. The workflow wrappers do not save this file.

  3. Duplicate IDs (output_path/duplicates/*.parquet):

    • IDs of documents identified as duplicates for removal
    • Format: Parquet file with columns: ["id"]
    • Important: Contains only the IDs of documents to remove, not the full document content
    • When perform_removal=True, clean dataset is saved to output_path/deduplicated/

Performance characteristics:

  • Computationally intensive, especially for large datasets
  • GPU acceleration required for embedding generation and clustering
  • Benefits often outweigh upfront cost (reduced training time, improved model performance)

GPU requirements:

  • NVIDIA GPU with CUDA support
  • Sufficient GPU memory (recommended: >8GB for medium datasets)
  • RAPIDS libraries (cuDF) for GPU-accelerated dataframe operations
  • CPU-only processing not supported

Performance tuning:

  • Adjust n_clusters based on dataset size and available resources
  • Use batched cosine similarity to reduce memory requirements
  • Consider distributed processing for very large datasets
Dataset SizeGPU MemoryProcessing TimeRecommended GPUs
<100K docs4-8 GB1-2 hoursRTX 3080, A100
100K-1M docs8-16 GB2-8 hoursRTX 4090, A100
>1M docs>16 GB8+ hoursA100, H100

For more details, see the SemDeDup paper by Abbas et al.

ID Generator for large-scale operations:

from nemo_curator.stages.deduplication.id_generator import (
create_id_generator_actor,
write_id_generator_to_disk,
kill_id_generator_actor
)
create_id_generator_actor()
id_generator_path = "semantic_id_generator.json"
write_id_generator_to_disk(id_generator_path)
kill_id_generator_actor()
# Use persisted ID generator in removal workflow
removal_workflow = TextDuplicatesRemovalWorkflow(
input_path=input_path,
ids_to_remove_path=duplicates_path,
output_path=output_path,
id_generator_path=id_generator_path,
input_files_per_partition=1, # Match partitioning as embedding generation
# ... other parameters
)

Critical requirements:

  • Use the same input configuration (file paths, partitioning) across all stages
  • ID consistency maintained by hashing filenames in each task
  • Mismatched partitioning causes ID lookup failures

Ray backend configuration:

from nemo_curator.core.client import RayClient
client = RayClient(
num_cpus=64, # Adjust based on available cores
num_gpus=4 # Should be roughly 2x the memory of embeddings
)
client.start()
try:
workflow = TextSemanticDeduplicationWorkflow(
input_path=input_path,
output_path=output_path,
cache_path=cache_path,
# ... other parameters
)
result = workflow.run()
# result.metadata contains: total_time, num_duplicates, num_duplicates_removed, embedding_time, identification_time, removal_time, final_output_path
finally:
client.stop()

Provides distributed processing, memory management, and fault tolerance.