diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 8170c24d..c43a801f 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -225,6 +225,41 @@ jobs: python -m build twine check dist/* + freshcore-native: + # Builds the optional FreshCore Rust extension (crates/freshcore) and runs + # the FreshCore tests against it. No other job installs the extension, so + # the native parity tests skip everywhere else. + runs-on: ubuntu-latest + timeout-minutes: 10 + env: + RUSTUP_TOOLCHAIN: "1.98.1" + steps: + - uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1 + - uses: actions/setup-python@5fda3b95a4ea91299a34e894583c3862153e4b97 # v7.0.0 + with: + python-version: "3.12" + cache: pip + - name: Install Rust toolchain + run: rustup toolchain install "$RUSTUP_TOOLCHAIN" --profile minimal + - name: Install + # maturin develop installs into a virtualenv, so create one in the + # workspace and put it first on PATH for the remaining steps. + run: | + python -m venv .venv + echo "VIRTUAL_ENV=$PWD/.venv" >> "$GITHUB_ENV" + echo "$PWD/.venv/bin" >> "$GITHUB_PATH" + .venv/bin/python -m pip install --upgrade pip + .venv/bin/pip install -c constraints/ci.txt -e ".[dev,freshcore]" + - name: Rust unit tests + run: cargo test --locked --manifest-path crates/freshcore/Cargo.toml + - name: Build the extension + run: maturin develop --manifest-path crates/freshcore/Cargo.toml --features extension-module + - name: FreshCore tests + # The coverage gate in addopts applies to the full suite, not this subset. + run: | + python -c "import freshdata_freshcore" + pytest tests/test_execution -k freshcore -o addopts="--strict-markers" -q + coverage-badge: # Publish a self-hosted shields endpoint badge to the `badges` branch # (no third-party account needed). Runs only on main. diff --git a/CHANGELOG.md b/CHANGELOG.md index 35e927f9..843855a8 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -387,6 +387,11 @@ adheres to [Semantic Versioning](https://semver.org/). - The FreshCore adapter reports native duplicate detections and applies `duplicate_ratio_action` to native drop counts, as the pandas pipeline does (#323, part). +- The FreshCore native module counts duplicate rows when `drop_duplicates` is + False, at the same stage as the pandas step, so `engine="freshcore"` records + the detection, warns above `duplicate_threshold` and raises + `DuplicateRatioError` under `duplicate_ratio_action="error"` without falling + back to pandas (#323). - The Spark engine renames columns without collisions, honours `duplicate_keep` and input order when deduplicating, reads float `NaN` as null with outlier fences from finite values only, and no longer treats diff --git a/crates/freshcore/src/kernels/duplicates.rs b/crates/freshcore/src/kernels/duplicates.rs index 04bda329..8ad93640 100644 --- a/crates/freshcore/src/kernels/duplicates.rs +++ b/crates/freshcore/src/kernels/duplicates.rs @@ -2,28 +2,43 @@ use std::collections::HashSet; use crate::arrays::{CellKey, Frame}; -pub fn drop_duplicates(frame: &mut Frame, keep_policy: &str) -> usize { - if frame.nrows == 0 { - return 0; - } +/// Marks full-row duplicates under `keep_policy` (`"last"` keeps the final +/// occurrence, anything else the first). +/// +/// Like pandas' `drop_duplicate_rows`, a frame with no rows or no columns has +/// no duplicates. +pub fn duplicated_mask(frame: &Frame, keep_policy: &str) -> Vec { let mut duplicated = vec![false; frame.nrows]; - if keep_policy == "last" { - let mut seen: HashSet> = HashSet::new(); - for row in (0..frame.nrows).rev() { - let key = row_key(frame, row); - if !seen.insert(key) { - duplicated[row] = true; - } + if frame.nrows == 0 || frame.columns.is_empty() { + return duplicated; + } + let mut seen: HashSet> = HashSet::new(); + let mut mark = |row: usize| { + if !seen.insert(row_key(frame, row)) { + duplicated[row] = true; } + }; + if keep_policy == "last" { + (0..frame.nrows).rev().for_each(&mut mark); } else { - let mut seen: HashSet> = HashSet::new(); - for row in 0..frame.nrows { - let key = row_key(frame, row); - if !seen.insert(key) { - duplicated[row] = true; - } - } + (0..frame.nrows).for_each(&mut mark); } + duplicated +} + +/// Number of full-row duplicates, without removing any rows. +/// +/// The count does not depend on the keep policy: every row beyond the first +/// occurrence of its key is a duplicate. +pub fn count_duplicates(frame: &Frame) -> usize { + duplicated_mask(frame, "first") + .iter() + .filter(|v| **v) + .count() +} + +pub fn drop_duplicates(frame: &mut Frame, keep_policy: &str) -> usize { + let duplicated = duplicated_mask(frame, keep_policy); let dropped = duplicated.iter().filter(|v| **v).count(); if dropped > 0 { let keep: Vec = duplicated.iter().map(|v| !*v).collect(); @@ -45,6 +60,36 @@ mod tests { use super::*; use crate::arrays::{Column, ColumnData, Frame}; + fn mixed_frame() -> Frame { + Frame::new(vec![ + Column { + name: "a".into(), + data: ColumnData::Float(vec![ + Some(1.0), + Some(1.0), + Some(0.0), + Some(-0.0), + None, + None, + Some(2.0), + ]), + }, + Column { + name: "b".into(), + data: ColumnData::Utf8(vec![ + Some("x".into()), + Some("x".into()), + Some("y".into()), + Some("y".into()), + None, + None, + Some("x".into()), + ]), + }, + ]) + .unwrap() + } + #[test] fn removes_duplicate_rows() { let mut frame = Frame::new(vec![Column { @@ -55,4 +100,54 @@ mod tests { assert_eq!(drop_duplicates(&mut frame, "first"), 1); assert_eq!(frame.nrows, 2); } + + #[test] + fn counts_duplicates_without_removing_rows() { + let frame = mixed_frame(); + // (1.0, x) twice, (0.0, y) == (-0.0, y), and (null, null) twice. + assert_eq!(count_duplicates(&frame), 3); + assert_eq!(frame.nrows, 7); + assert_eq!(frame, mixed_frame()); + } + + #[test] + fn count_matches_rows_dropped_for_every_keep_policy() { + for keep in ["first", "last"] { + let mut frame = mixed_frame(); + let counted = count_duplicates(&frame); + assert_eq!(drop_duplicates(&mut frame, keep), counted); + assert_eq!(count_duplicates(&frame), 0); + } + } + + #[test] + fn mask_marks_the_occurrences_the_keep_policy_drops() { + let frame = mixed_frame(); + assert_eq!( + duplicated_mask(&frame, "first"), + vec![false, true, false, true, false, true, false] + ); + assert_eq!( + duplicated_mask(&frame, "last"), + vec![true, false, true, false, true, false, false] + ); + } + + #[test] + fn frames_without_rows_or_columns_have_no_duplicates() { + let no_rows = Frame::new(vec![Column { + name: "x".into(), + data: ColumnData::Float(Vec::new()), + }]) + .unwrap(); + assert_eq!(count_duplicates(&no_rows), 0); + + let mut no_columns = Frame { + columns: Vec::new(), + nrows: 3, + }; + assert_eq!(count_duplicates(&no_columns), 0); + assert_eq!(drop_duplicates(&mut no_columns, "first"), 0); + assert_eq!(no_columns.nrows, 3); + } } diff --git a/crates/freshcore/src/python.rs b/crates/freshcore/src/python.rs index d24f26f3..92b1b3ab 100644 --- a/crates/freshcore/src/python.rs +++ b/crates/freshcore/src/python.rs @@ -25,6 +25,7 @@ fn execute_plan(py: Python<'_>, payload: &Bound<'_, PyDict>) -> PyResult = None; let mut outliers_handled = 0usize; let t = Instant::now(); @@ -123,6 +124,13 @@ fn execute_plan(py: Python<'_>, payload: &Bound<'_, PyDict>) -> PyResult, payload: &Bound<'_, PyDict>) -> PyResult dict: + return { + "actions": [ + (a.description if descriptions else None, a.count, a.risk) + for a in report.actions + if a.step == "drop_duplicates" + ], + "warnings": [w for w in report.warnings if "duplicate" in w], + "recommendations": [r for r in report.recommendations if "duplicated" in r], + } + + +def _detections(report) -> list[str]: + return [a.description for a in report.actions if a.step == "drop_duplicates"] + + +# -- the #323 reproduction --------------------------------------------------- + + +def _repro_frame() -> pd.DataFrame: + """10 rows, 9 of them duplicates.""" + return pd.DataFrame({"a": [1.0] * 10, "b": ["x"] * 10}) + + +REPRO = {"fix_dtypes": False} + + +def test_issue_323_repro_raises_natively_under_error(native): + with pytest.raises(DuplicateRatioError) as expected: + fd.clean(_repro_frame(), engine="pandas", duplicate_ratio_action="error", **KW, **REPRO) + # fallback_policy="error" turns any pandas fallback into FallbackError, so + # DuplicateRatioError here comes from the native duplicate count. + with pytest.raises(DuplicateRatioError) as raised: + fd.clean( + _repro_frame(), + engine="freshcore", + duplicate_ratio_action="error", + fallback_policy="error", + **KW, + **REPRO, + ) + assert native.calls == 1 + assert str(raised.value) == str(expected.value) + + +def test_issue_323_repro_warns_like_pandas_under_warn(native): + expected, report = _clean_both(native, _repro_frame(), duplicate_ratio_action="warn", **REPRO) + + assert _duplicate_view(report) == _duplicate_view(expected) + assert _detections(report) == ["detected 9 duplicate row(s) (90.0%), none removed"] + assert len(_duplicate_view(report)["warnings"]) == 1 + assert len(_duplicate_view(report)["recommendations"]) == 1 + + +def test_ignore_is_rejected_on_both_engines(native): + # duplicate_ratio_action accepts only "warn" and "error"; "ignore" fails + # config validation before either engine runs. + for engine in ("pandas", "freshcore"): + with pytest.raises(ValueError, match="duplicate_ratio_action must be one of"): + fd.clean( + _repro_frame(), engine=engine, duplicate_ratio_action="ignore", **KW, **REPRO + ) + assert native.calls == 0 + + +@pytest.mark.parametrize("action", ["warn", "error"]) +@pytest.mark.parametrize( + ("rows", "description"), + [ + pytest.param(20, "detected 1 duplicate row(s) (5.0%), none removed", id="below"), + pytest.param(10, "detected 1 duplicate row(s) (10.0%), none removed", id="at-threshold"), + ], +) +def test_ratio_within_threshold_reports_without_warning(native, action, rows, description): + df = pd.DataFrame({"a": [float(i) for i in range(rows - 1)] + [0.0], "b": ["x"] * rows}) + expected, report = _clean_both(native, df, duplicate_ratio_action=action, **REPRO) + + assert _duplicate_view(report) == _duplicate_view(expected) + assert _detections(report) == [description] + assert _duplicate_view(report)["warnings"] == [] + assert _duplicate_view(report)["recommendations"] == [] + + +def test_no_duplicates_records_no_detection(native): + df = pd.DataFrame({"a": [" x", "y", "z"], "v": [1.0, 2.0, 3.0]}) + expected, report = _clean_both(native, df, **REPRO) + assert _detections(report) == _detections(expected) == [] + + +# -- counts taken after cleaning, like the pandas step ---------------------- + + +CLEANING_CASES = [ + pytest.param( + pd.DataFrame( + {"name": [" alice", "alice ", "alice", "bob", "carol"], "v": [1.0, 1.0, 1.0, 2, 3]} + ), + {"fix_dtypes": False}, + 2, + id="whitespace", + ), + pytest.param( + pd.DataFrame({"code": ["N/A", "", "-", "x", "y"], "v": [1.0, 1.0, 1.0, 2, 3]}), + {"fix_dtypes": False}, + 2, + id="sentinels", + ), + pytest.param( + pd.DataFrame({"name": ["Alice", "ALICE", "alice", "Bob"], "v": [1.0, 1.0, 1.0, 1.0]}), + {"fix_dtypes": False, "string_case": "lower"}, + 2, + id="string-case", + ), + pytest.param( + pd.DataFrame( + { + "a": [1.0, 1.0, 1.0, 2, 3, 4, 5, 6, 7, 8, None, None], + "b": ["x", "x", "x", "y", "y", "y", "y", "y", "y", "y", None, None], + } + ), + {"fix_dtypes": False}, + 2, + id="empty-rows-leave-the-denominator", + ), + pytest.param( + pd.DataFrame({"a": [" x", "x", None, None], "b": ["N/A", None, "", None]}), + {"fix_dtypes": False}, + 1, + id="sentinels-make-empty-rows", + ), + pytest.param( + pd.DataFrame( + {"n": ["1", "1.0", "2", "2.00", "3", "4"], "k": ["a", "a", "b", "b", "c", "d"]} + ), + {"fix_dtypes": True}, + 2, + id="numeric-casts", + ), +] + + +@pytest.mark.parametrize(("df", "options", "n_dup"), CLEANING_CASES) +def test_detected_count_matches_pandas_after_cleaning(native, df, options, n_dup): + # The raw frame has a different count: only cleaning makes these rows equal. + assert int(df.duplicated().sum()) != n_dup + expected, report = _clean_both(native, df, **options) + + assert _duplicate_view(report) == _duplicate_view(expected) + (detection,) = _detections(report) + assert detection.startswith(f"detected {n_dup} duplicate row(s) (") + + +@pytest.mark.parametrize("keep", ["first", "last"]) +@pytest.mark.parametrize(("df", "options", "n_dup"), CLEANING_CASES) +def test_dropped_count_matches_pandas_after_cleaning(native, df, options, n_dup, keep): + expected, report = _clean_both( + native, df, drop_duplicates=True, duplicate_keep=keep, **options + ) + + assert report.duplicates_removed == expected.duplicates_removed == n_dup + # Drop descriptions spell the keep policy differently (keep="first" vs + # keep='first'); counts, risk, warnings and recommendations must match. + assert _duplicate_view(report, descriptions=False) == _duplicate_view( + expected, descriptions=False + ) + + +# -- the native result contract --------------------------------------------- + + +def _native_result(df: pd.DataFrame, **options) -> dict: + config = fd.CleanConfig(**KW, **options) + return native_module.execute_plan(FreshCoreEngine()._payload(df, config)) + + +def test_native_result_reports_detection_count_and_stage(): + df = pd.DataFrame({"a": [" x", "x", "y"], "v": [1.0, 1.0, None]}) + + detected = _native_result(df, fix_dtypes=False, impute="mean") + assert detected["duplicates_detected"] == 1 + assert detected["duplicates_removed"] == 0 + assert detected["rows_after"] == 3 + stages = [stage for stage, _ in detected["stage_timings"]] + assert "drop_duplicates" not in stages + # Counted after casts and before imputation, where the pandas step runs. + assert stages.index("fix_dtypes") < stages.index("detect_duplicates") < stages.index("impute") + + dropped = _native_result(df, fix_dtypes=False, impute="mean", drop_duplicates=True) + assert dropped["duplicates_detected"] is None + assert dropped["duplicates_removed"] == 1 + assert "detect_duplicates" not in [stage for stage, _ in dropped["stage_timings"]] + + +def test_detection_counts_before_imputation_like_pandas(native): + # Imputation would make rows 1 and 2 equal; neither engine counts them. + df = pd.DataFrame({"a": [1.0, 2.0, None, 2.0], "b": ["x", "y", "y", "z"]}) + expected, report = _clean_both(native, df, impute="mean", **REPRO) + assert _detections(report) == _detections(expected) == []