From 875ff8ea68bd2f79fa8a107d6db15bfabaf8bcad Mon Sep 17 00:00:00 2001 From: Kevin Costner <120246174+kevincostner17@users.noreply.github.com> Date: Tue, 15 Sep 2026 23:14:53 +0530 Subject: [PATCH] fix(polars): fall back to plain collect when streaming engine is rejected The native Polars backend requested collect(engine="streaming") and caught only TypeError, but polars 1.1-1.24 reject the value with ValueError (Invalid engine argument), which propagated and broke fd.clean on the declared polars floor. Broaden the catch to (TypeError, ValueError) so those versions fall back to a plain collect(); modern polars, which accepts the engine, is unchanged. Found by the production-readiness campaign (FDC-L1b-001). --- CHANGELOG.md | 5 ++++ src/freshdata/execution/backends/_polars.py | 5 +++- .../test_dispatch_and_config.py | 25 +++++++++++++++++++ 3 files changed, 34 insertions(+), 1 deletion(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 6d1e0cda..9f3bfe07 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 c2e308c2..7944d6ff 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 cadf192f..5cf1fd7b 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