diff --git a/CHANGELOG.md b/CHANGELOG.md index a159b6c..2d3d4ce 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,11 @@ adheres to [Semantic Versioning](https://semver.org/). ## [Unreleased] ### Fixed +- `fd.clean` on a Spark DataFrame no longer raises + `TypeError: cannot materialize source of type DataFrame`. Under the default + `strategy="balanced"` the pandas fallback now materializes a Spark source + through its `toPandas()` method, alongside the existing polars and DuckDB + paths. - `fd.evaluate_quality_debt` no longer scores a dimension a clean 0.0 when it was never measured. `type_instability` when profiling fails, `pii_risk` when the PII scan is unavailable or fails, and `schema_drift` and `category_churn` diff --git a/src/freshdata/execution/backends/_pandas.py b/src/freshdata/execution/backends/_pandas.py index f9af056..f36da1a 100644 --- a/src/freshdata/execution/backends/_pandas.py +++ b/src/freshdata/execution/backends/_pandas.py @@ -51,6 +51,11 @@ def materialize_to_pandas(source: Any) -> pd.DataFrame: to_df = getattr(source, "df", None) if callable(to_df): return to_df() + # Spark DataFrame (collects to the driver, matching the balanced-strategy + # pandas fallback for a Spark source). + to_pandas_spark = getattr(source, "toPandas", None) + if callable(to_pandas_spark): + return to_pandas_spark() raise TypeError(f"cannot materialize source of type {type(source).__name__}") diff --git a/tests/test_execution/test_dispatch_and_config.py b/tests/test_execution/test_dispatch_and_config.py index db83805..cadf192 100644 --- a/tests/test_execution/test_dispatch_and_config.py +++ b/tests/test_execution/test_dispatch_and_config.py @@ -61,3 +61,20 @@ def test_engine_config_spark_options(): spark_shuffle_partitions=4) assert cfg.spark_shuffle_partitions == 4 assert cfg.output_format == "spark" + + +def test_materialize_to_pandas_accepts_spark_style_frame(): + """A Spark DataFrame exposes ``toPandas`` (not ``to_pandas``/``df``); the + balanced-strategy pandas fallback must materialize it, not raise TypeError.""" + import pandas as pd + + from freshdata.execution.backends._pandas import materialize_to_pandas + + expected = pd.DataFrame({"a": [1, 2, 3], "b": ["x", "y", "z"]}) + + class _FakeSparkFrame: + def toPandas(self) -> pd.DataFrame: # noqa: N802 - matches pyspark's API + return expected.copy() + + out = materialize_to_pandas(_FakeSparkFrame()) + pd.testing.assert_frame_equal(out, expected)