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
35 changes: 35 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
5 changes: 5 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
131 changes: 113 additions & 18 deletions crates/freshcore/src/kernels/duplicates.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<bool> {
let mut duplicated = vec![false; frame.nrows];
if keep_policy == "last" {
let mut seen: HashSet<Vec<CellKey>> = 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<Vec<CellKey>> = 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<Vec<CellKey>> = 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<bool> = duplicated.iter().map(|v| !*v).collect();
Expand All @@ -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 {
Expand All @@ -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);
}
}
9 changes: 9 additions & 0 deletions crates/freshcore/src/python.rs
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ fn execute_plan(py: Python<'_>, payload: &Bound<'_, PyDict>) -> PyResult<PyObjec
let mut columns_dropped = Vec::new();
let mut columns_imputed = Vec::new();
let mut duplicates_removed = 0usize;
let mut duplicates_detected: Option<usize> = None;
let mut outliers_handled = 0usize;

let t = Instant::now();
Expand Down Expand Up @@ -123,6 +124,13 @@ fn execute_plan(py: Python<'_>, payload: &Bound<'_, PyDict>) -> PyResult<PyObjec
));
}
timings.push(("drop_duplicates".to_string(), t.elapsed().as_secs_f64()));
} else {
// Detection only: count at the same stage as the pandas step (after
// string cleaning, empty-row removal and casts; before imputation and
// outliers) so the adapter can report and escalate like pandas.
let t = Instant::now();
duplicates_detected = Some(duplicates::count_duplicates(&frame));
timings.push(("detect_duplicates".to_string(), t.elapsed().as_secs_f64()));
}

let t = Instant::now();
Expand Down Expand Up @@ -180,6 +188,7 @@ fn execute_plan(py: Python<'_>, payload: &Bound<'_, PyDict>) -> PyResult<PyObjec
result.set_item("missing_before", missing_before)?;
result.set_item("missing_after", frame.null_cells())?;
result.set_item("duplicates_removed", duplicates_removed)?;
result.set_item("duplicates_detected", duplicates_detected)?;
result.set_item("outliers_handled", outliers_handled)?;
result.set_item("columns_dropped", columns_dropped)?;
result.set_item("columns_imputed", columns_imputed)?;
Expand Down
8 changes: 5 additions & 3 deletions docs/backends.md
Original file line number Diff line number Diff line change
Expand Up @@ -176,9 +176,11 @@ Notes:
categorical, period and interval columns, for integer columns holding values
beyond ±2\*\*53, and for column labels that collide once stringified (such as
`1` and `"1"`). Other integer columns are cast back to their input dtype, and
non-string labels come back unchanged. With `drop_duplicates=False` and
`duplicate_ratio_action="error"`, it falls back unless the native module
reports a duplicate-row count (see the [fallback matrix](fallback-matrix.md)).
non-string labels come back unchanged. With `drop_duplicates=False`, the
native module counts duplicate rows at the pandas dedup stage, so detection,
the `duplicate_threshold` warning and `duplicate_ratio_action="error"` match
pandas. Native modules built before that count existed fall back under
`duplicate_ratio_action="error"` (see the [fallback matrix](fallback-matrix.md)).
- FreshCore v1 is a cleaning-first native engine, not an out-of-core engine. It
supports pandas-compatible materialized outputs and records per-stage timings
in `report.stage_timings`.
Expand Down
2 changes: 1 addition & 1 deletion docs/fallback-matrix.md
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@ FreshCore also runs its own config and data checks in
| full-row dedup (`keep="first"/"last"`) | native | native | native | native | streaming polars dedup drops row order (disclosed); `streaming_dedup=False` restores it. Spark keeps the first/last row in the input DataFrame's partition order |
| **subset dedup** (`duplicate_subset=`) | **native** | pandas | pandas | pandas | keep semantics are order-sensitive; Polars reproduces them via order-preserving `unique` (eager, not streaming — disclosed) |
| dedup `keep="drop"/"aggregate"` | pandas | pandas | pandas | pandas | group-wise resolution isn't expressed natively yet |
| detection-only dedup (`drop_duplicates=False`) with `duplicate_ratio_action="error"` | native | native | native | pandas unless the native module reports `duplicates_detected` | the escalation needs a duplicate-row count; FreshCore builds that don't report one would never raise |
| detection-only dedup (`drop_duplicates=False`) with `duplicate_ratio_action="error"` | native | native | native | native (pandas with native modules that don't report `duplicates_detected`) | the escalation needs the duplicate-row count at the pandas dedup stage; FreshCore counts it natively, but modules built before that count existed would never raise |
| global impute mean/median/mode | native | native | native | native | — |
| impute `mode`/`auto` with missing values in a nullable `boolean` column | native | native | native | pandas | FreshCore v1 kernels do not impute boolean columns |
| per-column `impute_strategy` | pandas | pandas | pandas | pandas | unimplemented natively (no fundamental blocker) |
Expand Down
29 changes: 26 additions & 3 deletions docs/freshcore.md
Original file line number Diff line number Diff line change
Expand Up @@ -55,9 +55,17 @@ aggregate/drop duplicate modes, model-based outliers, constant-column dropping,
memory downcasting, and the balanced/aggressive decision engine. It also falls
back for `impute="missforest"`, per-column `impute_strategy`, outlier handling
on float columns holding `±inf`, and mode/auto imputation of nullable boolean
columns with missing values. Detection-only dedup under
`duplicate_ratio_action="error"` falls back unless the native module reports
`duplicates_detected`.
columns with missing values.

With `drop_duplicates=False` (the default), the native module counts full-row
duplicates at the same stage as the pandas step: after string cleaning,
empty-row removal and casts, and before imputation and outliers. It returns
the count as `duplicates_detected` and records a `detect_duplicates` stage
timing. The adapter reports it like pandas: a "detected N duplicate row(s)"
action, a warning above `duplicate_threshold`, and `DuplicateRatioError` under
`duplicate_ratio_action="error"`. Native modules built before this count
existed don't report `duplicates_detected`; with those, detection is skipped,
and under `duplicate_ratio_action="error"` the adapter falls back to pandas.

The native arrays carry only float, bool and string values, so FreshCore also
falls back for datetime, timedelta, categorical, period and interval columns,
Expand All @@ -69,6 +77,21 @@ in range; otherwise they come back as `float64` and the change is recorded in
`report.backend_differences`. Non-string column labels such as `0` and `1`
come back unchanged.

## Building and testing

The CI `freshcore-native` job runs the same steps:

```bash
pip install -e ".[dev,freshcore]"
cargo test --manifest-path crates/freshcore/Cargo.toml
maturin develop --manifest-path crates/freshcore/Cargo.toml --features extension-module
pytest tests/test_execution -k freshcore
```

`tests/test_execution/test_freshcore_native_parity.py` runs the real extension
against the pandas reference and is skipped when `freshdata_freshcore` is not
installed. The other FreshCore tests use a fake native module.

## Benchmarking

Build the native module before benchmarking:
Expand Down
4 changes: 3 additions & 1 deletion src/freshdata/execution/backends/_freshcore.py
Original file line number Diff line number Diff line change
Expand Up @@ -93,7 +93,9 @@ def execute(
and native.get("duplicates_detected") is None
):
# Detection-only dedup needs the native duplicate count to honour
# the escalation; without it the error could never fire.
# the escalation. Current native modules always report it when
# drop_duplicates is False; modules built before #323 do not, and
# without it the error could never fire.
return self._fallback(
source,
config,
Expand Down
4 changes: 3 additions & 1 deletion tests/test_execution/test_freshcore_engine.py
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,9 @@ def test_engine_config_accepts_freshcore():
assert EngineSelector.get_engine("freshcore", cfg).name == "freshcore"


def test_missing_native_module_falls_back_to_pandas():
def test_missing_native_module_falls_back_to_pandas(monkeypatch):
# Simulate a missing extension so the test holds where it is built.
monkeypatch.setattr(FreshCoreEngine, "_load_native", staticmethod(lambda: None))
df = pd.DataFrame({"name": [" Alice ", "Bob"], "empty": [None, None]})
out, report = fd.clean(df, config=_cfg(), engine="freshcore", return_report=True)

Expand Down
Loading
Loading