Skip to content

Commit a723443

Browse files
fix(engines): disclose ingestion fallbacks; validate engine/output_format pairs (#226)
- #206: a pandas source whose object column mixes value types crashed the Polars engine during pl.from_pandas and was silently cast to text by DuckDB, and duplicate column labels crashed Polars and came back mis-renamed from DuckDB. execution/_ingest.py now flags those inputs before ingestion, and both engines take the recorded pandas fallback (fallback_policy="error" still raises first). The Polars engine now decides fallbacks before converting the source. - #205: requesting another engine's native handle silently returned a different type, even under fallback_policy="error" (e.g. engine="duckdb" with output_format="polars-lazy" returned a DuckDBPyRelation). EngineConfig now rejects the pairing with a ValueError; engine="auto" (and the default engine) picks the engine that owns the handle format. _convert_output returns a materialized frame in place of a handle only when the report records the pandas fallback. Docs: fallback-matrix lists the two input-driven fallbacks; backends.md says a handle format needs its own engine. Closes #205 Closes #206
1 parent b035482 commit a723443

9 files changed

Lines changed: 239 additions & 30 deletions

File tree

‎docs/backends.md‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -36,6 +36,11 @@ memory:
3636
| `"duckdb"` | `DuckDBPyRelation` (un-fetched) | **No** — you call `.fetchdf()`/`.arrow()` |
3737
| `"polars-lazy"` | `pl.LazyFrame` (un-collected) | **No** — you call `.collect()` |
3838

39+
A native handle comes from its own engine: `"duckdb"` needs `engine="duckdb"` and
40+
`"polars-lazy"` needs `engine="polars"`. With `engine="auto"` (or no `engine`),
41+
freshdata picks that engine for you; any other pairing raises `ValueError` rather
42+
than returning a different type.
43+
3944
Neither handle fetches/collects the *result* until you ask, but they are not
4045
equal during the pipeline: the DuckDB path keeps peak memory well below the
4146
eager equivalent, while the Polars pipeline currently collects intermediates

‎docs/fallback-matrix.md‎

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -8,8 +8,9 @@ engine delegates the whole pipeline to pandas** (recorded on
88
`strategy="conservative"` with `fix_dtypes=False`.
99

1010
This table is transcribed from the single source of truth,
11-
`PlanGenerator.fallback_reason()` in `src/freshdata/execution/_plan.py` —
12-
if you change that function, change this page.
11+
`PlanGenerator.fallback_reason()` in `src/freshdata/execution/_plan.py`, plus
12+
the input checks in `pandas_ingest_fallback_reason()`
13+
(`src/freshdata/execution/_ingest.py`) — if you change either, change this page.
1314

1415
| Operation / config | polars | duckdb | spark | freshcore | Why the fallback exists |
1516
|---|---|---|---|---|---|
@@ -28,6 +29,8 @@ if you change that function, change this page.
2829
| `drop_constant_columns` | pandas | pandas | pandas | pandas | needs a data scan before planning (two-phase plan not built) |
2930
| `optimize_memory` | pandas | pandas | pandas | pandas | pandas-specific downcasting — meaningless for other outputs, by design |
3031
| semantic cleaning | native-distinct | native-distinct | pandas | pandas | polars/duckdb run it over a natively extracted distinct table; non-default semantic backends force pandas |
32+
| pandas input with a mixed-type object column (e.g. numbers and strings) | pandas | pandas | — | pandas | native ingestion would reject the column (polars) or cast every value to text (duckdb) |
33+
| pandas input with duplicate column labels | pandas | pandas | — | pandas | native frames need unique column names; the pandas pipeline deduplicates them (`"x", "x"` → `"x", "x_2"`) |
3134
| contracts / validation / memory / profile replay | pandas | pandas | pandas | pandas | in-memory reference features (see [limitations](limitations.md)) |
3235

3336
“pandas” means the **whole pipeline** runs on the pandas reference (fallbacks

‎src/freshdata/api.py‎

Lines changed: 14 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@
1818
from .engine.context import build_contexts
1919
from .engine.model_select import EngineMode, rank_missing_models
2020
from .execution import run_with_engine
21+
from .execution._config import NATIVE_HANDLE_FORMATS
2122
from .parsers.registry import get_parser
2223
from .plan import suggest_plan
2324
from .profile import Profile, build_profile
@@ -55,6 +56,16 @@ def _is_native_engine_source(df: object) -> bool:
5556
return False
5657

5758

59+
def _auto_engine_for(df: object, engine: str, output_format: str) -> str:
60+
"""Resolve the default ``engine="pandas"`` to ``"auto"`` when only a native
61+
engine can serve the input or the requested native handle format."""
62+
if engine == "pandas" and (
63+
_is_native_engine_source(df) or output_format in NATIVE_HANDLE_FORMATS
64+
):
65+
return "auto"
66+
return engine
67+
68+
5869
def _fold_context_options(
5970
options: dict[str, object],
6071
*,
@@ -349,16 +360,14 @@ def clean(
349360
if fallback_policy is not None:
350361
from .execution import EngineConfig as _EngineConfig # noqa: PLC0415
351362

352-
if engine == "pandas" and engine_config is None and not _is_native_engine_source(df):
363+
resolved_engine = _auto_engine_for(df, engine, output_format)
364+
if resolved_engine == "pandas" and engine_config is None:
353365
raise TypeError(
354366
"fallback_policy applies to native engines; engine='pandas' "
355367
"cannot fall back (pass engine='polars'/'duckdb'/... or an "
356368
"engine_config)"
357369
)
358370
if engine_config is None:
359-
resolved_engine = (
360-
"auto" if engine == "pandas" and _is_native_engine_source(df) else engine
361-
)
362371
engine_config = _EngineConfig(
363372
engine=resolved_engine,
364373
output_format=output_format,
@@ -379,7 +388,7 @@ def clean(
379388
df,
380389
config,
381390
options,
382-
engine="auto" if native_source and engine == "pandas" else engine,
391+
engine=_auto_engine_for(df, engine, output_format),
383392
output_format=output_format,
384393
engine_config=engine_config,
385394
return_report=return_report,

‎src/freshdata/execution/__init__.py‎

Lines changed: 29 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@
2020
from ._base import ExecutionEngine
2121
from ._config import (
2222
FALLBACK_POLICIES,
23+
NATIVE_HANDLE_ENGINES,
2324
EngineConfig,
2425
EngineSelector,
2526
FallbackError,
@@ -58,7 +59,18 @@ def _is_spark_frame(frame: Any) -> bool:
5859
return isinstance(frame, SparkDataFrame)
5960

6061

61-
def _convert_output(frame: Any, output_format: str) -> Any:
62+
def _require_recorded_fallback(frame: Any, output_format: str, report: Any) -> None:
63+
"""Allow a materialized frame in place of a native handle only after a
64+
pandas fallback that the report discloses."""
65+
if report is not None and report.fallback_events:
66+
return
67+
raise RuntimeError(
68+
f"output_format={output_format!r} expects a native handle, but the backend "
69+
f"returned {type(frame).__name__} without recording a pandas fallback"
70+
)
71+
72+
73+
def _convert_output(frame: Any, output_format: str, report: Any = None) -> Any:
6274
"""Convert a backend-native frame to the requested output format."""
6375
import pandas as pd
6476

@@ -68,20 +80,20 @@ def _convert_output(frame: Any, output_format: str) -> Any:
6880
# untouched. The backend is responsible for *not* having collected/fetched
6981
# it (see the DuckDB/Polars engines). We never silently materialize here.
7082
#
71-
# A pandas frame at this point means the backend transparently fell back to
72-
# the pandas pipeline (e.g. the balanced decision engine, which only runs on
73-
# pandas). That fallback is already disclosed on the report
74-
# (``fallback_events`` + ``backend="pandas"``), so we return the materialized
75-
# frame rather than raising — the caller can read the report to see why the
76-
# native handle wasn't available.
83+
# EngineConfig only pairs a handle format with the engine that produces it,
84+
# so a pandas frame here means the backend fell back to the pandas pipeline
85+
# (e.g. the balanced decision engine). That is returned as-is only when the
86+
# fallback is recorded on the report (``fallback_events``); anything else is
87+
# a bug, not a silent substitution.
7788
if output_format == "duckdb":
7889
try:
7990
import duckdb
8091
except ImportError: # pragma: no cover - guarded upstream
8192
duckdb = None # type: ignore[assignment]
8293
if duckdb is not None and isinstance(frame, duckdb.DuckDBPyRelation):
8394
return frame
84-
return frame # disclosed pandas fallback
95+
_require_recorded_fallback(frame, output_format, report)
96+
return frame
8597
if output_format == "polars-lazy":
8698
from ._lazy import require_polars
8799

@@ -90,7 +102,8 @@ def _convert_output(frame: Any, output_format: str) -> Any:
90102
return frame
91103
if isinstance(frame, pl.DataFrame):
92104
return frame.lazy()
93-
return frame # disclosed pandas fallback
105+
_require_recorded_fallback(frame, output_format, report)
106+
return frame
94107

95108
if output_format == "spark":
96109
if is_spark:
@@ -157,7 +170,11 @@ def run_with_engine(
157170
requested = engine_config.engine
158171
resolved = engine_config.engine
159172
if resolved == "auto":
160-
resolved = EngineSelector.select(source, engine_config)
173+
# A native handle format can only come from its own engine.
174+
resolved = (
175+
NATIVE_HANDLE_ENGINES.get(engine_config.output_format)
176+
or EngineSelector.select(source, engine_config)
177+
)
161178
engine_config = replace(engine_config, engine=resolved)
162179

163180
# Semantic cleaning is scored on the pandas reference path. On a native
@@ -183,7 +200,7 @@ def run_with_engine(
183200
cleaned, report = run_pipeline(frame, config)
184201
report.backend = "pandas"
185202
report.record_fallback(resolved, "semantic", reason)
186-
result = _convert_output(cleaned, engine_config.output_format)
203+
result = _convert_output(cleaned, engine_config.output_format, report)
187204
_finish_report(report, requested, "pandas", result)
188205
return (result, report) if return_report else result
189206

@@ -195,7 +212,7 @@ def run_with_engine(
195212
from ..semantic.native import run_semantic_native
196213

197214
cleaned_native = run_semantic_native(cleaned_native, config, report, engine=resolved)
198-
result = _convert_output(cleaned_native, engine_config.output_format)
215+
result = _convert_output(cleaned_native, engine_config.output_format, report)
199216
_finish_report(report, requested, resolved, result)
200217
return (result, report) if return_report else result
201218

‎src/freshdata/execution/_config.py‎

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,8 @@
3232
MATERIALIZING_FORMATS = frozenset({"pandas", "polars", "arrow", "spark"})
3333
#: Output formats that hand back a native, lazy/streaming handle instead.
3434
NATIVE_HANDLE_FORMATS = frozenset({"duckdb", "polars-lazy"})
35+
#: The one engine that can produce each native handle format.
36+
NATIVE_HANDLE_ENGINES = {"duckdb": "duckdb", "polars-lazy": "polars"}
3537
#: What to do when a native backend must delegate to the pandas reference:
3638
#: ``"allow"`` (record silently on the report), ``"warn"`` (also emit a
3739
#: :class:`FallbackWarning`), ``"error"`` (raise :class:`FallbackError` before
@@ -117,6 +119,13 @@ def __post_init__(self) -> None:
117119
raise ValueError(
118120
f"output_format must be one of {OUTPUT_FORMATS}, got {self.output_format!r}"
119121
)
122+
handle_engine = NATIVE_HANDLE_ENGINES.get(self.output_format)
123+
if handle_engine is not None and self.engine not in (handle_engine, "auto"):
124+
raise ValueError(
125+
f"output_format={self.output_format!r} returns a native {handle_engine} "
126+
f"handle, which engine={self.engine!r} cannot produce; use "
127+
f"engine={handle_engine!r} or engine='auto'"
128+
)
120129
if self.fallback_policy not in FALLBACK_POLICIES:
121130
raise ValueError(
122131
f"fallback_policy must be one of {FALLBACK_POLICIES}, "

‎src/freshdata/execution/_ingest.py‎

Lines changed: 34 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,34 @@
1+
"""Input-driven reasons a pandas source must take the pandas reference path.
2+
3+
:meth:`~freshdata.execution.PlanGenerator.fallback_reason` looks only at the
4+
config. Some *inputs* cannot enter a native engine faithfully under any config:
5+
Arrow/Polars ingestion rejects duplicate column labels, and an object column
6+
holding mixed value types is either rejected (Polars) or silently cast to text
7+
(DuckDB). Those runs take the disclosed pandas fallback instead, which
8+
``fallback_policy="error"`` still blocks before any pandas work.
9+
"""
10+
11+
from __future__ import annotations
12+
13+
from typing import Any
14+
15+
#: ``infer_dtype`` kinds of an object column that native ingestion cannot keep as-is.
16+
_MIXED_KINDS = frozenset({"mixed", "mixed-integer"})
17+
18+
19+
def pandas_ingest_fallback_reason(source: Any) -> str | None:
20+
"""Return why *source* must be cleaned by the pandas reference, or ``None``."""
21+
import pandas as pd
22+
from pandas.api.types import infer_dtype, is_object_dtype
23+
24+
if not isinstance(source, pd.DataFrame):
25+
return None
26+
if source.columns.duplicated().any():
27+
return "duplicate input column labels require the pandas reference path"
28+
for i, dtype in enumerate(source.dtypes):
29+
if is_object_dtype(dtype) and infer_dtype(source.iloc[:, i], skipna=True) in _MIXED_KINDS:
30+
return (
31+
f"object column {source.columns[i]!r} mixes value types (e.g. numbers "
32+
"and strings), which native ingestion would reject or cast to text"
33+
)
34+
return None

‎src/freshdata/execution/backends/_duckdb.py‎

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@
2626
from ...steps.duplicates import check_duplicate_ratio, report_detected_duplicates
2727
from .._base import ExecutionEngine
2828
from .._config import NATIVE_HANDLE_FORMATS, enforce_fallback_policy
29+
from .._ingest import pandas_ingest_fallback_reason
2930
from .._lazy import has_duckdb, has_polars, require_duckdb
3031
from .._metadata import MetadataScanner, is_duckdb_float
3132
from .._native_steps import (
@@ -111,8 +112,10 @@ def execute(
111112

112113
plan_cols = self._peek_columns(source)
113114
plan = PlanGenerator(config).plan(plan_cols)
114-
if plan.needs_fallback or self._pandas_index_forces_fallback(source):
115-
reason = plan.fallback_reason or "pandas index semantics"
115+
reason = plan.fallback_reason or pandas_ingest_fallback_reason(source)
116+
if reason is None and self._pandas_index_forces_fallback(source):
117+
reason = "pandas index semantics"
118+
if reason is not None:
116119
enforce_fallback_policy(engine_config, "duckdb", "pipeline", reason)
117120
log.warning("freshdata DuckDBEngine: falling back to pandas (%s)", reason)
118121
cleaned, report = self._fallback(source, config)

‎src/freshdata/execution/backends/_polars.py‎

Lines changed: 24 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@
2121
from ...steps.duplicates import check_duplicate_ratio, report_detected_duplicates
2222
from .._base import ExecutionEngine
2323
from .._config import enforce_fallback_policy
24+
from .._ingest import pandas_ingest_fallback_reason
2425
from .._lazy import require_polars
2526
from .._metadata import MetadataScanner
2627
from .._native_steps import (
@@ -111,18 +112,22 @@ def execute(
111112
self._configure_threads(engine_config)
112113
started = time.perf_counter()
113114

115+
# Decide config- and input-driven fallbacks before ingestion: pl.from_pandas
116+
# raises on the inputs pandas_ingest_fallback_reason flags.
117+
reason = (
118+
PlanGenerator(config, backend=self.name).fallback_reason()
119+
or pandas_ingest_fallback_reason(source)
120+
)
121+
if reason is None and self._pandas_index_forces_fallback(source):
122+
reason = "pandas index semantics"
123+
if reason is not None:
124+
return self._delegate_to_pandas(source, config, engine_config, reason)
125+
114126
lf, memory_before = self._to_lazy(source, pl)
115127
names = list(lf.collect_schema().names())
116128
plan = PlanGenerator(config, backend=self.name).plan(names)
117-
118-
if plan.needs_fallback or self._pandas_index_forces_fallback(source):
119-
reason = plan.fallback_reason or "pandas index semantics"
120-
enforce_fallback_policy(engine_config, "polars", "pipeline", reason)
121-
log.warning("freshdata PolarsEngine: falling back to pandas (%s)", reason)
122-
cleaned, report = self._fallback(source, config)
123-
report.backend = "pandas"
124-
report.record_fallback("polars", "pipeline", reason)
125-
return cleaned, report
129+
if plan.fallback_reason is not None:
130+
return self._delegate_to_pandas(source, config, engine_config, plan.fallback_reason)
126131

127132
meta = MetadataScanner.from_polars_lazy(lf)
128133
report = init_report(meta, memory_before)
@@ -145,6 +150,16 @@ def execute(
145150
finalize_report(report, cleaned, started)
146151
return cleaned, report
147152

153+
def _delegate_to_pandas(
154+
self, source: Any, config: CleanConfig, engine_config: EngineConfig, reason: str
155+
) -> tuple[Any, CleanReport]:
156+
enforce_fallback_policy(engine_config, "polars", "pipeline", reason)
157+
log.warning("freshdata PolarsEngine: falling back to pandas (%s)", reason)
158+
cleaned, report = self._fallback(source, config)
159+
report.backend = "pandas"
160+
report.record_fallback("polars", "pipeline", reason)
161+
return cleaned, report
162+
148163
def _fallback(self, source: Any, config: CleanConfig) -> tuple[Any, CleanReport]:
149164
from ...cleaner import run_pipeline
150165

0 commit comments

Comments
 (0)