Skip to main content

Data Quality framework for Databricks — rule-based checks across six dimensions with Delta table persistence

Project description

NTB DQ Framework

Data Quality framework for Databricks — rule-based checks across six dimensions with Delta table persistence.

Installation

pip install ntb-dq-framework

Requires Python >= 3.9.

Quick Start

from pyspark.sql import SparkSession
from ntb_dq_framework import DQEngine

spark = SparkSession.builder.getOrCreate()

# Initialize — creates Delta tables if they don't exist yet
engine = DQEngine(spark, catalog="my_catalog")

# Register a rule
rule_id = engine.add_rule(
    rule_name="null_email",
    dimension="Completeness",
    target_table="my_catalog.my_schema.customers",
    target_column="email",
    rule_type="null_check",
    severity="critical",
    owner="data-team",
)

# Run all active rules for a table
summary = engine.run_checks("my_catalog.my_schema.customers")
print(f"Passed: {summary.total_rows_passed}, Failed: {summary.total_rows_failed}")

# Query results
results_df = engine.get_results(run_id=summary.run_id)
results_df.show()

Dimensions and Rule Types

How rules work: You always define what valid (good) data looks like. The framework then reports any rows that don't match as failures. For example, allowed_values=["Male", "Female"] means rows with those values pass — anything else fails. This applies to all rule types: regex patterns describe valid formats, range bounds describe acceptable values, business rules describe conditions that should be true, etc.

Accuracy

Rule Type Description Key Parameters
reference_match Joins target with a reference table and reports mismatched rows reference_table, key_column, comparison_column
tolerance_match Like reference_match but allows a numeric tolerance reference_table, key_column, comparison_column, tolerance
reconciliation Compares aggregate values (sum/count) between target and reference reference_table, aggregate_column, aggregate_function

Completeness

Rule Type Description Key Parameters
null_check Counts rows where the target column is null or empty
missing_partition Compares expected date partitions against actual values expected_partitions, partition_column
volume_check Checks row count against a baseline threshold baseline, threshold_percentage

Consistency

Rule Type Description Key Parameters
referential_integrity Left anti-join to find orphan rows missing from a reference table reference_table, key_column
cross_system_match Joins on key columns and reports rows where compared columns differ reference_table, key_columns, comparison_columns
business_rule Evaluates a SQL boolean expression and reports rows where it's false Uses rule_expression

Timeliness

Rule Type Description Key Parameters
freshness_check Computes delay in hours since the latest timestamp value
sla_check Fails when freshness delay exceeds a threshold sla_threshold_hours

Validity

Rule Type Description Key Parameters
regex_check Reports non-null rows not matching a regex pattern (nulls are excluded; column is auto-cast to string) Uses rule_expression
allowed_values Reports rows with values not in an allowed list allowed_values
range_check Reports rows outside min/max bounds min, max
type_check Reports rows that can't be cast to a target type target_type

Uniqueness

Rule Type Description Key Parameters
duplicate_check Groups by key columns and identifies duplicate groups key_columns

API Reference

DQEngine

engine = DQEngine(spark, catalog="my_catalog", schema="my_schema")
Parameter Type Default Description
spark SparkSession Active Spark session
catalog str Unity Catalog name
schema str "data_quality" Schema for DQ tables. Optional — defaults to data_quality if omitted
table_prefix str "dq_" Prefix for all Delta table names
max_failed_samples int 100 Max failing rows to sample per rule
Method Description
add_rule(...) Register a rule. Returns the rule_id. Deduplicates identical configs automatically. See add_rule Parameters below.
deactivate_rule(rule_id) Soft-delete a rule by setting is_active = False.
run_checks(target_table, df=None, run_name=None, triggered_by="manual") Execute all active rules for a table. Returns a RunSummary.
get_results(run_id=None, rule_id=None, dimension=None, target_table=None, date_range=None) Query the results table with optional filters. Returns a Spark DataFrame.
daily_summary(run_date=None) Aggregate run log data by date and target table.
dimension_summary(run_date=None) Aggregate results by date, dimension, and target table.
rule_trend(rule_id=None, days=30) Result data grouped by rule and date over a time window.
top_offenders(run_date=None, top_n=10) Top N rules with the lowest pass rate.

add_rule Parameters

engine.add_rule(
    rule_name="...",
    dimension="...",
    target_table="...",
    rule_type="...",
    rule_expression="",        # SQL expression or regex pattern (for regex_check, business_rule)
    target_column=None,        # Column to check
    parameters=None,           # Dict of rule-type-specific parameters (see below)
    severity="warning",        # "critical" or "warning"
    owner="",
)
Parameter Type Default Description
rule_name str Human-readable name for the rule
dimension str One of Accuracy, Completeness, Consistency, Timeliness, Validity, Uniqueness
target_table str Fully qualified table name (e.g., catalog.schema.table)
rule_type str Check type (see Dimensions and Rule Types)
rule_expression str "" SQL expression or regex pattern (used by regex_check, business_rule)
target_column str | None None Column to evaluate
parameters dict | None None Rule-type-specific parameters (see Key Parameters in each dimension table)
severity str "warning" "critical" or "warning"
owner str "" Owner or team responsible for this rule

Note: Rule-type-specific options like allowed_values, min, max, key_columns, etc. must be passed inside the parameters dict — not as top-level keyword arguments.

Examples:

# Validity — allowed_values
engine.add_rule(
    rule_name="valid_gender",
    dimension="Validity",
    target_table="my_catalog.my_schema.customers",
    target_column="gender",
    rule_type="allowed_values",
    severity="warning",
    owner="data-team",
    parameters={"allowed_values": ["Male", "Female"]},
)

# Validity — range_check
engine.add_rule(
    rule_name="valid_age",
    dimension="Validity",
    target_table="my_catalog.my_schema.customers",
    target_column="age",
    rule_type="range_check",
    parameters={"min": 0, "max": 120},
)

# Uniqueness — duplicate_check
engine.add_rule(
    rule_name="unique_order",
    dimension="Uniqueness",
    target_table="my_catalog.my_schema.orders",
    rule_type="duplicate_check",
    parameters={"key_columns": ["order_id"]},
)

Severity Levels

Rules support two severity levels: critical and warning.

  • critical — recorded in results with high priority. Does not block pipeline execution.
  • warning — recorded in results but does not block execution.

Notes on rule_expression

When using regex patterns in rule_expression (e.g., for regex_check or business_rule), use a raw string to preserve backslashes:

engine.add_rule(
    rule_name="13_digit_id",
    dimension="Validity",
    target_table="my_catalog.my_schema.employees",
    target_column="id_card",
    rule_type="regex_check",
    rule_expression=r"^\d{13}$",  # raw string — backslashes are preserved
    severity="warning",
    owner="data-team",
)

Project Structure

ntb_dq_framework/
├── __init__.py            # Public API — exports DQEngine
├── engine.py              # DQEngine entry point
├── models.py              # RuleConfig, CheckResult, RunSummary dataclasses
├── rule_manager.py        # CRUD operations on the rule registry
├── run_executor.py        # Orchestrates check execution
├── result_writer.py       # Persists results to Delta tables
├── table_initializer.py   # Creates Delta tables on first run
├── monitoring.py          # Aggregation views for dashboards
├── exceptions.py          # DQValidationError
└── checks/
    ├── base.py            # BaseCheck abstract class
    ├── accuracy.py
    ├── completeness.py
    ├── consistency.py
    ├── timeliness.py
    ├── validity.py
    └── uniqueness.py

Delta Tables

On initialization, the framework creates four Delta tables (prefixed with dq_ by default):

Table Purpose
dq_rule_registry Stores rule definitions with activation status
dq_run_log Tracks each run execution with summary stats
dq_results Per-rule results for every run
dq_failed_records Sampled failing rows (up to max_failed_samples)

Changelog

0.2.0

  • run_checks no longer raises DQRunError when critical-severity rules fail or error. The method now returns a RunSummary containing all results (including critical failures), allowing the pipeline to continue. Results are still persisted to Delta tables as before.

0.1.5

  • Initial release with six-dimension rule-based checks and Delta table persistence.

Project details


Download files

Download the file for your platform. If you're not sure which to choose, learn more about installing packages.

Source Distribution

ntb_dq_framework-0.2.0.tar.gz (26.5 kB view details)

Uploaded Source

Built Distribution

If you're not sure about the file name format, learn more about wheel file names.

ntb_dq_framework-0.2.0-py3-none-any.whl (35.0 kB view details)

Uploaded Python 3

File details

Details for the file ntb_dq_framework-0.2.0.tar.gz.

File metadata

  • Download URL: ntb_dq_framework-0.2.0.tar.gz
  • Upload date:
  • Size: 26.5 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? No
  • Uploaded via: twine/6.2.0 CPython/3.11.7

File hashes

Hashes for ntb_dq_framework-0.2.0.tar.gz
Algorithm Hash digest
SHA256 0bbb12cc7b3cb45ffa0552c73d0932398cd20748c97a9a7d9e71d5e7b442c644
MD5 9f1d3513ad7c320af749df0a3cc6adc2
BLAKE2b-256 1f9ee477452bce15f41151ca21cdf104b3fbd31d2db1b0e9225f3da5ff06507e

See more details on using hashes here.

File details

Details for the file ntb_dq_framework-0.2.0-py3-none-any.whl.

File metadata

File hashes

Hashes for ntb_dq_framework-0.2.0-py3-none-any.whl
Algorithm Hash digest
SHA256 8652dcf888d100fba905f7468a51d3b853cf84503f98a601cd7b8592598f990f
MD5 b9a2e92036bb89b39eaf9ced6b96b6f2
BLAKE2b-256 7753b5bf6275a38096bc90f69503386ef2afc87d55728567c82fbf4950a69e8e

See more details on using hashes here.

Supported by

AWS Cloud computing and Security Sponsor Datadog Monitoring Depot Continuous Integration Fastly CDN Google Download Analytics Pingdom Monitoring Sentry Error logging StatusPage Status page