Skip to content

Scalable execution backends

freshdata is pandas-first, but the same clean can run on Polars, DuckDB, Spark, or the optional FreshCore native engine. Every backend produces the same CleanReport audit contract — identical action schema (step, column, count, rationale, risk, confidence) — so downstream consumers (compliance, integrations, trust scoring) work unchanged.

import freshdata as fd

# in-memory pandas (default, unchanged)
clean = fd.clean(df)

# scale-out engines, result materialized to a pandas frame
clean = fd.clean("data.parquet", engine="duckdb", output_format="pandas")
clean = fd.clean(polars_df, engine="polars")
clean = fd.clean(spark_df,  engine="spark")          # or engine="auto"
clean = fd.clean(df,        engine="freshcore")      # optional native extension

# honest out-of-core: keep a native, un-materialized handle (you decide when to pull rows)
rel = fd.clean("data.parquet", engine="duckdb", output_format="duckdb")        # DuckDBPyRelation
lf  = fd.clean(polars_df,      engine="polars", output_format="polars-lazy")   # pl.LazyFrame

What "out-of-core" honestly means here

DuckDB and Polars can spill to disk during the cleaning pipeline, but the output format decides whether the cleaned result is then pulled fully into memory:

output_format Returns Materializes whole result?
"pandas" (default) pandas.DataFrame Yes — fetchdf() / collect()
"polars" / "arrow" / "spark" eager frame / Arrow table / Spark frame Yes
"duckdb" DuckDBPyRelation (un-fetched) No — you call .fetchdf()/.arrow()
"polars-lazy" pl.LazyFrame (un-collected) No — you call .collect()

A native handle comes from its own engine: "duckdb" needs engine="duckdb" and "polars-lazy" needs engine="polars". With engine="auto" (or no engine), freshdata picks that engine for you; any other pairing raises ValueError rather than returning a different type.

Neither handle fetches/collects the result until you ask, but they are not equal during the pipeline: the DuckDB path keeps peak memory well below the eager equivalent, while the Polars pipeline currently collects intermediates eagerly between stages, so "polars-lazy" defers only the final materialization — peak RSS during cleaning is comparable to eager output (measured; reproduce with python benchmarks/bench_outofcore.py). When a clean returns a native handle, report.materialized is False and report.summary() says so plainly. If the requested strategy needs the pandas decision engine (e.g. balanced/ aggressive imputation, dtype heuristics), the backend transparently falls back to pandas — that fallback is recorded in report.fallback_events, and the result is materialized. To keep the native handle use strategy="conservative" (deterministic representation repair + structural reduction) and fix_dtypes=False — dtype fixing relies on sampled pandas heuristics and forces the fallback even under conservative. Peak-RSS evidence for the four engine/output combinations is reproducible with python benchmarks/bench_outofcore.py.

Semantic cleaning stays native too. When semantic_mode is enabled on a Polars/DuckDB engine, the semantic stage runs over a natively extracted distinct table (a bounded GROUP BY) and maps repairs back with replace/SQL — the full frame is never pulled into pandas just to inspect values. See the native-engine semantic notes for the representation edges (e.g. partially-mapped boolean columns).

The StreamingCleaner micro-batch path (see Streaming) is the other genuinely out-of-core route: rows are processed one bounded batch at a time and never concatenated.

The pandas backend is the reference implementation. Native backends reproduce the deterministic subset directly; anything outside it is delegated to pandas and recorded in report.fallback_events. Every native step records a report.backend_differences entry when its statistics (e.g. quantile interpolation) can differ from the pandas reference.

Selecting a backend

engine="auto" resolves a concrete backend from the input:

Input auto picks
Spark DataFrame spark
.parquet / .csv path duckdb
Polars DataFrame/LazyFrame polars
Arrow Table / RecordBatch polars (else duckdb)
DuckDB relation duckdb
pandas DataFrame sized: pandas → polars → duckdb

EngineConfig controls execution (never what is cleaned):

from freshdata.execution import EngineConfig

import os

cfg = EngineConfig(engine="duckdb", memory_limit_gb=4,
                   temp_directory=os.path.expanduser("~/scratch/freshdata-spill"))
cfg = EngineConfig(engine="spark", spark_shuffle_partitions=200, output_format="spark")

DuckDB spill files contain rows of the data being cleaned, so each run spills into its own private (0700) subdirectory, removed when the run's connection closes (for output_format="duckdb", when the returned relation is released). By default the subdirectory is created under $FRESHDATA_SPILL_DIR, or under the per-user cache directory (~/.cache/freshdata/spill or $XDG_CACHE_HOME/freshdata/spill on Linux, ~/Library/Caches/freshdata/spill on macOS, %LOCALAPPDATA%\freshdata\spill on Windows), falling back to the system temp directory only when that is not writable. An explicit temp_directory is created with mode 0700 if missing; one that is not owned by you, or is group/other-writable without the sticky bit, raises PermissionError.

PySpark is an optional dependency (pip install 'freshdata-cleaner[spark]') and also needs a JVM at runtime. Importing freshdata never imports pyspark.

FreshCore is also optional. Install the Python package normally, then build the native extension from the repo checkout:

pip install -e ".[dev,freshcore]"
maturin develop --manifest-path crates/freshcore/Cargo.toml --features extension-module
python benchmarks/bench_freshcore.py --rows 10000 100000 --workload full

If engine="freshcore" is requested but the native module is not installed, FreshData delegates to the pandas reference pipeline and records the reason in report.fallback_events.

Refusing fallbacks: fallback_policy

fd.clean(..., engine="polars", fallback_policy="warn"|"error") turns a silent-but-recorded fallback into a FallbackWarning or a FallbackError raised before any pandas materialization — the strict out-of-core guarantee. Preview the verdict without executing anything via fd.plan(df, engine=...).fallback_reason. Full trigger list: fallback matrix.

Backend support matrix

native = run by the backend itself; fallback = delegated to the pandas reference (output identical, recorded in report.fallback_events); unsupported = not applicable to that engine.

Step (config) pandas polars duckdb spark freshcore
column_names (snake_case rename) native native native native native
strip_whitespace native native native native native
normalize_case (string_case) native fallback fallback fallback native
normalize_sentinels native native native native native
drop_empty_columns / drop_empty_rows native native native native native
drop_duplicates (full-row, keep first/last) native native native native native
impute = mean / median / mode / auto native native native native native
impute_method="missforest" / impute_strategy native fallback fallback fallback fallback
outliers with outlier_method="iqr"/"zscore" (clip/flag) native native native native native
outliers with outlier_method="isolation_forest" native fallback fallback fallback fallback
outliers with outlier_method="auto" (skew-based) native fallback fallback fallback fallback
drop_duplicates with a duplicate_subset native native (order-preserving, eager) fallback fallback fallback
duplicate_keep = drop / aggregate native fallback fallback fallback fallback
fix_dtypes (sampled heuristics) native fallback fallback fallback partial native
drop_constant_columns native fallback fallback fallback fallback
optimize_memory (downcasting) native fallback fallback fallback fallback
Decision engine (strategy="balanced"/"aggressive") native fallback fallback fallback fallback
Missing-indicator columns (missing_indicators) engine-only fallback fallback fallback fallback
output_format pandas pandas/polars/arrow pandas/arrow spark/pandas pandas/polars/arrow

Notes:

  • Imputation counts are exact across backends (the number of filled cells is unambiguous); the fill value for median/mode can differ slightly because each engine uses its own quantile interpolation / tie-breaking. Polars and DuckDB use linear-interpolated quantiles matching pandas; Spark uses approxQuantile. Such divergences are recorded in report.backend_differences.
  • Outlier counts match the pandas reference where the quantile statistics match (Polars/DuckDB linear interpolation); Spark may flag a different count. pandas, Polars, DuckDB and Spark compute the fences from finite values only; ±inf is still tested against them. For a non-constant column whose IQR is zero, pandas and FreshCore fall back to fences from the mean absolute deviation (see the cleaning engine). Polars, DuckDB and Spark still skip such columns, so they can flag fewer outliers there.
  • Float NaN is read as missing by Polars, DuckDB and Spark, as in pandas (Spark keeps NaN as a value in float/double columns, so it is converted to null on ingestion).
  • Spark full-row dedup keeps the first/last occurrence and the surviving rows' order, like pandas. Row order is the input DataFrame's partition order, which is the file order for a single read; a frame that was already shuffled has no stable order to preserve.
  • A non-default pandas index (e.g. a DatetimeIndex) forces a pandas fallback, since native frames carry no index.
  • FreshCore also falls back, based on the data, when outliers is set and a float column holds ±inf, or when impute="mode"/"auto" would fill a nullable boolean column. It falls back for datetime, timedelta, 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, 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).
  • 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.

Arrow interoperability

Arrow Table and RecordBatch are first-class inputs. DuckDB scans Arrow natively (zero-copy) and Polars uses from_arrow, so no pandas materialization happens on the way in. Round-trip Arrow in → Arrow out:

table = fd.clean(arrow_table, engine="duckdb", output_format="arrow")

Command line

freshdata clean input.parquet --engine spark
freshdata clean input.parquet --engine duckdb --memory-limit-gb 4
freshdata clean input.csv --engine polars

Non-pandas --engine values run the scalable path: the file is read by the backend (DuckDB/Polars scan in place; Spark uses its own readers), and the cleaned frame plus a CleanReport summary are emitted. --report report.json writes the full report (including backend, fallback_events, and backend_differences).