Skip to content

Commit e0563da

Browse files
committed
fix: restore run_pipeline signature
1 parent d967239 commit e0563da

1 file changed

Lines changed: 1 addition & 257 deletions

File tree

‎src/freshdata/cleaner.py‎

Lines changed: 1 addition & 257 deletions
Original file line numberDiff line numberDiff line change
@@ -23,260 +23,4 @@
2323
from .steps.prune import drop_constant_columns, drop_empty_columns, drop_empty_rows
2424
from .steps.strings import clean_strings
2525

26-
ProgressCallback = Callable[[dict[str, object]], None]
27-
28-
29-
def _validate_input(df: object, config: CleanConfig) -> pd.DataFrame:
30-
if isinstance(df, pd.Series):
31-
raise TypeError(
32-
"freshdata works on DataFrames; got a Series. Convert it first with s.to_frame()."
33-
)
34-
if not isinstance(df, pd.DataFrame) and not is_polars_frame(df):
35-
raise TypeError(f"expected a pandas or polars DataFrame, got {type(df).__name__}")
36-
frame = to_pandas(df)
37-
if frame.columns.duplicated().any() and not config.column_names:
38-
dupes = sorted({str(c) for c in frame.columns[frame.columns.duplicated()]})
39-
raise ValueError(
40-
f"DataFrame has duplicate column labels {dupes}, which makes "
41-
"column-wise cleaning ambiguous. Rename them, or leave "
42-
"column_names=True to deduplicate automatically."
43-
)
44-
return frame
45-
46-
47-
def _emit_progress(
48-
callback: ProgressCallback | None,
49-
step: str,
50-
status: str,
51-
frame: pd.DataFrame,
52-
) -> None:
53-
if callback is None:
54-
return
55-
callback(
56-
{
57-
"step": step,
58-
"status": status,
59-
"rows": len(frame),
60-
"columns": frame.shape[1],
61-
}
62-
)
63-
64-
65-
def run_pipeline(
66-
df: pd.DataFrame,
67-
def run_pipeline( # noqa: PLR0915 - fixed-order pipeline orchestration
68-
*,
69-
memory: object | None = None,
70-
profile: object | None = None,
71-
) -> tuple[pd.DataFrame, CleanReport]:
72-
"""Run every enabled step, in a fixed and documented order.
73-
74-
With ``preserve_original=True`` (the default) the input frame is never
75-
mutated: the pipeline works on a shallow copy and steps only rebind whole
76-
columns or build new frames, so the only extra memory used is for the
77-
columns that actually change. With ``preserve_original=False`` the
78-
pipeline may write into the input frame to save memory.
79-
80-
After representation repair (names, strings, sentinels, empties, dtypes,
81-
duplicates), ``strategy="auto"`` runs the decision engine for missing
82-
values and outliers; explicit ``impute=`` / ``outliers=`` settings always
83-
override the corresponding engine stage.
84-
85-
``memory`` (a :class:`~freshdata.CleaningMemory`, or ``None``) is passed
86-
through to the semantic stage so it can retrieve and replay compatible
87-
learned semantic repairs; every other step ignores it. ``profile`` (a
88-
:class:`~freshdata.learning.LearningProfile`, or ``None``) is likewise
89-
forwarded so the profile backend can replay learned value maps.
90-
"""
91-
df = _validate_input(df, config)
92-
progress_callback = config.progress_callback
93-
_emit_progress(progress_callback, "input", "after", df)
94-
report = CleanReport(
95-
rows_before=len(df),
96-
cols_before=df.shape[1],
97-
memory_before=memory_bytes(df),
98-
missing_before=int(df.isna().sum().sum()),
99-
)
100-
started = time.perf_counter()
101-
102-
if config.context is not None or config.policy is not None:
103-
# Compile/resolve the context policy against this frame's effective
104-
# schema and lower it into plain config fields. Lazily imported and
105-
# skipped entirely when no context is supplied (zero behaviour change).
106-
from .context import apply_policy_to_config # noqa: PLC0415
107-
108-
config = apply_policy_to_config(config, df=df, report=report)
109-
_emit_progress(progress_callback, "context", "after", df)
110-
111-
out = df.copy(deep=False) if config.preserve_original else df
112-
if config.column_names:
113-
out = normalize_column_names(out, report)
114-
_emit_progress(progress_callback, "column_names", "after", out)
115-
116-
# Hard protected-column guard (context policy / mutable=False): fold the
117-
# protected set into preserve_columns so drop/impute logic honors it, and
118-
# snapshot the columns now (post-rename) to verify byte-identity at the
119-
# end. Zero-cost when no context protection exists.
120-
from .guard import ( # noqa: PLC0415
121-
hard_protected_columns,
122-
snapshot_protected,
123-
verify_protected,
124-
)
125-
126-
hard_protected = hard_protected_columns(config, out.columns)
127-
if hard_protected:
128-
missing_preserve = tuple(c for c in hard_protected if c not in config.preserve_columns)
129-
if missing_preserve:
130-
config = dataclasses.replace(
131-
config, preserve_columns=config.preserve_columns + missing_preserve
132-
)
133-
guard_snapshot = snapshot_protected(out, hard_protected)
134-
else:
135-
guard_snapshot = {}
136-
out = clean_strings(out, config, report)
137-
_emit_progress(progress_callback, "strings", "after", out)
138-
if config.drop_empty_columns:
139-
out = drop_empty_columns(out, report, config)
140-
_emit_progress(progress_callback, "empty_columns", "after", out)
141-
if config.drop_empty_rows:
142-
out = drop_empty_rows(out, report)
143-
_emit_progress(progress_callback, "empty_rows", "after", out)
144-
if config.fix_dtypes:
145-
out = fix_dtypes(out, config, report)
146-
_emit_progress(progress_callback, "dtypes", "after", out)
147-
if config.drop_constant_columns:
148-
out = drop_constant_columns(out, config, report)
149-
_emit_progress(progress_callback, "constant_columns", "after", out)
150-
if config.drop_duplicates:
151-
out = drop_duplicate_rows(out, config, report)
152-
_emit_progress(progress_callback, "duplicates", "after", out)
153-
if config.semantic_enabled:
154-
# Semantic cleaning runs after representation repair and before the
155-
# statistical engine, so missing/outlier logic sees repaired values.
156-
# Lazily imported to keep ``import freshdata`` light.
157-
from .semantic.apply import run_semantic # noqa: PLC0415
158-
159-
out = run_semantic(out, config, report, memory=memory, profile=profile)
160-
_emit_progress(progress_callback, "semantic", "after", out)
161-
if config.engine_mode is not None:
162-
cache = build_engine_cache(out, config)
163-
_emit_progress(progress_callback, "engine_cache", "after", out)
164-
out = auto_missing(
165-
out, config, report, contexts=cache.contexts, numeric_corr=cache.numeric_corr
166-
)
167-
_emit_progress(progress_callback, "engine_missing", "after", out)
168-
out = auto_outliers(out, config, report, contexts=cache.contexts)
169-
_emit_progress(progress_callback, "engine_outliers", "after", out)
170-
out = impute_missing(out, config, report)
171-
_emit_progress(progress_callback, "missing", "after", out)
172-
out = handle_outliers(out, config, report)
173-
_emit_progress(progress_callback, "outliers", "after", out)
174-
out = optimize_memory(out, config, report)
175-
_emit_progress(progress_callback, "memory", "after", out)
176-
if guard_snapshot:
177-
# Physical byte-identity check, before reset_index so row survivors
178-
# can still be aligned by their original index labels.
179-
verify_protected(out, guard_snapshot, report)
180-
_emit_progress(progress_callback, "protected_columns", "after", out)
181-
if config.reset_index:
182-
out = out.reset_index(drop=True)
183-
_emit_progress(progress_callback, "index", "after", out)
184-
185-
report.rows_after = len(out)
186-
report.cols_after = out.shape[1]
187-
report.memory_after = memory_bytes(out)
188-
report.missing_after = int(out.isna().sum().sum())
189-
report.duration_seconds = time.perf_counter() - started
190-
_emit_progress(progress_callback, "complete", "after", out)
191-
return out, report
192-
193-
194-
class Cleaner:
195-
"""A configured, reusable cleaning pipeline.
196-
197-
Useful when the same settings are applied to many frames (e.g. every file
198-
in a directory), or when you want the report after the fact::
199-
200-
cleaner = fd.Cleaner(impute="median", drop_constant_columns=True)
201-
for path in paths:
202-
cleaned = cleaner.clean(pd.read_csv(path))
203-
print(cleaner.report_.summary())
204-
205-
Attributes
206-
----------
207-
config:
208-
The immutable :class:`~freshdata.CleanConfig` in effect.
209-
report_:
210-
The :class:`~freshdata.CleanReport` from the most recent
211-
:meth:`clean` call (``None`` before the first call).
212-
"""
213-
214-
def __init__(
215-
self,
216-
config: CleanConfig | Mapping[str, object] | None = None,
217-
**options: object,
218-
) -> None:
219-
if isinstance(config, Mapping):
220-
merged = dict(config)
221-
merged.update(options)
222-
self._profile = merged.pop("profile", None)
223-
self.config = merge_options(None, **merged)
224-
else:
225-
self._profile = options.pop("profile", None)
226-
self.config = merge_options(config, **options)
227-
self.report_: CleanReport | None = None
228-
229-
def clean(
230-
self,
231-
df: pd.DataFrame,
232-
*,
233-
report: bool = False,
234-
memory: object | None = None,
235-
profile: object | None = None,
236-
) -> pd.DataFrame | tuple[pd.DataFrame, CleanReport]:
237-
"""Clean *df* and return the result (the input is left unchanged
238-
unless ``preserve_original=False`` was configured).
239-
240-
With ``report=True``, returns ``(cleaned_df, CleanReport)`` instead.
241-
The latest report is always available as :attr:`report_`. ``memory``
242-
(a :class:`~freshdata.CleaningMemory`) lets the semantic stage replay
243-
compatible learned repairs; see :func:`freshdata.clean`'s ``memory=``.
244-
``profile`` (a :class:`~freshdata.learning.LearningProfile` or a path
245-
to a ``.fdprofile``) replays a learned profile; it overrides any
246-
profile the ``Cleaner`` was constructed with for this call.
247-
"""
248-
effective_profile = profile if profile is not None else self._profile
249-
gate = None
250-
if effective_profile is not None:
251-
from .learning.replay import ( # noqa: PLC0415 - lazy import
252-
check_profile_drift,
253-
resolve_profile,
254-
)
255-
256-
effective_profile = resolve_profile(effective_profile)
257-
gate = check_profile_drift(to_pandas(df), effective_profile)
258-
259-
cleaned, rep = run_pipeline(
260-
df,
261-
self.config,
262-
memory=memory,
263-
profile=effective_profile if gate is not None and gate.ok else None,
264-
)
265-
if effective_profile is not None and gate is not None:
266-
from .learning.replay import annotate_profile_report # noqa: PLC0415
267-
268-
annotate_profile_report(rep, effective_profile, gate)
269-
self.report_ = rep
270-
if self.config.verbose:
271-
print(rep.brief())
272-
return (cleaned, rep) if report else cleaned
273-
274-
def __repr__(self) -> str:
275-
defaults = CleanConfig()
276-
overrides = {
277-
f.name: getattr(self.config, f.name)
278-
for f in dataclasses.fields(CleanConfig)
279-
if getattr(self.config, f.name) != getattr(defaults, f.name)
280-
}
281-
inner = ", ".join(f"{k}={v!r}" for k, v in overrides.items())
282-
return f"Cleaner({inner})"
26+
ProgressCallback = Callable[[

0 commit comments

Comments
 (0)