Metadata-Version: 2.4
Name: spark-data-quality
Version: 1.1.0
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

### `execute_data_quality(df)`

The only method you need to call. You pass it your data (a Spark DataFrame) and it does everything else automatically.

**What happens when you call it:**

```
results = agent.execute_data_quality(df)
│
│  Step 1 — Fetch rules from DQ Engine
│  ├── Calls the DQ Engine API: "What checks exist for this table?"
│  └── DQ Engine responds with all assertions
│       (e.g., "customer_id not null", "amount between 1-500000")
│
│  Step 2 — Run every check against your data
│  ├── Each assertion is evaluated against every row in your DataFrame
│  ├── Checks run in parallel for speed (multiple checks at the same time)
│  └── Uses Great Expectations under the hood to validate
│
│  Step 3 — Save results to DQ Engine
│  ├── Each check's result (pass/fail/skipped) is sent to the DQ Engine
│  └── Results become visible in the DQ Engine web dashboard
│
│  Step 4 — Return summary to you
│  └── You get back a dictionary showing how many checks passed and failed
```

**What you get back:**

```python
{
    "sales_data_93": {
        "evaluated_expectations": 5,
        "successful_expectations": 4,
        "unsuccessful_expectations": 1,
        "success_percent": 80.0
    }
}
```

This tells you: 5 checks ran, 4 passed, 1 failed — 80% pass rate.

| Field | What It Means |
|-------|--------------|
| `evaluated_expectations` | Total number of checks that ran |
| `successful_expectations` | How many passed |
| `unsuccessful_expectations` | How many failed |
| `success_percent` | Pass rate (100.0 = all passed) |

**If all checks pass** → `success_percent` is `100.0` — your data is clean.

**If any check fails** → `unsuccessful_expectations` tells you how many failed. You can review the details in the DQ Engine dashboard to see exactly which column and which rule failed.

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