Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1,274 changes: 1,084 additions & 190 deletions Cargo.lock

Large diffs are not rendered by default.

3 changes: 2 additions & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -37,9 +37,10 @@ tokio-executor-trait = "2.1"
tokio-reactor-trait = "1.1"
chrono = { version = "0.4.41", features = ["serde"] }
serde_yaml = "0.9.34"
whatlang = "0.16.4"
anyhow = "1.0.98"
tokenizers = { version = "0.21", features = ["http"] } # Specify version and add http feature
html-escape = "0.2.13"
lingua = "1.7.2"

# {{ Define the library (implicitly done by having src/lib.rs) }}
# [lib]
Expand Down
26 changes: 13 additions & 13 deletions config/pipeline_config.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -2,18 +2,6 @@ pipeline:
- type: LanguageDetectionFilter
min_confidence: 0.65
allowed_languages: [ "dan" ]

- type: C4QualityFilter
split_paragraph: false
remove_citations: true
filter_no_terminal_punct: true
min_num_sentences: 10
min_words_per_line: 3
max_word_length: 15
filter_lorem_ipsum: true
filter_javascript: true
filter_curly_bracket: true
filter_policy: true

- type: GopherRepetitionFilter
dup_line_frac: 0.3
Expand All @@ -34,7 +22,7 @@ pipeline:

- type: GopherQualityFilter
min_doc_words: 50
max_doc_words: 1000000
max_doc_words: 100000
min_avg_word_length: 3.0
max_avg_word_length: 10.0
max_symbol_word_ratio: 0.1
Expand Down Expand Up @@ -65,6 +53,18 @@ pipeline:
"vore", "vores", "vær", "være", "været", "øh"
]

- type: C4QualityFilter
split_paragraph: true
remove_citations: true
filter_no_terminal_punct: true
min_num_sentences: 5
min_words_per_line: 3
max_word_length: 1000
filter_lorem_ipsum: true
filter_javascript: true
filter_curly_bracket: true
filter_policy: true

# - type: C4BadWordsFilter
# keep_fraction: 0.1 # Keep 10% of documents with bad words
# fail_on_missing_language: false # Continue if a language's bad word list is not found
Expand Down
2 changes: 1 addition & 1 deletion src/config/producer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,7 @@ pub struct Args {
pub text_column: String,

/// Optional ID column name in the Parquet file
#[arg(long)]
#[arg(long, default_value = "id")]
pub id_column: Option<String>,

/// RabbitMQ connection string (e.g., amqp://guest:guest@localhost:5672/%2f)
Expand Down
41 changes: 31 additions & 10 deletions src/pipeline/filters/c4_filters.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,26 +9,25 @@ use lazy_static::lazy_static;
use rand::rngs::StdRng;
use rand::{Rng, SeedableRng};
use regex::Regex;
use reqwest::blocking::Client; // For HTTP client
use reqwest::blocking::Client;
use std::collections::HashMap;
use std::fs::{self}; // Added fs::File
use std::path::PathBuf; // Added Path, PathBuf
use std::sync::Mutex; // Added Mutex for interior mutability
use std::sync::Mutex;
use tracing::info; // For HTTP client // Added Mutex for interior mutability

// Constants based on C4 and reference implementation
const END_PUNCTUATION: [char; 4] = ['.', '!', '?', '"'];
const END_PUNCTUATION: [char; 6] = ['.', '!', '?', '"', '\'', '”'];
const ELLIPSIS: &str = "...";

lazy_static! {
static ref POLICY_SUBSTRINGS: Vec<&'static str> = vec![
"cookie",
"cookies",
"terms of use",
"privacy policy",
"gdpr",
"ccpa",
"california consumer privacy act",
"do not sell my personal information",
"cookie policy",
"uses cookies",
"use of cookies",
"use cookies",
];
// Simple regex for Wikipedia-style citations like [1], [2, 3], [45]
static ref CITATION_REGEX: Regex = Regex::new(r"\[\d+(?:,\s*\d+)*\]").expect("Invalid regex");
Expand Down Expand Up @@ -156,14 +155,20 @@ impl ProcessingStep for C4QualityFilter {
split_into_sentences(&document.content) // Default to lines for now
};

if document.id == "danske-taler_2572" {
info!("{}", document.content);
}

let mut kept_lines: Vec<String> = Vec::new();
let mut filter_reasons: Vec<String> = Vec::new();

// Document-level checks that can fail early
if self.filter_lorem_ipsum && original_content.to_lowercase().contains("lorem ipsum") {
filter_reasons.push("lorem_ipsum".to_string());
}
if self.filter_curly_bracket && original_content.contains('{') {
if self.filter_curly_bracket
&& (original_content.contains('{') || original_content.contains('}'))
{
filter_reasons.push("curly_bracket".to_string());
}

Expand All @@ -181,6 +186,8 @@ impl ProcessingStep for C4QualityFilter {
});
}

let mut line_stats: HashMap<String, i32> = HashMap::new();

// Line-by-line filtering
for line in lines {
let current_line = line.trim().to_string();
Expand All @@ -202,6 +209,9 @@ impl ProcessingStep for C4QualityFilter {
.iter()
.any(|w| w.chars().count() > self.max_word_length)
{
*line_stats
.entry("line-filter-too_long_word".to_string())
.or_insert(0) += 1;
continue; // Drop line
}

Expand All @@ -214,12 +224,18 @@ impl ProcessingStep for C4QualityFilter {
let ends_with_ellipsis = processed_line.ends_with(ELLIPSIS);

if !ends_with_terminal_punct || ends_with_ellipsis {
*line_stats
.entry("line-filter-no_terminal_punc".to_string())
.or_insert(0) += 1;
continue; // Drop line
}
}

// min words per line
if self.min_words_per_line > 0 && words.len() < self.min_words_per_line {
*line_stats
.entry("line-filter-too_few_words".to_string())
.or_insert(0) += 1;
continue; // Drop line
}

Expand Down Expand Up @@ -261,6 +277,11 @@ impl ProcessingStep for C4QualityFilter {
document
.metadata
.insert("c4_filter_reasons".to_string(), reasons_string.clone());

for (key, value) in &line_stats {
document.metadata.insert(key.to_string(), value.to_string());
}

Err(PipelineError::DocumentFiltered {
document: Box::new(document),
reason: reasons_string,
Expand Down
113 changes: 84 additions & 29 deletions src/pipeline/filters/fineweb_quality.rs
Original file line number Diff line number Diff line change
Expand Up @@ -23,9 +23,7 @@ use std::collections::HashSet;
const DEFAULT_FILTER_NAME: &str = "FineWebQualityFilter";

// TERMINAL_PUNCTUATION in Python: (".", "?", "!", "。", "?", "!")
fn default_stop_chars() -> HashSet<char> {
vec!['.', '?', '!', '。', '?', '!'].into_iter().collect()
}
const END_PUNCTUATION: [char; 6] = ['.', '!', '?', '"', '\'', '”'];

#[derive(Debug)]
pub struct FineWebQualityFilter {
Expand All @@ -50,7 +48,7 @@ impl FineWebQualityFilter {
// language: String,
stop_chars: Option<HashSet<char>>,
) -> Self {
let sc = stop_chars.unwrap_or_else(|| default_stop_chars().iter().copied().collect());
let sc = stop_chars.unwrap_or_else(|| END_PUNCTUATION.into());
FineWebQualityFilter {
line_punct_thr,
line_punct_exclude_zero,
Expand All @@ -71,13 +69,20 @@ impl ProcessingStep for FineWebQualityFilter {
}

async fn process(&self, document: TextDocument) -> Result<TextDocument> {
let mut document = document;
let text_content = &document.content;
let lines: Vec<&str> = text_content
.lines()
.filter(|line| !line.trim().is_empty())
.collect();

if lines.is_empty() {
document
.metadata
.insert("fineweb_filter_status".into(), "filtered".into());
document
.metadata
.insert("fineweb_filter_reason".into(), "empty document".into());
return Err(PipelineError::DocumentFiltered {
document: Box::new(document),
reason: "empty".to_string(),
Expand All @@ -101,12 +106,19 @@ impl ProcessingStep for FineWebQualityFilter {
if line_punct_actual_ratio < self.line_punct_thr
&& !(line_punct_actual_ratio == 0.0 && self.line_punct_exclude_zero)
{
let reason = format!(
"line_punct_ratio: {:.4} < threshold {:.4} (exclude_zero: {})",
line_punct_actual_ratio, self.line_punct_thr, self.line_punct_exclude_zero
);
document
.metadata
.insert("fineweb_filter_status".into(), "filtered".into());
document
.metadata
.insert("fineweb_filter_reason".into(), reason.clone());
return Err(PipelineError::DocumentFiltered {
document: Box::new(document),
reason: format!(
"line_punct_ratio: {:.4} < threshold {:.4} (exclude_zero: {})",
line_punct_actual_ratio, self.line_punct_thr, self.line_punct_exclude_zero
),
reason,
});
}

Expand All @@ -117,30 +129,58 @@ impl ProcessingStep for FineWebQualityFilter {
.count();
let short_line_actual_ratio = short_lines_count as f64 / lines.len() as f64;
if short_line_actual_ratio > self.short_line_thr {
let reason = format!(
"short_line_ratio: {:.4} > threshold {:.4}",
short_line_actual_ratio, self.short_line_thr
);
document
.metadata
.insert("fineweb_filter_status".into(), "filtered".into());
document
.metadata
.insert("fineweb_filter_reason".into(), reason.clone());
return Err(PipelineError::DocumentFiltered {
document: Box::new(document),
reason: format!(
"short_line_ratio: {:.4} > threshold {:.4}",
short_line_actual_ratio, self.short_line_thr
),
reason,
});
}

// Character duplication ratio
let vec_line: Vec<String> = lines.iter().map(|l| l.to_string()).collect();
let (repeated_char_count, total_chars_for_dup_ratio) = find_duplicates(&vec_line);
let char_dup_actual_ratio = if total_chars_for_dup_ratio > 0 {
repeated_char_count as f64 / total_chars_for_dup_ratio as f64
// 1. Filter non-empty lines
let non_empty_lines: Vec<String> = lines
.iter()
.filter(|line| !line.trim().is_empty())
.map(|line| line.to_string())
.collect();

// 2. Count total characters in doc (excluding newlines)
let total_chars_no_newlines: usize =
document.content.chars().filter(|&c| c != '\n').count();

// 3. Find duplicates among non-empty lines
let (_dup_count, dup_char_count) = find_duplicates(&non_empty_lines);

// 4. Compute ratio
let char_dup_actual_ratio = if total_chars_no_newlines > 0 {
dup_char_count as f64 / total_chars_no_newlines as f64
} else {
0.0
};

if char_dup_actual_ratio > self.char_duplicates_ratio {
let reason = format!(
"char_dup_ratio: {:.4} > threshold {:.4}",
char_dup_actual_ratio, self.char_duplicates_ratio
);
document
.metadata
.insert("fineweb_filter_status".into(), "filtered".into());
document
.metadata
.insert("fineweb_filter_reason".into(), reason.clone());
return Err(PipelineError::DocumentFiltered {
document: Box::new(document),
reason: format!(
"char_dup_ratio: {:.4} > threshold {:.4}",
char_dup_actual_ratio, self.char_duplicates_ratio
),
reason,
});
}

Expand All @@ -150,20 +190,34 @@ impl ProcessingStep for FineWebQualityFilter {

if words.is_empty() {
if new_line_count > 0 {
let reason = "list_ratio_no_words (newlines present but no words)".to_string();
document
.metadata
.insert("fineweb_filter_status".into(), "filtered".into());
document
.metadata
.insert("fineweb_filter_reason".into(), reason.clone());
return Err(PipelineError::DocumentFiltered {
document: Box::new(document),
reason: "list_ratio_no_words (newlines present but no words)".to_string(),
reason,
});
}
} else {
let list_actual_ratio = new_line_count as f64 / words.len() as f64;
if list_actual_ratio > self.new_line_ratio {
let reason = format!(
"list_ratio: {:.4} > threshold {:.4}",
list_actual_ratio, self.new_line_ratio
);
document
.metadata
.insert("fineweb_filter_status".into(), "filtered".into());
document
.metadata
.insert("fineweb_filter_reason".into(), reason.clone());
return Err(PipelineError::DocumentFiltered {
document: Box::new(document),
reason: format!(
"list_ratio: {:.4} > threshold {:.4}",
list_actual_ratio, self.new_line_ratio
),
reason,
});
}
}
Expand All @@ -190,7 +244,7 @@ mod tests {
FineWebQualityFilter {
line_punct_thr: DEFAULT_LINE_PUNCT_THR,
line_punct_exclude_zero: DEFAULT_LINE_PUNCT_EXCLUDE_ZERO,
stop_chars: default_stop_chars(),
stop_chars: END_PUNCTUATION.into(),
short_line_thr: DEFAULT_SHORT_LINE_THR,
short_line_length: DEFAULT_SHORT_LINE_LENGTH,
char_duplicates_ratio: 0.95, // Temporarily very high to isolate other test failures
Expand Down Expand Up @@ -389,17 +443,18 @@ mod tests {
filter.short_line_thr = 1.0;
filter.new_line_ratio = 1.0; // Allow many newlines

filter.char_duplicates_ratio = 0.9; // Specific for this test to fail
let content = "a\na\na."; // Ends with '.', 31 chars, 29 'a' repeats. Ratio ~0.935
filter.char_duplicates_ratio = 0.66; // Specific for this test to fail
let content = "Hello World\nHello World\nHello World"; // Ends with '.', 31 chars, 29 'a' repeats. Ratio ~0.935
let doc = create_test_doc("char_dup_all_same", content);

let result = filter.process(doc).await;
println!("{:?}", result);
assert!(result.is_err());
match result.err().unwrap() {
PipelineError::DocumentFiltered { reason, .. } => {
// Ratio 29/31 = 0.93548...
assert!(
reason.starts_with("char_dup_ratio: 1.0000 > threshold 0.9000"),
reason.starts_with("char_dup_ratio: 0.6667 > threshold 0.6600"),
"Actual reason: {}",
reason
);
Expand Down
Loading
Loading