From 2fcf6e26324cc4f09e52073cffb79b0c2f5cb9ff Mon Sep 17 00:00:00 2001 From: TeddyCr <13626425+TeddyCr@users.noreply.github.com> Date: Wed, 23 Sep 2026 20:54:31 +0000 Subject: [PATCH 1/4] Fixes #33231: [DQ Threshold W2.4] row-level violation counts for the two between tests `columnValuesToBeBetween` and `columnValueLengthsToBeBetween` decided pass/fail from MIN/MAX aggregates, which cannot answer "how many rows are out of range" -- the one number a row tolerance is checked against. Both test definitions already declare the threshold as "number of failures tolerated ... or a share of the evaluated rows", so the verdict now comes from counting those rows. `BetweenBoundsChecker` already had both halves of the counting, used by the dimensional queries: `build_row_level_violations_sqa()` and `get_violations_mask()`. They are now wired into the overall result too, through a new `_run_violation_count()` on each engine -- one aggregate query for SQA, one vectorized pass per dataframe for pandas -- so the two count the same rows: a NULL is not a violation and an unset bound excludes nothing. The counting only runs when a tolerance is actually configured. With none, one value outside the window and a MIN/MAX outside it are the same verdict, so the result is bit for bit what it was and no extra query is paid for. Since the tolerance is now spent on rows, these two no longer widen their bounds with it: the window is evaluated exactly as configured, and a threshold applied to both the window and the rows would be applied twice. Bounds are still resolved through `get_min_bound`/`get_max_bound`, so dynamic assertion keeps working, and `run_validation()` is untouched. MIN/MAX stay in `testResultValue` -- users read them. When a tolerance decided the verdict the result message leads with the count it was checked against and keeps the window and the extremes as context, and the counts are reported as passed/failed rows since they are already computed. Two fixes the counting needed: an infinite-bound check that only asks `math.isinf` of a float, so a datetime window does not raise, and a COALESCE around the SUM, so an empty table counts zero violations rather than NULL. Co-Authored-By: Claude Opus 5 --- .../checkers/between_bounds_checker.py | 22 +- .../base/columnValueLengthsToBeBetween.py | 92 +++++- .../column/base/columnValuesToBeBetween.py | 101 +++++- .../pandas/columnValueLengthsToBeBetween.py | 30 +- .../column/pandas/columnValuesToBeBetween.py | 30 +- .../columnValueLengthsToBeBetween.py | 19 +- .../sqlalchemy/columnValuesToBeBetween.py | 19 +- .../validations/mixins/sqa_validator_mixin.py | 25 ++ .../validations/result_messages.py | 6 +- .../test_between_row_violations.py | 311 ++++++++++++++++++ 10 files changed, 606 insertions(+), 49 deletions(-) create mode 100644 ingestion/tests/unit/observability/data_quality/validations/test_between_row_violations.py diff --git a/ingestion/src/metadata/data_quality/validations/checkers/between_bounds_checker.py b/ingestion/src/metadata/data_quality/validations/checkers/between_bounds_checker.py index 80c0a71facd8..d44e0374d9e5 100644 --- a/ingestion/src/metadata/data_quality/validations/checkers/between_bounds_checker.py +++ b/ingestion/src/metadata/data_quality/validations/checkers/between_bounds_checker.py @@ -32,6 +32,15 @@ def __init__(self, min_bound: float, max_bound: float): self.min_bound = min_bound self.max_bound = max_bound + def _is_unbounded(self, bound: Any) -> bool: + """Whether that side of the window lets everything through + + An unset bound resolves to ∓inf and needs no condition at all. Only a float can be + infinite: a datetime window -- what a between test on a date column resolves to -- + compares fine and must not be handed to `math.isinf`, which only takes a number. + """ + return isinstance(bound, float) and math.isinf(bound) + def _check_violations(self, values): """Core violation check logic - works for both scalar and Series. @@ -85,9 +94,9 @@ def build_violation_sqa(self, metrics: list["ClauseElement"]) -> "ClauseElement" for expr in metrics: expr_conditions = [] - if not math.isinf(self.min_bound): + if not self._is_unbounded(self.min_bound): expr_conditions.append(and_(expr.isnot(None), expr < self.min_bound)) - if not math.isinf(self.max_bound): + if not self._is_unbounded(self.max_bound): expr_conditions.append(and_(expr.isnot(None), expr > self.max_bound)) if expr_conditions: @@ -115,9 +124,9 @@ def build_row_level_violations_sqa(self, column: "ClauseElement") -> "ClauseElem # Build condition: value NOT NULL AND (value < min OR value > max) conditions = [] - if not math.isinf(self.min_bound): + if not self._is_unbounded(self.min_bound): conditions.append(and_(column.isnot(None), column < self.min_bound)) - if not math.isinf(self.max_bound): + if not self._is_unbounded(self.max_bound): conditions.append(and_(column.isnot(None), column > self.max_bound)) if not conditions: @@ -125,5 +134,6 @@ def build_row_level_violations_sqa(self, column: "ClauseElement") -> "ClauseElem violation_condition = or_(*conditions) if len(conditions) > 1 else conditions[0] - # Return SUM(CASE WHEN violation THEN 1 ELSE 0 END) - return func.sum(case((violation_condition, literal(1)), else_=literal(0))) + # Return SUM(CASE WHEN violation THEN 1 ELSE 0 END). SUM over no row at all is NULL, + # which is not a count: a table with nothing in it has zero violations. + return func.coalesce(func.sum(case((violation_condition, literal(1)), else_=literal(0))), literal(0)) diff --git a/ingestion/src/metadata/data_quality/validations/column/base/columnValueLengthsToBeBetween.py b/ingestion/src/metadata/data_quality/validations/column/base/columnValueLengthsToBeBetween.py index 77fbc6a222c3..3795f0999554 100644 --- a/ingestion/src/metadata/data_quality/validations/column/base/columnValueLengthsToBeBetween.py +++ b/ingestion/src/metadata/data_quality/validations/column/base/columnValueLengthsToBeBetween.py @@ -75,6 +75,11 @@ def _run_validation(self) -> TestCaseResult: Metrics.minLength.name: min_res, } + if self._needs_violation_count(): + total_rows, violating_rows = self._run_violation_count(column, test_params) + metric_values[DIMENSION_TOTAL_COUNT_KEY] = total_rows + metric_values[DIMENSION_FAILED_COUNT_KEY] = violating_rows + except (ValueError, RuntimeError) as exc: msg = f"Error computing {self.test_case.fullyQualifiedName}: {exc}" # type: ignore logger.debug(traceback.format_exc()) @@ -94,7 +99,9 @@ def _run_validation(self) -> TestCaseResult: column, test_params[self.MIN_BOUND], test_params[self.MAX_BOUND] ) else: - row_count, failed_rows = None, None + # A row tolerance already counted both, so report them rather than counting twice. + row_count = metric_values.get(DIMENSION_TOTAL_COUNT_KEY) + failed_rows = metric_values.get(DIMENSION_FAILED_COUNT_KEY) evaluation = self._evaluate_test_condition(metric_values, test_params) result_message = self._format_result_message(metric_values, test_params=test_params) @@ -120,18 +127,27 @@ def _get_validation_checker(self, test_params: dict) -> BetweenBoundsChecker: def _get_test_parameters(self) -> dict: """Get test parameters for this validator + The window is left exactly as the test case configured it. This test reads every row, + so its failure threshold is a row tolerance -- it is spent on how many values may fall + outside the length window, not on widening the window itself. Widening here as well + would apply the same tolerance twice. + Returns: - dict: Test parameters including min and max bounds, widened by the failure threshold + dict: Test parameters including min and max bounds """ - min_bound, max_bound = self.get_bounds(self.MIN_BOUND, self.MAX_BOUND) return { - self.MIN_BOUND: min_bound, - self.MAX_BOUND: max_bound, + self.MIN_BOUND: self.get_min_bound(self.MIN_BOUND), + self.MAX_BOUND: self.get_max_bound(self.MAX_BOUND), } def _get_metrics_to_compute(self, test_params: dict | None = None) -> dict: """Get metrics that need to be computed for this test + The values whose length falls outside the window are not a registry metric -- there is + no aggregate that answers "how many values are too short or too long". They are built + from the bounds by `BetweenBoundsChecker` instead, in `_run_violation_count()` for the + overall result and in `_execute_dimensional_validation()` for each dimension row. + Args: test_params: Optional test parameters (unused for max validator) @@ -143,11 +159,30 @@ def _get_metrics_to_compute(self, test_params: dict | None = None) -> dict: Metrics.minLength.name: Metrics.minLength, } + def _needs_violation_count(self) -> bool: + """Whether the values outside the length window have to be counted + + Only a configured tolerance needs the count: with no tolerance, one value outside the + window and a shortest or longest value outside it are the same verdict, and counting + would cost a query nobody reads. + """ + return bool(self.get_failure_threshold().value) + + def _has_violation_count(self, metric_values: dict) -> bool: + """Whether this result was decided by counting rows rather than by the two extremes + + The dimensional query counts violations for every test case, tolerance or not, so the + count alone does not mean a row tolerance applies. Without one the two agree anyway. + """ + return self._needs_violation_count() and metric_values.get(DIMENSION_FAILED_COUNT_KEY) is not None + def _evaluate_test_condition(self, metric_values: dict, test_params: dict) -> TestEvaluation: """Evaluate the max-to-be-between test condition - For dimensional validation, computes row-level passed/failed counts. - For non-dimensional validation, row counts are not applicable. + Without a tolerance the verdict is read off the two extremes: the shortest and the + longest value inside the window means every length is. That cannot answer "how many + rows are out of range", which is what a row tolerance is checked against, so a test + case that configures one is decided on the counted violations instead. Args: metric_values: Dictionary with keys from Metrics enum names @@ -169,11 +204,15 @@ def _evaluate_test_condition(self, metric_values: dict, test_params: dict) -> Te min_bound = test_params[self.MIN_BOUND] max_bound = test_params[self.MAX_BOUND] - matched = min_bound <= min_length_value and max_length_value <= max_bound - - # Extract row counts if available (dimensional validation) + # Extract row counts if available (row tolerance or dimensional validation) total_rows = metric_values.get(DIMENSION_TOTAL_COUNT_KEY) failed_rows = metric_values.get(DIMENSION_FAILED_COUNT_KEY) + + if self._has_violation_count(metric_values): + matched = self._apply_row_threshold(failed_rows, total_rows) + else: + matched = min_bound <= min_length_value and max_length_value <= max_bound + passed_rows = None if total_rows is not None and failed_rows is not None: passed_rows = total_rows - failed_rows @@ -211,6 +250,22 @@ def _format_result_message( column = self.column_label() + if self._has_violation_count(metric_values): + # A row tolerance is a verdict on rows, so the message leads with the count it was + # checked against and keeps the window and its extremes as context. + counted = self.format_violation_message( + violations=metric_values.get(DIMENSION_FAILED_COUNT_KEY), + population=metric_values.get(DIMENSION_TOTAL_COUNT_KEY), + violation_noun=f"values in {column} outside the expected length", + matched=self._matched(metric_values, test_params), + dimension_info=dimension_info, + ) + return ( + f"{counted} Expected lengths {result_messages.bounds_phrase(min_bound, max_bound)}, " + f"with a shortest value of {result_messages.format_value(min_length_value)} characters " + f"and a longest of {result_messages.format_value(max_length_value)}." + ) + # Both extremes are checked against the same window, so the message reports both and # states the verdict once, on the pair. return self.format_statistic_message( @@ -245,6 +300,23 @@ def _get_test_result_values(self, metric_values: dict) -> list[TestResultValue]: def _run_results(self, metric: Metrics, column: SQALikeColumn | Column): raise NotImplementedError + @abstractmethod + def _run_violation_count(self, column: SQALikeColumn | Column, test_params: dict) -> tuple[int | None, int | None]: + """Count the rows read and the ones whose length falls outside the window + + Both halves are built by `BetweenBoundsChecker` so the SQL and the pandas engines + count the same rows: a NULL has no length and is not a violation, and an unset bound + excludes nothing. + + Args: + column: the column under test + test_params: test parameters including min and max bounds + + Returns: + tuple[int | None, int | None]: rows evaluated, rows outside the window + """ + raise NotImplementedError + @abstractmethod def _execute_dimensional_validation( self, diff --git a/ingestion/src/metadata/data_quality/validations/column/base/columnValuesToBeBetween.py b/ingestion/src/metadata/data_quality/validations/column/base/columnValuesToBeBetween.py index a28b14bf4801..617015e2e11d 100644 --- a/ingestion/src/metadata/data_quality/validations/column/base/columnValuesToBeBetween.py +++ b/ingestion/src/metadata/data_quality/validations/column/base/columnValuesToBeBetween.py @@ -80,6 +80,11 @@ def _run_validation(self) -> TestCaseResult: Metrics.min.name: min_res, Metrics.max.name: max_res, } + + if self._needs_violation_count(): + total_rows, violating_rows = self._run_violation_count(column, test_params) + metric_values[DIMENSION_TOTAL_COUNT_KEY] = total_rows + metric_values[DIMENSION_FAILED_COUNT_KEY] = violating_rows except (ValueError, RuntimeError) as exc: msg = f"Error computing {self.test_case.fullyQualifiedName}: {exc}" # type: ignore logger.debug(traceback.format_exc()) @@ -99,7 +104,9 @@ def _run_validation(self) -> TestCaseResult: column, test_params[self.MIN_BOUND], test_params[self.MAX_BOUND] ) else: - row_count, failed_rows = None, None + # A row tolerance already counted both, so report them rather than counting twice. + row_count = metric_values.get(DIMENSION_TOTAL_COUNT_KEY) + failed_rows = metric_values.get(DIMENSION_FAILED_COUNT_KEY) evaluation = self._evaluate_test_condition(metric_values, test_params) result_message = self._format_result_message(metric_values, test_params=test_params) @@ -123,9 +130,10 @@ def _run_validation(self) -> TestCaseResult: def _get_test_parameters(self) -> dict: """Get Test Parameters - A datetime window is left as the test case configured it: the failure threshold is a - number, and there is no meaningful way to widen a date by one. The result message says - as much rather than claiming a tolerance that never applied. + The window is left exactly as the test case configured it. This test reads every row, + so its failure threshold is a row tolerance -- it is spent on how many values may fall + outside the window, not on widening the window itself. Widening here as well would + apply the same tolerance twice. """ column = self.get_column() @@ -146,7 +154,8 @@ def _get_test_parameters(self) -> dict: pre_processor=convert_timestamp, ) else: - min_bound, max_bound = self.get_bounds(self.MIN_BOUND, self.MAX_BOUND) + min_bound = self.get_min_bound(self.MIN_BOUND) + max_bound = self.get_max_bound(self.MAX_BOUND) return { self.MIN_BOUND: min_bound, @@ -154,27 +163,54 @@ def _get_test_parameters(self) -> dict: } def _get_metrics_to_compute(self, test_params: dict | None = None) -> dict: - """Get Metrics needed to compute""" + """Get Metrics needed to compute + + The out-of-range rows a tolerance is counted against are not a registry metric -- there + is no aggregate that answers "how many values fall outside this window". They are built + from the bounds by `BetweenBoundsChecker` instead, in `_run_violation_count()` for the + overall result and in `_execute_dimensional_validation()` for each dimension row. + """ return {Metrics.min.name: Metrics.min, Metrics.max.name: Metrics.max} + def _needs_violation_count(self) -> bool: + """Whether the rows outside the window have to be counted + + Only a configured tolerance needs the count: with no tolerance, one value outside the + window and a MIN/MAX outside it are the same verdict, and counting would cost a query + nobody reads. + """ + return bool(self.get_failure_threshold().value) + + def _has_violation_count(self, metric_values: dict) -> bool: + """Whether this result was decided by counting rows rather than by MIN/MAX + + The dimensional query counts violations for every test case, tolerance or not, so the + count alone does not mean a row tolerance applies. Without one the two agree anyway. + """ + return self._needs_violation_count() and metric_values.get(DIMENSION_FAILED_COUNT_KEY) is not None + def _evaluate_test_condition(self, metric_values: dict, test_params: dict) -> TestEvaluation: """Evaluate the values-to-be-between test condition - For this test, the condition passes if both min and max values are within bounds. - Since this is a statistical validator (group-level), passed/failed row counts - are not applicable at the test level (only for computePassedFailedRowCount). + Without a tolerance the verdict is read off MIN/MAX: both extremes inside the window + means every value is. That cannot answer "how many rows are out of range", which is + what a row tolerance is checked against, so a test case that configures one is decided + on the counted violations instead. Args: metric_values: Dictionary with keys from Metrics enum names e.g., {"MIN": 10, "MAX": 100} + With a row tolerance, and for every dimension row, also includes: + - DIMENSION_TOTAL_COUNT_KEY: rows evaluated + - DIMENSION_FAILED_COUNT_KEY: rows outside the window test_params: Dictionary with 'minValue' and 'maxValue' Returns: dict with keys: - - matched: bool - whether test passed (both min >= min_bound and max <= max_bound) - - passed_rows: None - not applicable for statistical validators - - failed_rows: None - not applicable for statistical validators - - total_rows: None - not applicable for statistical validators + - matched: bool - whether test passed + - passed_rows: Optional[int] - rows inside the window, when counted + - failed_rows: Optional[int] - rows outside the window, when counted + - total_rows: Optional[int] - rows evaluated, when counted """ min_value = metric_values[Metrics.min.name] @@ -182,9 +218,14 @@ def _evaluate_test_condition(self, metric_values: dict, test_params: dict) -> Te min_bound = test_params[self.MIN_BOUND] max_bound = test_params[self.MAX_BOUND] - matched = min_value >= min_bound and max_value <= max_bound total_rows = metric_values.get(DIMENSION_TOTAL_COUNT_KEY) failed_rows = metric_values.get(DIMENSION_FAILED_COUNT_KEY) + + if self._has_violation_count(metric_values): + matched = self._apply_row_threshold(failed_rows, total_rows) + else: + matched = min_value >= min_bound and max_value <= max_bound + passed_rows = total_rows - failed_rows if (total_rows is not None and failed_rows is not None) else None return { @@ -221,6 +262,22 @@ def _format_result_message( matched = self._matched(metric_values, test_params) column = self.column_label() + if self._has_violation_count(metric_values): + # A row tolerance is a verdict on rows, so the message leads with the count it was + # checked against and keeps the window and its extremes as context. + counted = self.format_violation_message( + violations=metric_values.get(DIMENSION_FAILED_COUNT_KEY), + population=metric_values.get(DIMENSION_TOTAL_COUNT_KEY), + violation_noun=f"values in {column} outside the expected range", + matched=matched, + dimension_info=dimension_info, + ) + return ( + f"{counted} Expected {result_messages.bounds_phrase(min_bound, max_bound)}, " + f"with a minimum of {result_messages.format_value(min_value)} and a maximum of " + f"{result_messages.format_value(max_value)}." + ) + # Both extremes are checked against the same window, so the message reports both and # states the verdict once, on the pair. return self.format_statistic_message( @@ -286,6 +343,22 @@ def _execute_dimensional_validation( def _run_results(self, metric: Metrics, column: SQALikeColumn | Column): raise NotImplementedError + @abstractmethod + def _run_violation_count(self, column: SQALikeColumn | Column, test_params: dict) -> tuple[int | None, int | None]: + """Count the rows read and the ones whose value falls outside the window + + Both halves are built by `BetweenBoundsChecker` so the SQL and the pandas engines + count the same rows: a NULL is not a violation, and an unset bound excludes nothing. + + Args: + column: the column under test + test_params: test parameters including min and max bounds + + Returns: + tuple[int | None, int | None]: rows evaluated, rows outside the window + """ + raise NotImplementedError + @abstractmethod def compute_row_count(self, column: SQALikeColumn | Column, min_bound, max_bound): """Compute row count for the given column diff --git a/ingestion/src/metadata/data_quality/validations/column/pandas/columnValueLengthsToBeBetween.py b/ingestion/src/metadata/data_quality/validations/column/pandas/columnValueLengthsToBeBetween.py index 5959e9742508..7e17a34e8602 100644 --- a/ingestion/src/metadata/data_quality/validations/column/pandas/columnValueLengthsToBeBetween.py +++ b/ingestion/src/metadata/data_quality/validations/column/pandas/columnValueLengthsToBeBetween.py @@ -63,6 +63,27 @@ def _run_results(self, metric: Metrics, column: SQALikeColumn) -> int | None: """ return self.run_dataframe_results(self.runner, metric, column) + def _run_violation_count(self, column: SQALikeColumn, test_params: dict) -> tuple[int, int]: + """Count the rows read and the values whose length falls outside the window + + The dataframes are walked one at a time rather than concatenated, like every other + aggregate this validator computes: a dataset split across many files does not have to + fit in memory to be counted. + + Args: + column: column under test + test_params: test parameters including min and max bounds + """ + checker = self._get_validation_checker(test_params) + + total_rows = 0 + violating_rows = 0 + for df in self.runner: + total_rows += len(df) + violating_rows += int(checker.get_violations_mask(df[column.name].str.len()).sum()) + + return total_rows, violating_rows + def _build_dimension_metric_values(self, row, metrics_to_compute, test_params=None): metric_values = self._build_metric_values_from_row(row, metrics_to_compute, test_params) metric_values[DIMENSION_TOTAL_COUNT_KEY] = row.get(DIMENSION_TOTAL_COUNT_KEY) @@ -229,10 +250,11 @@ def compute_row_count(self, column: SQALikeColumn, min_bound: int, max_bound: in return row_count, failed_rows def filter(self): - # The verdict is taken against the length window the failure threshold widened into, so the - # failed rows are filtered with it too: a value the tolerance accepted is not a failure and - # has no business showing up in the sample. - min_bound, max_bound = self.get_bounds(self.MIN_BOUND, self.MAX_BOUND) + # The window is the one the test case configured: the failure threshold is a row tolerance + # here, and a row it tolerates is still a value whose length fell outside the window, so it + # belongs in the sample of failing rows. + min_bound = self.get_min_bound(self.MIN_BOUND) + max_bound = self.get_max_bound(self.MAX_BOUND) filters = [] if min_bound is not None and min_bound > float("-inf"): filters.append(f"{self.get_column().name}.astype('str').str.len() < {min_bound}") diff --git a/ingestion/src/metadata/data_quality/validations/column/pandas/columnValuesToBeBetween.py b/ingestion/src/metadata/data_quality/validations/column/pandas/columnValuesToBeBetween.py index 4c52083511c1..5dc14d60de01 100644 --- a/ingestion/src/metadata/data_quality/validations/column/pandas/columnValuesToBeBetween.py +++ b/ingestion/src/metadata/data_quality/validations/column/pandas/columnValuesToBeBetween.py @@ -66,6 +66,27 @@ def _run_results(self, metric: Metrics, column: SQALikeColumn) -> int | None: """ return self.run_dataframe_results(self.runner, metric, column) + def _run_violation_count(self, column: SQALikeColumn, test_params: dict) -> tuple[int, int]: + """Count the rows read and the values falling outside the window + + The dataframes are walked one at a time rather than concatenated, like every other + aggregate this validator computes: a dataset split across many files does not have to + fit in memory to be counted. + + Args: + column: column under test + test_params: test parameters including min and max bounds + """ + checker = self._get_validation_checker(test_params) + + total_rows = 0 + violating_rows = 0 + for df in self.runner: + total_rows += len(df) + violating_rows += int(checker.get_violations_mask(df[column.name]).sum()) + + return total_rows, violating_rows + def _build_dimension_metric_values(self, row, metrics_to_compute, test_params=None): metric_values = self._build_metric_values_from_row(row, metrics_to_compute, test_params) metric_values[DIMENSION_TOTAL_COUNT_KEY] = row.get(DIMENSION_TOTAL_COUNT_KEY) @@ -253,10 +274,11 @@ def filter(self): pre_processor=convert_timestamp, ) else: - # The verdict is taken against the window the failure threshold widened into, so the - # failed rows are filtered with it too: a value the tolerance accepted is not a failure - # and has no business showing up in the sample. - min_bound, max_bound = self.get_bounds(self.MIN_BOUND, self.MAX_BOUND) + # The window is the one the test case configured: the failure threshold is a row + # tolerance here, and a row it tolerates is still a row that fell outside the window, + # so it belongs in the sample of failing rows. + min_bound = self.get_min_bound(self.MIN_BOUND) + max_bound = self.get_max_bound(self.MAX_BOUND) filters = [] if min_bound is not None: filters.append(f"{column.name} < {min_bound}") diff --git a/ingestion/src/metadata/data_quality/validations/column/sqlalchemy/columnValueLengthsToBeBetween.py b/ingestion/src/metadata/data_quality/validations/column/sqlalchemy/columnValueLengthsToBeBetween.py index 615e7151c772..e88eb52ee370 100644 --- a/ingestion/src/metadata/data_quality/validations/column/sqlalchemy/columnValueLengthsToBeBetween.py +++ b/ingestion/src/metadata/data_quality/validations/column/sqlalchemy/columnValueLengthsToBeBetween.py @@ -59,6 +59,16 @@ def _run_results(self, metric: Metrics, column: Column) -> int | None: """ return self.run_query_results(self.runner, metric, column) + def _run_violation_count(self, column: Column, test_params: dict) -> tuple[int | None, int | None]: + """Count the rows read and the values whose length falls outside the window + + Args: + column: column under test + test_params: test parameters including min and max bounds + """ + checker = self._get_validation_checker(test_params) + return self._compute_row_violations(self.runner, checker.build_row_level_violations_sqa(LenFn(column))) + def compute_row_count(self, column: Column, min_bound: int, max_bound: int): """Compute row count for the given column @@ -153,10 +163,11 @@ def _execute_dimensional_validation( return dimension_results def filter(self): - # The verdict is taken against the length window the failure threshold widened into, so the - # failed rows are filtered with it too: a value the tolerance accepted is not a failure and - # has no business showing up in the sample. - min_bound, max_bound = self.get_bounds(self.MIN_BOUND, self.MAX_BOUND) + # The window is the one the test case configured: the failure threshold is a row tolerance + # here, and a row it tolerates is still a value whose length fell outside the window, so it + # belongs in the sample of failing rows. + min_bound = self.get_min_bound(self.MIN_BOUND) + max_bound = self.get_max_bound(self.MAX_BOUND) filters = [] if min_bound is not None and min_bound > float("-inf"): filters.append((LenFn(self.get_column()), "lt", min_bound)) diff --git a/ingestion/src/metadata/data_quality/validations/column/sqlalchemy/columnValuesToBeBetween.py b/ingestion/src/metadata/data_quality/validations/column/sqlalchemy/columnValuesToBeBetween.py index d6b8942cb367..779ac64f9987 100644 --- a/ingestion/src/metadata/data_quality/validations/column/sqlalchemy/columnValuesToBeBetween.py +++ b/ingestion/src/metadata/data_quality/validations/column/sqlalchemy/columnValuesToBeBetween.py @@ -61,6 +61,16 @@ def _run_results(self, metric: Metrics, column: Column) -> int | None: """ return self.run_query_results(self.runner, metric, column) + def _run_violation_count(self, column: Column, test_params: dict) -> tuple[int | None, int | None]: + """Count the rows read and the values falling outside the window + + Args: + column: column under test + test_params: test parameters including min and max bounds + """ + checker = self._get_validation_checker(test_params) + return self._compute_row_violations(self.runner, checker.build_row_level_violations_sqa(column)) + def _build_dimension_metric_values(self, row, metrics_to_compute, test_params=None): min_value = row.get(Metrics.min.name) max_value = row.get(Metrics.max.name) @@ -172,10 +182,11 @@ def filter(self): pre_processor=convert_timestamp, ) else: - # The verdict is taken against the window the failure threshold widened into, so the - # failed rows are filtered with it too: a value the tolerance accepted is not a failure - # and has no business showing up in the sample. - min_bound, max_bound = self.get_bounds(self.MIN_BOUND, self.MAX_BOUND) + # The window is the one the test case configured: the failure threshold is a row + # tolerance here, and a row it tolerates is still a row that fell outside the window, + # so it belongs in the sample of failing rows. + min_bound = self.get_min_bound(self.MIN_BOUND) + max_bound = self.get_max_bound(self.MAX_BOUND) filters = [] if min_bound is not None: diff --git a/ingestion/src/metadata/data_quality/validations/mixins/sqa_validator_mixin.py b/ingestion/src/metadata/data_quality/validations/mixins/sqa_validator_mixin.py index 5e5064001acb..77dcd5fabbb8 100644 --- a/ingestion/src/metadata/data_quality/validations/mixins/sqa_validator_mixin.py +++ b/ingestion/src/metadata/data_quality/validations/mixins/sqa_validator_mixin.py @@ -185,6 +185,31 @@ def _compute_row_count_between( return res + def _compute_row_violations(self, runner: QueryRunner, violations_expr: ClauseElement) -> tuple[Any, Any]: + """Count the rows read and the violating ones among them, in a single aggregate query + + Both counts come from the same scan so they are always counted against each other: a + violation count read from one query and a population from another can disagree on a + table that changed in between. + + Args: + runner: runner with the sqlalchemy session + violations_expr: aggregate expression counting the violating rows, built by a checker + + Returns: + tuple[Any, Any]: rows evaluated, violating rows + """ + try: + row = runner.dispatch_query_select_first( + Metrics.rowCount().fn().label(DIMENSION_TOTAL_COUNT_KEY), + violations_expr.label(DIMENSION_FAILED_COUNT_KEY), + ) + values = dict(row._mapping) + except Exception as exc: + raise SQLAlchemyError(exc) # noqa: B904 + + return values.get(DIMENSION_TOTAL_COUNT_KEY), values.get(DIMENSION_FAILED_COUNT_KEY) + def _compute_row_count(self, runner: QueryRunner, column: Column, **kwargs): """compute row count diff --git a/ingestion/src/metadata/data_quality/validations/result_messages.py b/ingestion/src/metadata/data_quality/validations/result_messages.py index e0122dcceb05..2a663fdd404a 100644 --- a/ingestion/src/metadata/data_quality/validations/result_messages.py +++ b/ingestion/src/metadata/data_quality/validations/result_messages.py @@ -156,7 +156,7 @@ def violation_sentence( ) -def _bounds_phrase(min_bound: float | None, max_bound: float | None, lead: bool = True) -> str: +def bounds_phrase(min_bound: float | None, max_bound: float | None, lead: bool = True) -> str: """ "between 90 and 110", "at least 90", "at most 110" -- whichever bounds are set `lead` drops the "between" so the phrase can follow one that already introduced the window, @@ -196,13 +196,13 @@ def statistic_sentence( str: e.g. "Mean of `amount` is 87.4. Expected between 90 and 110, widened by a 5% tolerance to 85.5 and 115.5, so this test passed." """ - expected = f"Expected {_bounds_phrase(*configured_bounds)}" + expected = f"Expected {bounds_phrase(*configured_bounds)}" if not threshold.value: tolerance = ", with no tolerance applied" elif effective_bounds != configured_bounds: tolerance = ( - f", widened by a {format_threshold(threshold)} tolerance to {_bounds_phrase(*effective_bounds, lead=False)}" + f", widened by a {format_threshold(threshold)} tolerance to {bounds_phrase(*effective_bounds, lead=False)}" ) else: # A threshold that widened nothing: no finite bound to widen, or bounds a numeric diff --git a/ingestion/tests/unit/observability/data_quality/validations/test_between_row_violations.py b/ingestion/tests/unit/observability/data_quality/validations/test_between_row_violations.py new file mode 100644 index 000000000000..0ebb8d63e712 --- /dev/null +++ b/ingestion/tests/unit/observability/data_quality/validations/test_between_row_violations.py @@ -0,0 +1,311 @@ +# Copyright 2025 Collate +# Licensed under the Collate Community License, Version 1.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# https://github.com/open-metadata/OpenMetadata/blob/main/ingestion/LICENSE +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. +""" +Row-level violation counts for the two between tests. + +`columnValuesToBeBetween` and `columnValueLengthsToBeBetween` read every row, so their failure +threshold is a row tolerance: it is checked against how many values fall outside the window, which +the MIN/MAX these tests report cannot answer. The counting is built by `BetweenBoundsChecker`, so +the SQL and the pandas engines have to agree on the same fixture. +""" + +import os +from datetime import datetime +from unittest.mock import patch +from uuid import uuid4 + +import pytest +import sqlalchemy as sqa +from pandas import DataFrame +from sqlalchemy.orm import DeclarativeBase + +from metadata.data_quality.builders.validator_builder import ValidatorBuilder +from metadata.data_quality.interface.sqlalchemy.sqa_test_suite_interface import ( + SQATestSuiteInterface, +) +from metadata.data_quality.validations.column.pandas.columnValueLengthsToBeBetween import ( + ColumnValueLengthsToBeBetweenValidator as PandasLengthsValidator, +) +from metadata.data_quality.validations.column.pandas.columnValuesToBeBetween import ( + ColumnValuesToBeBetweenValidator as PandasValuesValidator, +) +from metadata.data_quality.validations.column.sqlalchemy.columnValueLengthsToBeBetween import ( + ColumnValueLengthsToBeBetweenValidator as SQALengthsValidator, +) +from metadata.data_quality.validations.column.sqlalchemy.columnValuesToBeBetween import ( + ColumnValuesToBeBetweenValidator as SQAValuesValidator, +) +from metadata.generated.schema.entity.data.table import Column, DataType, Table +from metadata.generated.schema.entity.services.connections.database.sqliteConnection import ( + SQLiteConnection, + SQLiteScheme, +) +from metadata.generated.schema.tests.basic import TestCaseStatus +from metadata.generated.schema.tests.testCase import TestCase, TestCaseParameterValue +from metadata.generated.schema.type.entityReference import EntityReference +from metadata.profiler.processor.runner import PandasRunner +from metadata.sampler.sqlalchemy.sampler import SQASampler + +EXECUTION_DATE = datetime.strptime("2021-07-03", "%Y-%m-%d") + +ENTITY_LINK_VALUE = "<#E::table::service.db.measurements::columns::value>" +ENTITY_LINK_LABEL = "<#E::table::service.db.measurements::columns::label>" + +# Both tests are run against the same window, [3, 8], so the fixture is hand countable once: +# `value` and `label`'s length are the same seven numbers. Three of them -- 1, 9 and 12 -- fall +# outside the window, and the NULL row violates nothing: it has no value to compare and no length. +ROWS = [ + (1, "a"), + (3, "abc"), + (5, "abcde"), + (8, "abcdefgh"), + (9, "abcdefghi"), + (12, "abcdefghijkl"), + (None, None), +] + +MIN_BOUND = 3 +MAX_BOUND = 8 +EXPECTED_VIOLATIONS = 3 +EXPECTED_ROWS = len(ROWS) + + +class Base(DeclarativeBase): + pass + + +class Measurement(Base): + __tablename__ = "measurements" + id = sqa.Column(sqa.Integer, primary_key=True) + value = sqa.Column(sqa.Integer) + label = sqa.Column(sqa.String(64)) + + +TABLE = Table( + id=uuid4(), + name="measurements", + fullyQualifiedName="service.db.measurements", + columns=[ + Column(name="id", dataType=DataType.INT), # type: ignore + Column(name="value", dataType=DataType.INT), # type: ignore + Column(name="label", dataType=DataType.STRING), # type: ignore + ], + database=EntityReference(id=uuid4(), name="db", type="database"), # type: ignore +) # type: ignore + + +@pytest.fixture +def sqa_runner(worker_id): + """A sqlite table holding the fixture, and the runner reading it""" + worker_suffix = f"_{worker_id}" if worker_id != "master" else "" + db_path = os.path.join( # noqa: PTH118 + os.path.dirname(__file__), # noqa: PTH120 + f"{os.path.splitext(os.path.basename(__file__))[0]}{worker_suffix}.db", # noqa: PTH119, PTH122 + ) + sqlite_conn = SQLiteConnection( + scheme=SQLiteScheme.sqlite_pysqlite, + databaseMode=db_path + "?check_same_thread=False", + ) # type: ignore + + with patch.object(SQASampler, "build_table_orm", return_value=Measurement): + sampler = SQASampler( + service_connection_config=sqlite_conn, + ometa_client=None, + entity=TABLE, + ) + interface = SQATestSuiteInterface( + sqlite_conn, + None, + sampler, + TABLE, + validator_builder=ValidatorBuilder, + ) + + engine = interface.session.get_bind() + Measurement.__table__.create(bind=engine) + interface.session.add_all([Measurement(value=value, label=label) for value, label in ROWS]) + interface.session.commit() + + yield interface.runner + + Measurement.__table__.drop(bind=engine) + if os.path.exists(db_path): # noqa: PTH110 + os.remove(db_path) # noqa: PTH107 + + +@pytest.fixture +def pandas_runner(): + """The same fixture, split over two dataframes so the counting has to accumulate""" + frames = ( + DataFrame(ROWS[:4], columns=["value", "label"]), + DataFrame(ROWS[4:], columns=["value", "label"]), + ) + return PandasRunner(dataset=lambda: iter(frames), raw_dataset=None) + + +def values_test_case(threshold=None, unit=None) -> TestCase: + """A values-to-be-between test case on the fixture window""" + return _test_case( + ENTITY_LINK_VALUE, + [ + TestCaseParameterValue(name="minValue", value=str(MIN_BOUND)), + TestCaseParameterValue(name="maxValue", value=str(MAX_BOUND)), + ], + threshold, + unit, + ) + + +def lengths_test_case(threshold=None, unit=None) -> TestCase: + """A lengths-to-be-between test case on the fixture window""" + return _test_case( + ENTITY_LINK_LABEL, + [ + TestCaseParameterValue(name="minLength", value=str(MIN_BOUND)), + TestCaseParameterValue(name="maxLength", value=str(MAX_BOUND)), + ], + threshold, + unit, + ) + + +def _test_case(entity_link, parameter_values, threshold, unit) -> TestCase: + if threshold is not None: + parameter_values = [*parameter_values, TestCaseParameterValue(name="threshold", value=str(threshold))] + if unit is not None: + parameter_values = [*parameter_values, TestCaseParameterValue(name="thresholdUnit", value=unit)] + return TestCase( + name="my_test_case", + entityLink=entity_link, + testSuite=EntityReference(id=uuid4(), type="TestSuite"), # type: ignore + testDefinition=EntityReference(id=uuid4(), type="TestDefinition"), # type: ignore + parameterValues=parameter_values, + ) # type: ignore + + +def test_sqa_counts_the_rows_outside_the_value_window(sqa_runner): + validator = SQAValuesValidator(sqa_runner, values_test_case(threshold=1), EXECUTION_DATE.timestamp()) + column = validator.get_column() + + assert validator._run_violation_count(column, validator._get_test_parameters()) == ( + EXPECTED_ROWS, + EXPECTED_VIOLATIONS, + ) + + +def test_pandas_counts_the_rows_outside_the_value_window(pandas_runner): + validator = PandasValuesValidator(pandas_runner, values_test_case(threshold=1), EXECUTION_DATE.timestamp()) + column = validator.get_column() + + assert validator._run_violation_count(column, validator._get_test_parameters()) == ( + EXPECTED_ROWS, + EXPECTED_VIOLATIONS, + ) + + +def test_sqa_counts_the_rows_outside_the_length_window(sqa_runner): + validator = SQALengthsValidator(sqa_runner, lengths_test_case(threshold=1), EXECUTION_DATE.timestamp()) + column = validator.get_column() + + assert validator._run_violation_count(column, validator._get_test_parameters()) == ( + EXPECTED_ROWS, + EXPECTED_VIOLATIONS, + ) + + +def test_pandas_counts_the_rows_outside_the_length_window(pandas_runner): + validator = PandasLengthsValidator(pandas_runner, lengths_test_case(threshold=1), EXECUTION_DATE.timestamp()) + column = validator.get_column() + + assert validator._run_violation_count(column, validator._get_test_parameters()) == ( + EXPECTED_ROWS, + EXPECTED_VIOLATIONS, + ) + + +@pytest.mark.parametrize( + "threshold,unit,status", + [ + (EXPECTED_VIOLATIONS, "ABSOLUTE", TestCaseStatus.Success), + (EXPECTED_VIOLATIONS - 1, "ABSOLUTE", TestCaseStatus.Failed), + (50, "PERCENTAGE", TestCaseStatus.Success), # 3 out of 7 rows is 42.86% + (40, "PERCENTAGE", TestCaseStatus.Failed), + ], +) +def test_row_tolerance_decides_the_value_verdict(sqa_runner, pandas_runner, threshold, unit, status): + """The verdict is the violation count against the threshold, not the MIN/MAX""" + sqa_result = SQAValuesValidator( + sqa_runner, values_test_case(threshold, unit), EXECUTION_DATE.timestamp() + ).run_validation() + pandas_result = PandasValuesValidator( + pandas_runner, values_test_case(threshold, unit), EXECUTION_DATE.timestamp() + ).run_validation() + + for result in (sqa_result, pandas_result): + assert result.testCaseStatus == status + assert result.failedRows == EXPECTED_VIOLATIONS + assert result.passedRows == EXPECTED_ROWS - EXPECTED_VIOLATIONS + + +@pytest.mark.parametrize( + "threshold,unit,status", + [ + (EXPECTED_VIOLATIONS, "ABSOLUTE", TestCaseStatus.Success), + (EXPECTED_VIOLATIONS - 1, "ABSOLUTE", TestCaseStatus.Failed), + (50, "PERCENTAGE", TestCaseStatus.Success), + (40, "PERCENTAGE", TestCaseStatus.Failed), + ], +) +def test_row_tolerance_decides_the_length_verdict(sqa_runner, pandas_runner, threshold, unit, status): + sqa_result = SQALengthsValidator( + sqa_runner, lengths_test_case(threshold, unit), EXECUTION_DATE.timestamp() + ).run_validation() + pandas_result = PandasLengthsValidator( + pandas_runner, lengths_test_case(threshold, unit), EXECUTION_DATE.timestamp() + ).run_validation() + + assert sqa_result.testCaseStatus == status + assert pandas_result.testCaseStatus == status + + +def test_min_and_max_are_still_reported(sqa_runner): + """Users read the extremes, so a row tolerance does not take them off the result""" + result = SQAValuesValidator(sqa_runner, values_test_case(threshold=5), EXECUTION_DATE.timestamp()).run_validation() + + assert [(value.name, value.value) for value in result.testResultValue] == [("min", "1"), ("max", "12")] + + lengths = SQALengthsValidator( + sqa_runner, lengths_test_case(threshold=5), EXECUTION_DATE.timestamp() + ).run_validation() + + assert [(value.name, value.value) for value in lengths.testResultValue] == [ + ("minValueLength", "1"), + ("maxValueLength", "12"), + ] + + +def test_no_tolerance_keeps_the_min_max_verdict(sqa_runner): + """Without a tolerance nothing is counted: the extremes alone answer the same question""" + validator = SQAValuesValidator(sqa_runner, values_test_case(), EXECUTION_DATE.timestamp()) + + with patch.object(SQAValuesValidator, "_run_violation_count") as counted: + result = validator.run_validation() + + counted.assert_not_called() + assert result.testCaseStatus == TestCaseStatus.Failed + assert result.failedRows is None + + +def test_the_window_is_never_widened_by_the_threshold(sqa_runner): + """The tolerance is spent on rows here, so spending it on the bounds too would double it""" + validator = SQAValuesValidator(sqa_runner, values_test_case(threshold=5), EXECUTION_DATE.timestamp()) + + assert validator._get_test_parameters() == {"minValue": MIN_BOUND, "maxValue": MAX_BOUND} From c56b0cf48311a29ff57a3c27c2e84d61f7682dc4 Mon Sep 17 00:00:00 2001 From: TeddyCr <13626425+TeddyCr@users.noreply.github.com> Date: Wed, 23 Sep 2026 21:36:03 +0000 Subject: [PATCH 2/4] Review: report the counted violations, and treat a None bound as unbounded MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The between validators counted their violations and then, with computePassedFailedRowCount set, counted again through compute_row_count to fill passedRows/failedRows. The second scan reads the table again -- or, on a percentage sample, a different set of rows -- so the reported rows could contradict the status the first count decided. Both now report the counts the verdict was taken on and only fall back to compute_row_count when no violation count was taken at all. BetweenBoundsChecker also read only ∓inf as an unset bound, so a validator that resolves its bounds dynamically and finds none on one side would compare against None: junk SQL on the SQL half, a TypeError on the pandas one. Both halves now build a condition only for the sides the window actually sets. Co-Authored-By: Claude Opus 5 --- .../checkers/between_bounds_checker.py | 24 ++++++-- .../base/columnValueLengthsToBeBetween.py | 15 +++-- .../column/base/columnValuesToBeBetween.py | 15 +++-- .../test_between_row_violations.py | 58 ++++++++++++++++++- 4 files changed, 94 insertions(+), 18 deletions(-) diff --git a/ingestion/src/metadata/data_quality/validations/checkers/between_bounds_checker.py b/ingestion/src/metadata/data_quality/validations/checkers/between_bounds_checker.py index d44e0374d9e5..4daea3e758d9 100644 --- a/ingestion/src/metadata/data_quality/validations/checkers/between_bounds_checker.py +++ b/ingestion/src/metadata/data_quality/validations/checkers/between_bounds_checker.py @@ -35,11 +35,13 @@ def __init__(self, min_bound: float, max_bound: float): def _is_unbounded(self, bound: Any) -> bool: """Whether that side of the window lets everything through - An unset bound resolves to ∓inf and needs no condition at all. Only a float can be - infinite: a datetime window -- what a between test on a date column resolves to -- - compares fine and must not be handed to `math.isinf`, which only takes a number. + An unset bound resolves to ∓inf, and to None for a validator that resolves its bounds + dynamically and found none on that side. Neither excludes a value, so neither needs a + condition -- and None cannot be compared against at all. Only a float can be infinite: + a datetime window -- what a between test on a date column resolves to -- compares fine + and must not be handed to `math.isinf`, which only takes a number. """ - return isinstance(bound, float) and math.isinf(bound) + return bound is None or (isinstance(bound, float) and math.isinf(bound)) def _check_violations(self, values): """Core violation check logic - works for both scalar and Series. @@ -55,7 +57,19 @@ def _check_violations(self, values): """ import pandas as pd - return ~pd.isna(values) & ((values < self.min_bound) | (values > self.max_bound)) + # Only the sides the window actually sets are compared against, like the SQL half + # builds only those conditions: an unset bound excludes nothing, and a None one cannot + # be compared against at all. + outside = None + if not self._is_unbounded(self.min_bound): + outside = values < self.min_bound + if not self._is_unbounded(self.max_bound): + above = values > self.max_bound + outside = above if outside is None else (outside | above) + + # `& False` over a window with neither side set keeps the shape of `values`, so a caller + # masking a Series with the result still gets a Series back. + return ~pd.isna(values) & (False if outside is None else outside) def _value_violates(self, value: Any) -> bool: """Check violation of one value (scalar). diff --git a/ingestion/src/metadata/data_quality/validations/column/base/columnValueLengthsToBeBetween.py b/ingestion/src/metadata/data_quality/validations/column/base/columnValueLengthsToBeBetween.py index 3795f0999554..b0a1f6d1da9a 100644 --- a/ingestion/src/metadata/data_quality/validations/column/base/columnValueLengthsToBeBetween.py +++ b/ingestion/src/metadata/data_quality/validations/column/base/columnValueLengthsToBeBetween.py @@ -76,6 +76,9 @@ def _run_validation(self) -> TestCaseResult: } if self._needs_violation_count(): + # The counts go under the same two keys a dimension row carries its own under: + # the verdict and the message are read off `metric_values` by code shared with + # the dimensional path, which only finds them by those names. total_rows, violating_rows = self._run_violation_count(column, test_params) metric_values[DIMENSION_TOTAL_COUNT_KEY] = total_rows metric_values[DIMENSION_FAILED_COUNT_KEY] = violating_rows @@ -94,14 +97,16 @@ def _run_validation(self) -> TestCaseResult: ], ) - if self.test_case.computePassedFailedRowCount: + # A row tolerance already counted both, so report the counts the verdict was taken on + # rather than counting twice: a second scan reads the table again -- or, on a percentage + # sample, a different set of rows -- and could report rows that contradict the status. + row_count = metric_values.get(DIMENSION_TOTAL_COUNT_KEY) + failed_rows = metric_values.get(DIMENSION_FAILED_COUNT_KEY) + + if failed_rows is None and self.test_case.computePassedFailedRowCount: row_count, failed_rows = self.compute_row_count( column, test_params[self.MIN_BOUND], test_params[self.MAX_BOUND] ) - else: - # A row tolerance already counted both, so report them rather than counting twice. - row_count = metric_values.get(DIMENSION_TOTAL_COUNT_KEY) - failed_rows = metric_values.get(DIMENSION_FAILED_COUNT_KEY) evaluation = self._evaluate_test_condition(metric_values, test_params) result_message = self._format_result_message(metric_values, test_params=test_params) diff --git a/ingestion/src/metadata/data_quality/validations/column/base/columnValuesToBeBetween.py b/ingestion/src/metadata/data_quality/validations/column/base/columnValuesToBeBetween.py index 617015e2e11d..91d2396e672d 100644 --- a/ingestion/src/metadata/data_quality/validations/column/base/columnValuesToBeBetween.py +++ b/ingestion/src/metadata/data_quality/validations/column/base/columnValuesToBeBetween.py @@ -82,6 +82,9 @@ def _run_validation(self) -> TestCaseResult: } if self._needs_violation_count(): + # The counts go under the same two keys a dimension row carries its own under: + # the verdict and the message are read off `metric_values` by code shared with + # the dimensional path, which only finds them by those names. total_rows, violating_rows = self._run_violation_count(column, test_params) metric_values[DIMENSION_TOTAL_COUNT_KEY] = total_rows metric_values[DIMENSION_FAILED_COUNT_KEY] = violating_rows @@ -99,14 +102,16 @@ def _run_validation(self) -> TestCaseResult: ], ) - if self.test_case.computePassedFailedRowCount: + # A row tolerance already counted both, so report the counts the verdict was taken on + # rather than counting twice: a second scan reads the table again -- or, on a percentage + # sample, a different set of rows -- and could report rows that contradict the status. + row_count = metric_values.get(DIMENSION_TOTAL_COUNT_KEY) + failed_rows = metric_values.get(DIMENSION_FAILED_COUNT_KEY) + + if failed_rows is None and self.test_case.computePassedFailedRowCount: row_count, failed_rows = self.compute_row_count( column, test_params[self.MIN_BOUND], test_params[self.MAX_BOUND] ) - else: - # A row tolerance already counted both, so report them rather than counting twice. - row_count = metric_values.get(DIMENSION_TOTAL_COUNT_KEY) - failed_rows = metric_values.get(DIMENSION_FAILED_COUNT_KEY) evaluation = self._evaluate_test_condition(metric_values, test_params) result_message = self._format_result_message(metric_values, test_params=test_params) diff --git a/ingestion/tests/unit/observability/data_quality/validations/test_between_row_violations.py b/ingestion/tests/unit/observability/data_quality/validations/test_between_row_violations.py index 0ebb8d63e712..f84a6969a1ba 100644 --- a/ingestion/tests/unit/observability/data_quality/validations/test_between_row_violations.py +++ b/ingestion/tests/unit/observability/data_quality/validations/test_between_row_violations.py @@ -31,6 +31,9 @@ from metadata.data_quality.interface.sqlalchemy.sqa_test_suite_interface import ( SQATestSuiteInterface, ) +from metadata.data_quality.validations.checkers.between_bounds_checker import ( + BetweenBoundsChecker, +) from metadata.data_quality.validations.column.pandas.columnValueLengthsToBeBetween import ( ColumnValueLengthsToBeBetweenValidator as PandasLengthsValidator, ) @@ -151,7 +154,7 @@ def pandas_runner(): return PandasRunner(dataset=lambda: iter(frames), raw_dataset=None) -def values_test_case(threshold=None, unit=None) -> TestCase: +def values_test_case(threshold=None, unit=None, compute_row_count=False) -> TestCase: """A values-to-be-between test case on the fixture window""" return _test_case( ENTITY_LINK_VALUE, @@ -161,10 +164,11 @@ def values_test_case(threshold=None, unit=None) -> TestCase: ], threshold, unit, + compute_row_count, ) -def lengths_test_case(threshold=None, unit=None) -> TestCase: +def lengths_test_case(threshold=None, unit=None, compute_row_count=False) -> TestCase: """A lengths-to-be-between test case on the fixture window""" return _test_case( ENTITY_LINK_LABEL, @@ -174,10 +178,11 @@ def lengths_test_case(threshold=None, unit=None) -> TestCase: ], threshold, unit, + compute_row_count, ) -def _test_case(entity_link, parameter_values, threshold, unit) -> TestCase: +def _test_case(entity_link, parameter_values, threshold, unit, compute_row_count=False) -> TestCase: if threshold is not None: parameter_values = [*parameter_values, TestCaseParameterValue(name="threshold", value=str(threshold))] if unit is not None: @@ -188,6 +193,7 @@ def _test_case(entity_link, parameter_values, threshold, unit) -> TestCase: testSuite=EntityReference(id=uuid4(), type="TestSuite"), # type: ignore testDefinition=EntityReference(id=uuid4(), type="TestDefinition"), # type: ignore parameterValues=parameter_values, + computePassedFailedRowCount=compute_row_count, ) # type: ignore @@ -304,6 +310,52 @@ def test_no_tolerance_keeps_the_min_max_verdict(sqa_runner): assert result.failedRows is None +@pytest.mark.parametrize( + "validator_class,test_case_builder", + [ + (SQAValuesValidator, values_test_case), + (SQALengthsValidator, lengths_test_case), + ], +) +def test_row_reporting_reuses_the_counted_violations(sqa_runner, validator_class, test_case_builder): + """The reported rows are the ones the verdict was taken on, not a second scan's + + A table can change between two scans, and a percentage sample draws different rows each + time, so counting again could report rows that contradict the status. + """ + test_case = test_case_builder(threshold=EXPECTED_VIOLATIONS, unit="ABSOLUTE", compute_row_count=True) + validator = validator_class(sqa_runner, test_case, EXECUTION_DATE.timestamp()) + + with patch.object(validator_class, "compute_row_count") as counted_again: + result = validator.run_validation() + + counted_again.assert_not_called() + assert result.testCaseStatus == TestCaseStatus.Success + assert result.failedRows == EXPECTED_VIOLATIONS + assert result.passedRows == EXPECTED_ROWS - EXPECTED_VIOLATIONS + + +@pytest.mark.parametrize( + "min_bound,max_bound,violations", + [ + (None, None, 0), # a window with neither side set excludes nothing + (MIN_BOUND, None, 1), # only the 1 is below 3 + (None, MAX_BOUND, 2), # only the 9 and the 12 are above 8 + ], +) +def test_an_unset_bound_excludes_nothing(pandas_runner, min_bound, max_bound, violations): + """A bound a validator resolved to None is a side of the window that lets everything through + + Both engines have to read it that way: the SQL half builds no condition for it, so the + pandas half must not compare against it either. + """ + checker = BetweenBoundsChecker(min_bound=min_bound, max_bound=max_bound) + + counted = sum(int(checker.get_violations_mask(df["value"]).sum()) for df in pandas_runner) + + assert counted == violations + + def test_the_window_is_never_widened_by_the_threshold(sqa_runner): """The tolerance is spent on rows here, so spending it on the bounds too would double it""" validator = SQAValuesValidator(sqa_runner, values_test_case(threshold=5), EXECUTION_DATE.timestamp()) From 2bc7ca8e5006ca0364486db9519bdbc93a0b78cd Mon Sep 17 00:00:00 2001 From: TeddyCr Date: Mon, 28 Sep 2026 15:55:12 -0700 Subject: [PATCH 3/4] Fix basedpyright errors in the between row-violation count Keep the base signature on the _run_violation_count overrides, narrow the runner to QueryRunner, type the checker's count as a ColumnElement so it can be labelled, and fail loudly on an empty result row. Pin that an explicit zero threshold does not pay for the count query. --- .../validations/checkers/between_bounds_checker.py | 4 ++-- .../validations/column/base/columnValuesToBeBetween.py | 3 ++- .../column/pandas/columnValueLengthsToBeBetween.py | 3 ++- .../column/pandas/columnValuesToBeBetween.py | 3 ++- .../column/sqlalchemy/columnValueLengthsToBeBetween.py | 10 ++++++++-- .../column/sqlalchemy/columnValuesToBeBetween.py | 10 ++++++++-- .../validations/mixins/sqa_validator_mixin.py | 7 +++++-- .../validations/test_between_row_violations.py | 10 +++++++--- 8 files changed, 36 insertions(+), 14 deletions(-) diff --git a/ingestion/src/metadata/data_quality/validations/checkers/between_bounds_checker.py b/ingestion/src/metadata/data_quality/validations/checkers/between_bounds_checker.py index 4daea3e758d9..ed3d77b627c2 100644 --- a/ingestion/src/metadata/data_quality/validations/checkers/between_bounds_checker.py +++ b/ingestion/src/metadata/data_quality/validations/checkers/between_bounds_checker.py @@ -22,7 +22,7 @@ ) if TYPE_CHECKING: - from sqlalchemy.sql.elements import ClauseElement + from sqlalchemy.sql.elements import ClauseElement, ColumnElement class BetweenBoundsChecker(BaseValidationChecker): @@ -119,7 +119,7 @@ def build_violation_sqa(self, metrics: list["ClauseElement"]) -> "ClauseElement" return literal(False) return or_(*conditions) if len(conditions) > 1 else conditions[0] - def build_row_level_violations_sqa(self, column: "ClauseElement") -> "ClauseElement": + def build_row_level_violations_sqa(self, column: "ClauseElement") -> "ColumnElement": """Build SQL expression to count row-level violations. Returns a SUM(CASE...) expression that counts individual rows where diff --git a/ingestion/src/metadata/data_quality/validations/column/base/columnValuesToBeBetween.py b/ingestion/src/metadata/data_quality/validations/column/base/columnValuesToBeBetween.py index 91d2396e672d..4f6c7509e6d7 100644 --- a/ingestion/src/metadata/data_quality/validations/column/base/columnValuesToBeBetween.py +++ b/ingestion/src/metadata/data_quality/validations/column/base/columnValuesToBeBetween.py @@ -16,6 +16,7 @@ import traceback from abc import abstractmethod from datetime import date, datetime, time +from typing import Any from sqlalchemy import Column @@ -76,7 +77,7 @@ def _run_validation(self) -> TestCaseResult: min_res = self._normalize_metric_value(min_res, is_min=True) max_res = self._normalize_metric_value(max_res, is_min=False) - metric_values = { + metric_values: dict[str, Any] = { Metrics.min.name: min_res, Metrics.max.name: max_res, } diff --git a/ingestion/src/metadata/data_quality/validations/column/pandas/columnValueLengthsToBeBetween.py b/ingestion/src/metadata/data_quality/validations/column/pandas/columnValueLengthsToBeBetween.py index 7e17a34e8602..2bb8ab23894d 100644 --- a/ingestion/src/metadata/data_quality/validations/column/pandas/columnValueLengthsToBeBetween.py +++ b/ingestion/src/metadata/data_quality/validations/column/pandas/columnValueLengthsToBeBetween.py @@ -17,6 +17,7 @@ from typing import cast import pandas as pd +from sqlalchemy import Column from metadata.data_quality.validations.base_test_handler import ( DIMENSION_FAILED_COUNT_KEY, @@ -63,7 +64,7 @@ def _run_results(self, metric: Metrics, column: SQALikeColumn) -> int | None: """ return self.run_dataframe_results(self.runner, metric, column) - def _run_violation_count(self, column: SQALikeColumn, test_params: dict) -> tuple[int, int]: + def _run_violation_count(self, column: SQALikeColumn | Column, test_params: dict) -> tuple[int, int]: """Count the rows read and the values whose length falls outside the window The dataframes are walked one at a time rather than concatenated, like every other diff --git a/ingestion/src/metadata/data_quality/validations/column/pandas/columnValuesToBeBetween.py b/ingestion/src/metadata/data_quality/validations/column/pandas/columnValuesToBeBetween.py index 5dc14d60de01..56d8ae65dffe 100644 --- a/ingestion/src/metadata/data_quality/validations/column/pandas/columnValuesToBeBetween.py +++ b/ingestion/src/metadata/data_quality/validations/column/pandas/columnValuesToBeBetween.py @@ -18,6 +18,7 @@ from typing import cast import pandas as pd +from sqlalchemy import Column from metadata.data_quality.validations.base_test_handler import ( DIMENSION_FAILED_COUNT_KEY, @@ -66,7 +67,7 @@ def _run_results(self, metric: Metrics, column: SQALikeColumn) -> int | None: """ return self.run_dataframe_results(self.runner, metric, column) - def _run_violation_count(self, column: SQALikeColumn, test_params: dict) -> tuple[int, int]: + def _run_violation_count(self, column: SQALikeColumn | Column, test_params: dict) -> tuple[int, int]: """Count the rows read and the values falling outside the window The dataframes are walked one at a time rather than concatenated, like every other diff --git a/ingestion/src/metadata/data_quality/validations/column/sqlalchemy/columnValueLengthsToBeBetween.py b/ingestion/src/metadata/data_quality/validations/column/sqlalchemy/columnValueLengthsToBeBetween.py index e88eb52ee370..b61688d80895 100644 --- a/ingestion/src/metadata/data_quality/validations/column/sqlalchemy/columnValueLengthsToBeBetween.py +++ b/ingestion/src/metadata/data_quality/validations/column/sqlalchemy/columnValueLengthsToBeBetween.py @@ -14,6 +14,7 @@ """ import math +from typing import cast from sqlalchemy import Column @@ -37,7 +38,9 @@ from metadata.generated.schema.tests.dimensionResult import DimensionResult from metadata.profiler.metrics.registry import Metrics from metadata.profiler.orm.functions.length import LenFn +from metadata.profiler.processor.runner import QueryRunner from metadata.utils.logger import test_suite_logger +from metadata.utils.sqa_like_column import SQALikeColumn logger = test_suite_logger() @@ -59,7 +62,7 @@ def _run_results(self, metric: Metrics, column: Column) -> int | None: """ return self.run_query_results(self.runner, metric, column) - def _run_violation_count(self, column: Column, test_params: dict) -> tuple[int | None, int | None]: + def _run_violation_count(self, column: SQALikeColumn | Column, test_params: dict) -> tuple[int | None, int | None]: """Count the rows read and the values whose length falls outside the window Args: @@ -67,7 +70,10 @@ def _run_violation_count(self, column: Column, test_params: dict) -> tuple[int | test_params: test parameters including min and max bounds """ checker = self._get_validation_checker(test_params) - return self._compute_row_violations(self.runner, checker.build_row_level_violations_sqa(LenFn(column))) + return self._compute_row_violations( + cast(QueryRunner, self.runner), # noqa: TC006 + checker.build_row_level_violations_sqa(LenFn(column)), + ) def compute_row_count(self, column: Column, min_bound: int, max_bound: int): """Compute row count for the given column diff --git a/ingestion/src/metadata/data_quality/validations/column/sqlalchemy/columnValuesToBeBetween.py b/ingestion/src/metadata/data_quality/validations/column/sqlalchemy/columnValuesToBeBetween.py index 779ac64f9987..86043fff864f 100644 --- a/ingestion/src/metadata/data_quality/validations/column/sqlalchemy/columnValuesToBeBetween.py +++ b/ingestion/src/metadata/data_quality/validations/column/sqlalchemy/columnValuesToBeBetween.py @@ -15,6 +15,7 @@ import math from datetime import datetime +from typing import cast from sqlalchemy import Column @@ -38,7 +39,9 @@ from metadata.generated.schema.tests.dimensionResult import DimensionResult from metadata.profiler.metrics.registry import Metrics from metadata.profiler.orm.registry import is_date_time +from metadata.profiler.processor.runner import QueryRunner from metadata.utils.logger import test_suite_logger +from metadata.utils.sqa_like_column import SQALikeColumn from metadata.utils.time_utils import convert_timestamp logger = test_suite_logger() @@ -61,7 +64,7 @@ def _run_results(self, metric: Metrics, column: Column) -> int | None: """ return self.run_query_results(self.runner, metric, column) - def _run_violation_count(self, column: Column, test_params: dict) -> tuple[int | None, int | None]: + def _run_violation_count(self, column: SQALikeColumn | Column, test_params: dict) -> tuple[int | None, int | None]: """Count the rows read and the values falling outside the window Args: @@ -69,7 +72,10 @@ def _run_violation_count(self, column: Column, test_params: dict) -> tuple[int | test_params: test parameters including min and max bounds """ checker = self._get_validation_checker(test_params) - return self._compute_row_violations(self.runner, checker.build_row_level_violations_sqa(column)) + return self._compute_row_violations( + cast(QueryRunner, self.runner), # noqa: TC006 + checker.build_row_level_violations_sqa(cast(Column, column)), # noqa: TC006 + ) def _build_dimension_metric_values(self, row, metrics_to_compute, test_params=None): min_value = row.get(Metrics.min.name) diff --git a/ingestion/src/metadata/data_quality/validations/mixins/sqa_validator_mixin.py b/ingestion/src/metadata/data_quality/validations/mixins/sqa_validator_mixin.py index 77dcd5fabbb8..43961e48a87e 100644 --- a/ingestion/src/metadata/data_quality/validations/mixins/sqa_validator_mixin.py +++ b/ingestion/src/metadata/data_quality/validations/mixins/sqa_validator_mixin.py @@ -185,7 +185,7 @@ def _compute_row_count_between( return res - def _compute_row_violations(self, runner: QueryRunner, violations_expr: ClauseElement) -> tuple[Any, Any]: + def _compute_row_violations(self, runner: QueryRunner, violations_expr: ColumnElement) -> tuple[Any, Any]: """Count the rows read and the violating ones among them, in a single aggregate query Both counts come from the same scan so they are always counted against each other: a @@ -204,10 +204,13 @@ def _compute_row_violations(self, runner: QueryRunner, violations_expr: ClauseEl Metrics.rowCount().fn().label(DIMENSION_TOTAL_COUNT_KEY), violations_expr.label(DIMENSION_FAILED_COUNT_KEY), ) - values = dict(row._mapping) except Exception as exc: raise SQLAlchemyError(exc) # noqa: B904 + if row is None: + raise SQLAlchemyError("The row violation count query returned no row") + + values = dict(row._mapping) return values.get(DIMENSION_TOTAL_COUNT_KEY), values.get(DIMENSION_FAILED_COUNT_KEY) def _compute_row_count(self, runner: QueryRunner, column: Column, **kwargs): diff --git a/ingestion/tests/unit/observability/data_quality/validations/test_between_row_violations.py b/ingestion/tests/unit/observability/data_quality/validations/test_between_row_violations.py index f84a6969a1ba..8c799e4348c5 100644 --- a/ingestion/tests/unit/observability/data_quality/validations/test_between_row_violations.py +++ b/ingestion/tests/unit/observability/data_quality/validations/test_between_row_violations.py @@ -298,9 +298,13 @@ def test_min_and_max_are_still_reported(sqa_runner): ] -def test_no_tolerance_keeps_the_min_max_verdict(sqa_runner): - """Without a tolerance nothing is counted: the extremes alone answer the same question""" - validator = SQAValuesValidator(sqa_runner, values_test_case(), EXECUTION_DATE.timestamp()) +@pytest.mark.parametrize("threshold", [None, 0, 0.0]) +def test_no_tolerance_keeps_the_min_max_verdict(sqa_runner, threshold): + """Without a tolerance nothing is counted: the extremes alone answer the same question + + A threshold explicitly set to zero tolerates nothing either, so it must not pay for the count. + """ + validator = SQAValuesValidator(sqa_runner, values_test_case(threshold), EXECUTION_DATE.timestamp()) with patch.object(SQAValuesValidator, "_run_violation_count") as counted: result = validator.run_validation() From 59581d6226feff182909c914a49af5202390e0a2 Mon Sep 17 00:00:00 2001 From: TeddyCr Date: Mon, 28 Sep 2026 18:02:30 -0700 Subject: [PATCH 4/4] Point the Playwright impact map at the moved login specs #34056 moved Login.spec.ts and LoginConfiguration.spec.ts from e2e/Pages to e2e/Auth without updating impact-map.json. Every targeted plan picks the smoke entry up, so the planner fails on a selector whose spec no longer exists. Auth/** is delegated to the SSO lane, so the moved paths are dropped from the main lanes instead of failing the plan. --- .github/playwright/impact-map.json | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/.github/playwright/impact-map.json b/.github/playwright/impact-map.json index 9b207eb7808e..d1402f9680f5 100644 --- a/.github/playwright/impact-map.json +++ b/.github/playwright/impact-map.json @@ -12,7 +12,7 @@ { "projects": ["Basic"], "specs": [ - "playwright/e2e/Pages/Login.spec.ts", + "playwright/e2e/Auth/Login.spec.ts", "playwright/e2e/Flow/Navbar.spec.ts", "playwright/e2e/Features/Dashboards.spec.ts", "playwright/e2e/Pages/Lineage/LineageControls.spec.ts", @@ -354,7 +354,7 @@ "specs": [ "playwright/e2e/Features/SettingsNavigationPage.spec.ts", "playwright/e2e/Flow/LineageSettings.spec.ts", - "playwright/e2e/Pages/LoginConfiguration.spec.ts", + "playwright/e2e/Auth/LoginConfiguration.spec.ts", "playwright/e2e/Pages/OmdURLConfiguration.spec.ts", "playwright/e2e/Pages/ProfilerConfigurationPage.spec.ts", "playwright/e2e/Pages/SearchSettings.spec.ts", @@ -393,7 +393,7 @@ "projects": ["chromium", "Basic", "Ingestion", "SearchRBAC", "ImportExport"], "specs": [ "playwright/e2e/**/*Permission*.spec.ts", - "playwright/e2e/Pages/Login*.spec.ts", + "playwright/e2e/Auth/Login*.spec.ts", "playwright/e2e/Flow/SearchRBAC.spec.ts" ] }, @@ -467,7 +467,7 @@ ], "projects": ["chromium", "Basic"], "specs": [ - "playwright/e2e/Pages/Login.spec.ts", + "playwright/e2e/Auth/Login.spec.ts", "playwright/e2e/Flow/Navbar.spec.ts" ] },