Skip to content

Commit f7ec88a

Browse files
Merge pull request #62 from FreshCode-Org/feat/quality-ops-interop
Quality-Ops interoperability: dbt / Great Expectations / exception tables / lineage exporters
2 parents dc042cc + 9d7fef0 commit f7ec88a

20 files changed

Lines changed: 1678 additions & 4 deletions

‎src/freshdata/__init__.py‎

Lines changed: 19 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -39,6 +39,7 @@
3939
from .config import CleanConfig
4040
from .execution import EngineConfig
4141
from .explain import ExplainReport, explain_clean
42+
from .findings import QualityFinding
4243
from .plan import CleanPlan, ColumnPlan, compare_clean, compare_plans
4344
from .profile import ColumnProfile, Profile
4445
from .report import Action, CleanReport
@@ -59,6 +60,7 @@
5960
"EngineConfig",
6061
"ExplainReport",
6162
"Profile",
63+
"QualityFinding",
6264
"StreamingCleanConfig",
6365
"StreamingCleaner",
6466
"StreamingState",
@@ -135,6 +137,18 @@
135137
"BlockingRule",
136138
})
137139

140+
#: Quality-ops exporters served lazily from :mod:`freshdata.integrations`. Like the
141+
#: enterprise exports they are kept out of ``__all__`` so ``import freshdata`` stays
142+
#: light — importing the integrations layer pulls in the enterprise layer and its
143+
#: optional deps, which should only happen when an exporter is actually used.
144+
_INTEGRATION_EXPORTS = {
145+
"export_quality_ops": "freshdata.integrations.quality_ops",
146+
"QualityOpsResult": "freshdata.integrations.quality_ops",
147+
"export_dbt_tests": "freshdata.integrations.dbt",
148+
"export_gx_suite": "freshdata.integrations.great_expectations",
149+
"build_exception_table": "freshdata.integrations.exceptions",
150+
}
151+
138152

139153
def __getattr__(name: str) -> object:
140154
"""Lazily resolve the ``enterprise`` submodule and its key exports (PEP 562)."""
@@ -146,8 +160,12 @@ def __getattr__(name: str) -> object:
146160
import importlib
147161

148162
return getattr(importlib.import_module("freshdata.enterprise"), name)
163+
if name in _INTEGRATION_EXPORTS:
164+
import importlib
165+
166+
return getattr(importlib.import_module(_INTEGRATION_EXPORTS[name]), name)
149167
raise AttributeError(f"module 'freshdata' has no attribute {name!r}")
150168

151169

152170
def __dir__() -> list:
153-
return sorted([*__all__, "enterprise", *_ENTERPRISE_EXPORTS])
171+
return sorted([*__all__, "enterprise", *_ENTERPRISE_EXPORTS, *_INTEGRATION_EXPORTS])

‎src/freshdata/enterprise/cli.py‎

Lines changed: 59 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@
1515

1616
import argparse
1717
import json
18+
from pathlib import Path
1819
from typing import Any
1920

2021
import pandas as pd
@@ -153,6 +154,42 @@ def cmd_trust(args: argparse.Namespace) -> int:
153154
return 0
154155

155156

157+
def cmd_quality_ops(args: argparse.Namespace) -> int:
158+
from ..findings import findings_from_dict
159+
from ..integrations.quality_ops import export_quality_ops
160+
161+
with open(args.input, encoding="utf-8") as fh:
162+
report_dict = json.load(fh)
163+
findings = findings_from_dict(report_dict)
164+
165+
df = _read_frame(args.data, None) if args.data else None
166+
model_name = args.model_name or Path(args.input).stem
167+
suite_name = args.suite_name or f"{model_name}_suite"
168+
169+
result = export_quality_ops(
170+
findings,
171+
model_name=model_name,
172+
suite_name=suite_name,
173+
dbt_path=args.dbt,
174+
gx_path=args.gx,
175+
exception_table_path=args.exceptions,
176+
df=df,
177+
include_pii=args.include_pii,
178+
exceptions_format=args.exceptions_format,
179+
)
180+
if args.lineage:
181+
with open(args.lineage, "w", encoding="utf-8") as fh:
182+
json.dump(result.lineage_event, fh, indent=2, default=str)
183+
if not args.quiet:
184+
print(f"freshdata quality-ops: {len(result.findings)} finding(s)")
185+
for label, dest in (("dbt", result.dbt_path), ("gx", result.gx_path),
186+
("exceptions", result.exception_table_path),
187+
("lineage", args.lineage)):
188+
if dest:
189+
print(f" {label}: {dest}")
190+
return 0
191+
192+
156193
def build_parser() -> argparse.ArgumentParser:
157194
parser = argparse.ArgumentParser(prog="freshdata", description="freshdata enterprise CLI")
158195
subparsers = parser.add_subparsers(dest="command", required=True)
@@ -190,6 +227,28 @@ def build_parser() -> argparse.ArgumentParser:
190227
help="exit non-zero if the trust score is below this")
191228
trust.set_defaults(func=cmd_trust)
192229

230+
qops = subparsers.add_parser(
231+
"quality-ops",
232+
help="export findings from a report.json to dbt/GX/exception/lineage artifacts",
233+
)
234+
qops.add_argument("input", help="path to a freshdata report JSON (CleanReport.to_dict)")
235+
qops.add_argument("--dbt", metavar="schema.yml", help="write dbt generic tests YAML here")
236+
qops.add_argument("--gx", metavar="suite.json", help="write a Great Expectations suite here")
237+
qops.add_argument("--exceptions", metavar="PATH",
238+
help="write an exception table here (.csv/.parquet/.duckdb)")
239+
qops.add_argument("--exceptions-format", choices=("csv", "parquet", "duckdb"),
240+
help="exception-table format (else inferred from the extension)")
241+
qops.add_argument("--lineage", metavar="lineage.json",
242+
help="write the OpenLineage event (with artifact facets) here")
243+
qops.add_argument("--data", metavar="PATH",
244+
help="optional source file to enrich exception observed_values")
245+
qops.add_argument("--model-name", help="dbt model name (default: report filename stem)")
246+
qops.add_argument("--suite-name", help="GX suite name (default: <model>_suite)")
247+
qops.add_argument("--include-pii", action="store_true",
248+
help="reveal observed values in the exception table (default: redacted)")
249+
qops.add_argument("--quiet", action="store_true")
250+
qops.set_defaults(func=cmd_quality_ops)
251+
193252
from ..streaming._cli import add_stream_subparsers
194253

195254
add_stream_subparsers(subparsers)

‎src/freshdata/enterprise/contracts.py‎

Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -41,6 +41,7 @@
4141
import pandas as pd
4242

4343
from ..adapters.polars import to_pandas
44+
from ..findings import QualityFinding
4445
from .config import (
4546
AnonymizationConfig, # noqa: F401 (re-exported for discoverability)
4647
DriftConfig,
@@ -464,6 +465,34 @@ def to_dict(self) -> dict[str, Any]:
464465
"contract_results": self.contract_results,
465466
}
466467

468+
def to_findings(self, *, lineage_run_id: str | None = None) -> list:
469+
"""Project warned/failed drift findings into :class:`~freshdata.QualityFinding`."""
470+
out: list = []
471+
for f in self.findings:
472+
if f.status == "passed":
473+
continue
474+
expected = None
475+
if f.metric is not None:
476+
expected = str(f.metric)
477+
if f.baseline_value is not None:
478+
expected += f" ~ baseline {f.baseline_value}"
479+
if f.threshold is not None:
480+
expected += f" (threshold {f.threshold})"
481+
out.append(QualityFinding.create(
482+
severity=f.level,
483+
step="drift",
484+
column=f.column,
485+
rule_name=f.check_id,
486+
message=f.message,
487+
observed_value=f.current_value,
488+
expected_condition=expected,
489+
action_taken=f.status,
490+
lineage_run_id=lineage_run_id,
491+
extra={"metric": f.metric, "baseline_value": f.baseline_value,
492+
"threshold": f.threshold, **(f.details or {})},
493+
))
494+
return out
495+
467496
def to_json(self, *, indent: int | None = 2) -> str:
468497
return json.dumps(self.to_dict(), indent=indent, default=str, sort_keys=True)
469498

‎src/freshdata/enterprise/entity_resolution.py‎

Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -228,6 +228,35 @@ def to_dict(self) -> dict[str, Any]:
228228
"runtime_metadata": self.runtime_metadata,
229229
}
230230

231+
def to_findings(self, *, lineage_run_id: str | None = None) -> list:
232+
"""Project match / possible-match pairs into :class:`~freshdata.QualityFinding`.
233+
234+
Each surviving pair is a candidate duplicate: ``match`` maps to ``error``,
235+
``possible_match`` to ``warning``; non-matches are dropped.
236+
"""
237+
from ..findings import QualityFinding
238+
239+
out: list = []
240+
for p in self.pairs:
241+
if p.decision == "non_match":
242+
continue
243+
out.append(QualityFinding.create(
244+
severity=p.decision,
245+
step="entity_resolution",
246+
column=None,
247+
rule_name="duplicate_match",
248+
message=(f"records {p.left_id} & {p.right_id} {p.decision} "
249+
f"(p={p.match_probability:.3f})"),
250+
row_selector=f"{p.left_id} <-> {p.right_id}",
251+
observed_value=p.comparison_vector,
252+
expected_condition="distinct entities",
253+
action_taken=p.decision,
254+
lineage_run_id=lineage_run_id,
255+
extra={"match_probability": round(p.match_probability, 4),
256+
"match_weight": round(p.match_weight, 4)},
257+
))
258+
return out
259+
231260
def summary(self) -> str:
232261
return (
233262
f"entity resolution ({self.backend}): {self.n_records} record(s), "

‎src/freshdata/enterprise/lineage.py‎

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -112,6 +112,13 @@ def __init__(
112112
self.output_name = output_name
113113
self.events: list[LineageEvent] = []
114114
self.started_at = _now_iso()
115+
#: Extra OpenLineage run facets attached by callers (e.g. quality-ops artifact
116+
#: paths). Merged into every emitted event alongside ``freshdata_transformations``.
117+
self.extra_run_facets: dict[str, Any] = {}
118+
119+
def add_run_facet(self, name: str, facet: dict[str, Any]) -> None:
120+
"""Attach a custom OpenLineage run facet (e.g. ``"dbt_tests_path"``)."""
121+
self.extra_run_facets[name] = facet
115122

116123
def record(
117124
self,
@@ -204,7 +211,8 @@ def _event(self, event_type: str, *, with_facets: bool) -> dict[str, Any]:
204211
}
205212
for e in self.events
206213
],
207-
}
214+
},
215+
**self.extra_run_facets,
208216
}
209217
output_facets: dict[str, Any] = {}
210218
input_facets: dict[str, Any] = {}

‎src/freshdata/enterprise/privacy.py‎

Lines changed: 42 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -638,6 +638,48 @@ def to_dict(self) -> dict[str, Any]:
638638
"metadata": self.metadata,
639639
}
640640

641+
def to_findings(self, *, lineage_run_id: str | None = None) -> list:
642+
"""Project masking events (and any k-anonymity breach) into findings.
643+
644+
``observed_value`` is the already-redacted ``original_preview`` and the
645+
findings are flagged ``sensitive``, so PII never leaks through an export
646+
unless the caller explicitly opts in.
647+
"""
648+
from ..findings import QualityFinding
649+
650+
out: list = []
651+
for e in self.events:
652+
out.append(QualityFinding.create(
653+
severity=e.risk_level,
654+
step="privacy",
655+
column=e.column,
656+
rule_name=e.entity_type,
657+
message=f"{e.entity_type} detected via {e.source}",
658+
row_index=e.row,
659+
observed_value=e.original_preview,
660+
expected_condition="no PII",
661+
action_taken=e.strategy,
662+
lineage_run_id=lineage_run_id,
663+
sensitive=True,
664+
extra={"reversible": e.reversible, "format_preserving": e.format_preserving,
665+
"hipaa_tag": e.hipaa_tag, "gdpr_tag": e.gdpr_tag, "score": e.score},
666+
))
667+
ka = self.k_anonymity or {}
668+
if ka and not ka.get("ok", True):
669+
out.append(QualityFinding.create(
670+
severity="error",
671+
step="privacy",
672+
column=None,
673+
rule_name="k_anonymity",
674+
message=(f"k-anonymity violated: {ka.get('rows_violating_k')} row(s) "
675+
f"below k={ka.get('k')}"),
676+
row_selector=f"quasi_identifiers: {list(ka.get('quasi_identifiers', []))}",
677+
expected_condition=f"every group >= k={ka.get('k')}",
678+
lineage_run_id=lineage_run_id,
679+
extra={"violation_ratio": ka.get("violation_ratio")},
680+
))
681+
return out
682+
641683
def to_json(self, *, indent: int | None = 2) -> str:
642684
return json.dumps(self.to_dict(), indent=indent, default=str)
643685

0 commit comments

Comments
 (0)