From 1ada3f9e9c4903d0e4e60a1c7f952e94b26672f5 Mon Sep 17 00:00:00 2001 From: Danila Pechenev Date: Sun, 6 Sep 2026 19:16:21 +0200 Subject: [PATCH 1/2] Add option to persist only invalid rules --- dataframely/filter_result.py | 32 ++++++- docs/guides/features/serialization.md | 11 +++ tests/failure_info/test_parquet.py | 117 ++++++++++++++++++++++++++ 3 files changed, 156 insertions(+), 4 deletions(-) diff --git a/dataframely/filter_result.py b/dataframely/filter_result.py index 48f58d4..e5c21c0 100644 --- a/dataframely/filter_result.py +++ b/dataframely/filter_result.py @@ -128,6 +128,9 @@ def details(self) -> pl.DataFrame: filters in addition to member-level rules, or when calling :meth:`Schema.filter` with `cast=True` and dtype-casting fails for a value. """ + if len(self._rule_columns) == 0: + return self.invalid() + return self._df.select( pl.exclude(self._rule_columns), pl.col(*self._rule_columns).replace_strict( @@ -166,24 +169,45 @@ def __len__(self) -> int: # ---------------------------------- PERSISTENCE --------------------------------- # - def write_parquet(self, file: str | Path | IO[bytes], **kwargs: Any) -> None: + def write_parquet( + self, + file: str | Path | IO[bytes], + *, + only_invalid_rules: bool = False, + **kwargs: Any, + ) -> None: """Write the failure info to a single parquet file. Writes the invalid rows along with additional boolean rule columns indicating which validation rules failed. Unlike :meth:`invalid`, this includes columns - for each rule, where ``False`` indicates the rule failed for that row. + for each rule by default, where ``False`` indicates the rule failed for that row. + + Setting ``only_invalid_rules`` produces a reduced, intentionally lossy + representation that omits rule columns containing only successful or unknown + outcomes. Args: file: The file path or writable file-like object to which to write the parquet file. + only_invalid_rules: Whether to write only rule columns containing at least + one validation failure. kwargs: Additional keyword arguments passed directly to :meth:`polars.write_parquet`. `metadata` may only be provided if it is a dictionary. """ metadata = kwargs.pop("metadata", {}) or {} - self._df.write_parquet( + df = self._df + rule_columns = self._rule_columns + if only_invalid_rules: + counts = _compute_counts(df, self._rule_columns) + rule_columns = [column for column in self._rule_columns if column in counts] + df = df.drop( + column for column in self._rule_columns if column not in counts + ) + + df.write_parquet( file, - metadata={**metadata, "rule_columns": json.dumps(self._rule_columns)}, + metadata={**metadata, "rule_columns": json.dumps(rule_columns)}, **kwargs, ) diff --git a/docs/guides/features/serialization.md b/docs/guides/features/serialization.md index 1f25ee3..c70e59c 100644 --- a/docs/guides/features/serialization.md +++ b/docs/guides/features/serialization.md @@ -65,3 +65,14 @@ failure.write_parquet("failures.parquet") # ...and read it back. failure = dy.FailureInfo.read_parquet("failures.parquet") ``` + +By default, all rule-output columns are persisted. To create a narrower parquet file +for debugging, keep only rule columns that contain at least one validation failure: + +```python +failure.write_parquet("failures.parquet", only_invalid_rules=True) +``` + +All data columns present in the failure information are still written. This reduced +representation is intentionally lossy: rule columns containing only successful or +unknown outcomes are omitted and unavailable after reading the file. diff --git a/tests/failure_info/test_parquet.py b/tests/failure_info/test_parquet.py index a4186ac..8c50f39 100644 --- a/tests/failure_info/test_parquet.py +++ b/tests/failure_info/test_parquet.py @@ -1,6 +1,7 @@ # Copyright (c) QuantCo 2025-2026 # SPDX-License-Identifier: BSD-3-Clause +import json from pathlib import Path import polars as pl @@ -29,6 +30,37 @@ def failure() -> FailureInfo: return failure +@pytest.fixture() +def reducible_failure() -> FailureInfo: + df = pl.DataFrame( + { + "a": [1, 2], + "b": ["foo", "bar"], + "failing_first": [False, True], + "successful": [True, True], + "failing_second": [None, False], + "unknown": [True, None], + } + ) + return FailureInfo( + df.lazy(), + rule_columns=[ + "failing_first", + "successful", + "failing_second", + "unknown", + ], + ) + + +@pytest.fixture() +def empty_failure() -> FailureInfo: + return FailureInfo( + pl.LazyFrame(schema={"a": pl.Int64, "rule": pl.Boolean}), + rule_columns=["rule"], + ) + + @pytest.mark.parametrize("lazy", [True, False]) def test_read_write_parquet(tmp_path: Path, failure: FailureInfo, lazy: bool) -> None: # Arrange @@ -80,3 +112,88 @@ def test_write_parquet_custom_metadata(tmp_path: Path, failure: FailureInfo) -> # The rule columns must still be persisted alongside the custom metadata. read = FailureInfo.read_parquet(path) assert read._rule_columns == failure._rule_columns + + +@pytest.mark.parametrize("lazy", [True, False]) +def test_write_parquet_only_invalid_rules( + tmp_path: Path, reducible_failure: FailureInfo, lazy: bool +) -> None: + # Arrange + path = tmp_path / "failure.parquet" + expected_rule_columns = ["failing_first", "failing_second"] + expected = reducible_failure._df.drop("successful", "unknown") + original = reducible_failure._df.clone() + original_rule_columns = reducible_failure._rule_columns.copy() + + # Act + reducible_failure.write_parquet( + path, + only_invalid_rules=True, + metadata={"custom": "test"}, + ) + read = FailureInfo.scan_parquet(path) if lazy else FailureInfo.read_parquet(path) + + # Assert + assert_frame_equal(pl.read_parquet(path), expected) + metadata = pl.read_parquet_metadata(path) + assert json.loads(metadata["rule_columns"]) == expected_rule_columns + assert metadata["custom"] == "test" + + assert_frame_equal(read._df, expected) + assert read._rule_columns == expected_rule_columns + assert_frame_equal(read.invalid(), reducible_failure.invalid()) + assert read.counts() == reducible_failure.counts() + assert read.cooccurrence_counts() == reducible_failure.cooccurrence_counts() + + assert_frame_equal(reducible_failure._df, original) + assert reducible_failure._rule_columns == original_rule_columns + + +@pytest.mark.parametrize("lazy", [True, False]) +def test_write_parquet_only_invalid_rules_empty( + tmp_path: Path, empty_failure: FailureInfo, lazy: bool +) -> None: + # Arrange + path = tmp_path / "failure.parquet" + expected = pl.DataFrame(schema={"a": pl.Int64}) + + # Act + empty_failure.write_parquet(path, only_invalid_rules=True) + read = FailureInfo.scan_parquet(path) if lazy else FailureInfo.read_parquet(path) + + # Assert + assert_frame_equal(pl.read_parquet(path), expected) + assert json.loads(pl.read_parquet_metadata(path)["rule_columns"]) == [] + assert read._rule_columns == [] + assert_frame_equal(read.invalid(), expected) + assert read.counts() == {} + assert read.cooccurrence_counts() == {} + assert_frame_equal(read.details(), expected) + assert len(read) == 0 + + +@pytest.mark.parametrize("lazy", [True, False]) +def test_write_parquet_only_invalid_rules_repeated_empty( + tmp_path: Path, empty_failure: FailureInfo, lazy: bool +) -> None: + # Arrange + first_path = tmp_path / "first.parquet" + second_path = tmp_path / "second.parquet" + expected = pl.DataFrame(schema={"a": pl.Int64}) + empty_failure.write_parquet(first_path, only_invalid_rules=True) + reduced = FailureInfo.read_parquet(first_path) + + # Act + reduced.write_parquet(second_path, only_invalid_rules=True) + read = ( + FailureInfo.scan_parquet(second_path) + if lazy + else FailureInfo.read_parquet(second_path) + ) + + # Assert + assert_frame_equal(pl.read_parquet(second_path), expected) + assert json.loads(pl.read_parquet_metadata(second_path)["rule_columns"]) == [] + assert read._rule_columns == [] + assert_frame_equal(read.invalid(), expected) + assert read.counts() == {} From 9f36f3df8a5e45cea7b496a0c7f36690e92635c1 Mon Sep 17 00:00:00 2001 From: Danila Pechenev Date: Mon, 7 Sep 2026 17:19:34 +0200 Subject: [PATCH 2/2] Apply reviewer's suggestions --- dataframely/filter_result.py | 10 +++------- docs/guides/features/serialization.md | 2 +- tests/failure_info/test_parquet.py | 14 +++++++------- 3 files changed, 11 insertions(+), 15 deletions(-) diff --git a/dataframely/filter_result.py b/dataframely/filter_result.py index e5c21c0..40b54cc 100644 --- a/dataframely/filter_result.py +++ b/dataframely/filter_result.py @@ -173,7 +173,7 @@ def write_parquet( self, file: str | Path | IO[bytes], *, - only_invalid_rules: bool = False, + only_failing_rules: bool = False, **kwargs: Any, ) -> None: """Write the failure info to a single parquet file. @@ -182,14 +182,10 @@ def write_parquet( which validation rules failed. Unlike :meth:`invalid`, this includes columns for each rule by default, where ``False`` indicates the rule failed for that row. - Setting ``only_invalid_rules`` produces a reduced, intentionally lossy - representation that omits rule columns containing only successful or unknown - outcomes. - Args: file: The file path or writable file-like object to which to write the parquet file. - only_invalid_rules: Whether to write only rule columns containing at least + only_failing_rules: Whether to write only rule columns containing at least one validation failure. kwargs: Additional keyword arguments passed directly to :meth:`polars.write_parquet`. `metadata` may only be provided if it @@ -198,7 +194,7 @@ def write_parquet( metadata = kwargs.pop("metadata", {}) or {} df = self._df rule_columns = self._rule_columns - if only_invalid_rules: + if only_failing_rules: counts = _compute_counts(df, self._rule_columns) rule_columns = [column for column in self._rule_columns if column in counts] df = df.drop( diff --git a/docs/guides/features/serialization.md b/docs/guides/features/serialization.md index c70e59c..ed4f09f 100644 --- a/docs/guides/features/serialization.md +++ b/docs/guides/features/serialization.md @@ -70,7 +70,7 @@ By default, all rule-output columns are persisted. To create a narrower parquet for debugging, keep only rule columns that contain at least one validation failure: ```python -failure.write_parquet("failures.parquet", only_invalid_rules=True) +failure.write_parquet("failures.parquet", only_failing_rules=True) ``` All data columns present in the failure information are still written. This reduced diff --git a/tests/failure_info/test_parquet.py b/tests/failure_info/test_parquet.py index 8c50f39..a6f3605 100644 --- a/tests/failure_info/test_parquet.py +++ b/tests/failure_info/test_parquet.py @@ -115,7 +115,7 @@ def test_write_parquet_custom_metadata(tmp_path: Path, failure: FailureInfo) -> @pytest.mark.parametrize("lazy", [True, False]) -def test_write_parquet_only_invalid_rules( +def test_write_parquet_only_failing_rules( tmp_path: Path, reducible_failure: FailureInfo, lazy: bool ) -> None: # Arrange @@ -128,7 +128,7 @@ def test_write_parquet_only_invalid_rules( # Act reducible_failure.write_parquet( path, - only_invalid_rules=True, + only_failing_rules=True, metadata={"custom": "test"}, ) read = FailureInfo.scan_parquet(path) if lazy else FailureInfo.read_parquet(path) @@ -150,7 +150,7 @@ def test_write_parquet_only_invalid_rules( @pytest.mark.parametrize("lazy", [True, False]) -def test_write_parquet_only_invalid_rules_empty( +def test_write_parquet_only_failing_rules_empty( tmp_path: Path, empty_failure: FailureInfo, lazy: bool ) -> None: # Arrange @@ -158,7 +158,7 @@ def test_write_parquet_only_invalid_rules_empty( expected = pl.DataFrame(schema={"a": pl.Int64}) # Act - empty_failure.write_parquet(path, only_invalid_rules=True) + empty_failure.write_parquet(path, only_failing_rules=True) read = FailureInfo.scan_parquet(path) if lazy else FailureInfo.read_parquet(path) # Assert @@ -173,18 +173,18 @@ def test_write_parquet_only_invalid_rules_empty( @pytest.mark.parametrize("lazy", [True, False]) -def test_write_parquet_only_invalid_rules_repeated_empty( +def test_write_parquet_only_failing_rules_repeated_empty( tmp_path: Path, empty_failure: FailureInfo, lazy: bool ) -> None: # Arrange first_path = tmp_path / "first.parquet" second_path = tmp_path / "second.parquet" expected = pl.DataFrame(schema={"a": pl.Int64}) - empty_failure.write_parquet(first_path, only_invalid_rules=True) + empty_failure.write_parquet(first_path, only_failing_rules=True) reduced = FailureInfo.read_parquet(first_path) # Act - reduced.write_parquet(second_path, only_invalid_rules=True) + reduced.write_parquet(second_path, only_failing_rules=True) read = ( FailureInfo.scan_parquet(second_path) if lazy