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/modecan differ slightly because each engine uses its own quantile interpolation / tie-breaking. Polars and DuckDB use linear-interpolated quantiles matching pandas; Spark usesapproxQuantile. Such divergences are recorded inreport.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;
±infis 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
NaNis read as missing by Polars, DuckDB and Spark, as in pandas (Spark keepsNaNas 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
outliersis set and a float column holds±inf, or whenimpute="mode"/"auto"would fill a nullablebooleancolumn. 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 as1and"1"). Other integer columns are cast back to their input dtype, and non-string labels come back unchanged. Withdrop_duplicates=False, the native module counts duplicate rows at the pandas dedup stage, so detection, theduplicate_thresholdwarning andduplicate_ratio_action="error"match pandas. Native modules built before that count existed fall back underduplicate_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:
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).