spark-data-quality
What Is This For?
Imagine your pipeline loads 100,000 rows of sales data into a table every day. Before you use that data to make business decisions, you'd want to check:
- Does the
customer_idcolumn have any blank or missing values? - Did the table receive all 100,000 rows, or did some get lost and only 50,000 arrived?
- Is every value in the
amountcolumn between 1 and 500,000 (no negatives, no unrealistic numbers)? - Does the
regioncolumn only contain expected values like "North", "South", "East", "West" — or did something like "XYZ123" sneak in? - Are there exactly 4 unique regions, not 3 or 50?
Doing these checks manually is impossible at scale. This package automates that process. You define your rules once (in the DQ Engine dashboard), and every time new data arrives, the package checks it against all those rules and tells you what passed and what failed.
If something fails, you know immediately — before bad data reaches your reports, dashboards, or downstream systems.
How It Works
Your Data (DataFrame)
│
▼
┌──────────────────┐ ┌──────────────┐
│ SparkDQAgent │──GET──>│ DQ Engine │ "What rules should I check?"
│ │<───────│ │ (returns assertions)
│ Runs checks │ │ │
│ using Great │ │ │
│ Expectations │──POST─>│ │ "Here are the results"
│ │ │ Stores & │
│ Returns summary │ │ displays │
└──────────────────┘ └──────────────┘
- You create a
SparkDQAgentand point it to your table using catalog, schema, and table name - The agent fetches rules from the DQ Engine (e.g., "column X must not be null")
- You load your data into a Spark DataFrame
- The agent runs every rule against your data
- Results are saved to the DQ Engine and returned to you
Installation
pip install spark-data-quality
Quick Start
from pyspark.sql import SparkSession
from spark_dq.quality import SparkDQAgent
spark = SparkSession.builder.master("local[*]").getOrCreate()
# 1. Create the agent — just pass your table's catalog, schema, and name
agent = SparkDQAgent(
catalog="mycatalog",
schema="myschema",
table="sales_data",
data_quality_url="https://your-dq-engine.example.com/api/v1/spark",
catalog_type="unmanaged",
trino_host="trino.example.com:443",
trino_user="service_account",
trino_pwd="your_password",
)
# 2. Load your data
df = spark.read.format("jdbc") \
.option("url", "jdbc:trino://trino.example.com:443?SSL=true") \
.option("driver", "io.trino.jdbc.TrinoDriver") \
.option("user", "service_account") \
.option("password", "your_password") \
.option("query", "SELECT * FROM mycatalog.myschema.sales_data") \
.load()
# 3. Run checks and review results
results = agent.execute_data_quality(df)
for suite_name, stats in results.items():
passed = stats["successful_expectations"]
total = stats["evaluated_expectations"]
print(f"{suite_name}: {passed}/{total} passed")
spark.stop()
Key Methods
fetch_table_config()
Fetches the table's configuration from the DQ Engine — this includes all the assertions (rules) defined for your table, the SQL query to load data, and scan limits. The result is cached after the first call, so calling it multiple times does not make extra API requests.
execute_data_quality(df)
The main method. Runs all assertions against your Spark DataFrame using Great Expectations:
- Fetches config from DQ Engine (if not already fetched)
- Runs all assertions in parallel for speed
- Saves results to the DQ Engine (
POST /save-results) - Returns a summary with pass/fail counts
# Example output
{
"sales_data_93": {
"evaluated_expectations": 5,
"successful_expectations": 5,
"unsuccessful_expectations": 0,
"success_percent": 100.0
}
}
| Field | What It Means |
|---|---|
evaluated_expectations |
Total number of checks run |
successful_expectations |
How many passed |
unsuccessful_expectations |
How many failed |
success_percent |
Pass rate (100.0 = all passed) |
Common Assertion Types
| Check | Assertion Type |
|---|---|
| No null values | expect_column_values_to_not_be_null |
| All values unique | expect_column_values_to_be_unique |
| Values in a range | expect_column_values_to_be_between |
| Row count in a range | expect_table_row_count_to_be_between |
| Distinct value count | expect_column_unique_value_count_to_be_between |
| Median in a range | expect_column_median_to_be_between |
| Text matches pattern | expect_column_values_to_match_regex |
| Text length in range | expect_column_value_lengths_to_be_between |
Requirements
- Python 3.8+
- PySpark 3.1.1+
- Great Expectations 0.18.12
- A running DQ Engine instance with assertions configured
License
MIT
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distribution
Built Distribution
Filter files by name, interpreter, ABI, and platform.
If you're not sure about the file name format, learn more about wheel file names.
Copy a direct link to the current filters
File details
Details for the file spark_data_quality-1.0.14.tar.gz.
File metadata
- Download URL: spark_data_quality-1.0.14.tar.gz
- Upload date:
- Size: 14.2 kB
- Tags: Source
- Uploaded using Trusted Publishing? Yes
- Uploaded via:
twine/7.0.0 CPython/3.13.14
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
558bf2d436dd3e728c3200473b03d09db6b67567f2b40610aa6a1cb5498342fa
|
|
| MD5 |
2695d68ebf1e9964a05a615b2abde2ed
|
|
| BLAKE2b-256 |
cc19ec263e20c5a3285564cc208f1c30a49fd491473643a4c71c3099c0ce8124
|
Provenance
The following attestation bundles were made for spark_data_quality-1.0.14.tar.gz:
Publisher:
ci.yml on saal-core/digixt-quality-package
-
Statement:
-
Statement type:
https://in-toto.io/Statement/v1 -
Predicate type:
https://docs.pypi.org/attestations/publish/v1 -
Subject name:
spark_data_quality-1.0.14.tar.gz -
Subject digest:
558bf2d436dd3e728c3200473b03d09db6b67567f2b40610aa6a1cb5498342fa - Sigstore transparency entry: 2343848294
- Sigstore integration time:
-
Permalink:
saal-core/digixt-quality-package@b51c64a0f4014e7a804a21839649976d9e7a3cd0 -
Branch / Tag:
refs/tags/v1.0.14 - Owner: https://github.com/saal-core
-
Access:
private
-
Token Issuer:
https://token.actions.githubusercontent.com -
Runner Environment:
github-hosted -
Publication workflow:
ci.yml@b51c64a0f4014e7a804a21839649976d9e7a3cd0 -
Trigger Event:
push
-
Statement type:
File details
Details for the file spark_data_quality-1.0.14-py3-none-any.whl.
File metadata
- Download URL: spark_data_quality-1.0.14-py3-none-any.whl
- Upload date:
- Size: 18.3 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? Yes
- Uploaded via:
twine/7.0.0 CPython/3.13.14
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
d84a62c25a442b2d1c95842c9693a9ba01823c5ad933e5558490ff096435eadf
|
|
| MD5 |
698c23e23b1ebb6965b0c9f60b5566a7
|
|
| BLAKE2b-256 |
3fa163d0a63c003c289cab87e1e594fbe8ca839152b638a03f3c8b62f690de59
|
Provenance
The following attestation bundles were made for spark_data_quality-1.0.14-py3-none-any.whl:
Publisher:
ci.yml on saal-core/digixt-quality-package
-
Statement:
-
Statement type:
https://in-toto.io/Statement/v1 -
Predicate type:
https://docs.pypi.org/attestations/publish/v1 -
Subject name:
spark_data_quality-1.0.14-py3-none-any.whl -
Subject digest:
d84a62c25a442b2d1c95842c9693a9ba01823c5ad933e5558490ff096435eadf - Sigstore transparency entry: 2343848325
- Sigstore integration time:
-
Permalink:
saal-core/digixt-quality-package@b51c64a0f4014e7a804a21839649976d9e7a3cd0 -
Branch / Tag:
refs/tags/v1.0.14 - Owner: https://github.com/saal-core
-
Access:
private
-
Token Issuer:
https://token.actions.githubusercontent.com -
Runner Environment:
github-hosted -
Publication workflow:
ci.yml@b51c64a0f4014e7a804a21839649976d9e7a3cd0 -
Trigger Event:
push
-
Statement type: