Metadata-Version: 2.4
Name: spark-data-quality
Version: 1.0.14
Summary: SparkDQAgent — Data Quality validation package for K8s Spark pods
Author-email: khailas <khailas.rangath@saal.ai>
License-Expression: MIT
Classifier: Programming Language :: Python :: 3
Classifier: Operating System :: OS Independent
Requires-Python: >=3.8
Description-Content-Type: text/markdown
License-File: LICENSE
Requires-Dist: requests>=2.28.0
Requires-Dist: great-expectations==0.18.12
Requires-Dist: trino>=0.320.0
Requires-Dist: PyYAML>=6.0
Provides-Extra: spark
Requires-Dist: pyspark>=3.1.1; extra == "spark"
Provides-Extra: dev
Requires-Dist: pytest; extra == "dev"
Requires-Dist: ruff; extra == "dev"
Dynamic: license-file

# 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_id` column 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 `amount` column **between 1 and 500,000** (no negatives, no unrealistic numbers)?
- Does the `region` column 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    │
└──────────────────┘        └──────────────┘
```

1. You create a `SparkDQAgent` and point it to your table using catalog, schema, and table name
2. The agent fetches rules from the DQ Engine (e.g., "column X must not be null")
3. You load your data into a Spark DataFrame
4. The agent runs every rule against your data
5. Results are saved to the DQ Engine and returned to you

## Installation

```bash
pip install spark-data-quality
```

## Quick Start

```python
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

```python
# 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
