diff --git a/CHANGELOG.md b/CHANGELOG.md index 6d1e0cd..9f3bfe0 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,11 @@ adheres to [Semantic Versioning](https://semver.org/). ## [Unreleased] ### Fixed +- The native Polars backend no longer raises on polars versions that reject + `collect(engine="streaming")` with a `ValueError` (polars 1.1–1.24): the + streaming collect now falls back to a plain `collect()` on those versions, + as it already did for the older `streaming=` keyword. Modern polars is + unaffected. - Explicit imputation (`impute="mean"`, `"median"`, `"mode"`, `"auto"`, `"missforest"` or an `impute_strategy` entry) no longer fills the declared `id_columns` or `target_column`, as documented. Those columns keep their diff --git a/src/freshdata/execution/backends/_polars.py b/src/freshdata/execution/backends/_polars.py index c2e308c..7944d6f 100644 --- a/src/freshdata/execution/backends/_polars.py +++ b/src/freshdata/execution/backends/_polars.py @@ -557,7 +557,10 @@ def _collect(self, lf: Any, engine_config: EngineConfig | None, pl: Any) -> Any: # keyword first, then the legacy one, then a plain collect. try: return lf.collect(engine="streaming") - except TypeError: + except (TypeError, ValueError): + # Older polars either lacks the ``engine`` keyword (TypeError) or + # rejects the ``"streaming"`` value (ValueError: Invalid engine + # argument); fall through to the legacy switch, then plain collect. pass try: return lf.collect(streaming=True) diff --git a/tests/test_execution/test_dispatch_and_config.py b/tests/test_execution/test_dispatch_and_config.py index cadf192..5cf1fd7 100644 --- a/tests/test_execution/test_dispatch_and_config.py +++ b/tests/test_execution/test_dispatch_and_config.py @@ -78,3 +78,28 @@ def toPandas(self) -> pd.DataFrame: # noqa: N802 - matches pyspark's API out = materialize_to_pandas(_FakeSparkFrame()) pd.testing.assert_frame_equal(out, expected) + + +def test_polars_collect_falls_back_when_streaming_engine_unsupported(): + """Older polars rejects ``collect(engine="streaming")`` with ValueError (not + TypeError); the backend must fall back to a plain collect, not propagate it.""" + from freshdata.execution._config import EngineConfig + from freshdata.execution.backends._polars import PolarsEngine + + calls: list = [] + + class _FakeLazy: + def collect(self, **kwargs): + calls.append(kwargs) + if kwargs.get("engine") == "streaming": + raise ValueError("Invalid engine argument engine='streaming'") + if kwargs.get("streaming"): + raise TypeError("unexpected keyword argument 'streaming'") + return "collected" + + class _Pl: + pass + + out = PolarsEngine()._collect(_FakeLazy(), EngineConfig(), _Pl()) + assert out == "collected" + assert {} in calls # a plain collect() was reached