> ## Documentation Index
> Fetch the complete documentation index at: https://astronomer.io/docs/llms.txt
> Use this file to discover all available pages before exploring further.

# Data quality and Airflow

> Check the quality of your data using Airflow SQL check operators, pre-load checks, and third-party frameworks.

Data quality is fundamental to trustworthy analytics and AI. Silent data issues, like missing records, schema changes, or duplicate entries, can go undetected until dashboards break or models fail. By then, the damage is done.

This guide covers Dag-level checks with SQL check operators, pre-load checks against files in object storage, and third-party frameworks.

<CardGroup cols={3}>
  <Card title="Dag-level checks" icon="code" iconType="light" href="#dag-level-checks-with-airflow">
    SQL check operators that run as part of your pipeline, with full control over failure behavior and notifications.
  </Card>

  <Card title="Pre-load checks" icon="file-magnifying-glass" iconType="light" href="#pre-load-checks-with-the-analyticsoperator">
    The `AnalyticsOperator` runs SQL against files in object storage before they are loaded.
  </Card>

  <Card title="Third-party frameworks" icon="puzzle-piece" iconType="light" href="#third-party-frameworks">
    Great Expectations, dbt tests, and Soda Core for additional validation capabilities.
  </Card>
</CardGroup>

<Info>
  This guide uses Airflow 3.3 and version 2.1 of the Common SQL provider. If you use a different version, especially a different major version, some details might not apply to your setup.
</Info>

## Types of data quality challenges

* **Volume anomalies**: Unexpected spikes or drops in row counts.
* **Schema drift**: Column types changing without notice.
* **Completeness issues**: Null values in critical fields.
* **Duplicate records**: Compromising uniqueness constraints.
* **Business rule violations**: Negative amounts, invalid dates, orphaned records.

## Where to run your checks

Data quality checks can be run at multiple points in a pipeline. Often teams want two different checks: against files before they are loaded, and against tables after the load.

| | Pre-load checks | In-warehouse checks |
| - | - | - |
| **Runs against** | Files in object storage | Tables in a database or warehouse |
| **Runs when** | Before the load | During or after the load |
| **Query engine** | Apache DataFusion | Your database or warehouse |
| **Operators** | `AnalyticsOperator` | SQL check operators |
| **Failure behavior** | Returns query results, so a downstream task raises the exception | Fails the task directly (see the following note) |
| **Compute** | No warehouse compute, and the queried table does not have to exist | Warehouse compute |
| **Use for** | Row counts, null counts, and format validation on incoming files | Business rules, multi-table joins, and comparisons against previous runs |

<Info>
  **Dag-level checks do not have to stop your pipeline**

  You have several options to control failure behavior:

  1. **Trigger rules**: Use trigger rules like `all_done` or `none_failed_min_one_success` on downstream tasks to continue despite failed checks. See [Airflow trigger rules](/docs/airflow/airflow-trigger-rules).
  2. **Branching**: Use the `BranchSQLOperator` to route to a quarantine path instead of failing.
  3. **Shell exit code handling**: When calling third-party frameworks through `@task.bash`, append `|| true` or `|| exit 0` to prevent non-zero exit codes from failing the task.
  4. **Skip instead of fail**: Set `skip_on_exit_code` in `@task.bash(...)` to mark tasks as skipped rather than failed.

  **Important**: Even when a task fails but the Dag succeeds (through trigger rules), task-level `on_failure_callback` still fires, ensuring you can get notified about check failures without blocking your pipeline.
</Info>

## Choose a tool

Which tool you choose is determined by the needs and preferences of your organization. Astronomer recommends using Dag-level checks with SQL check operators if you want to:

* Write checks without needing to set up software in addition to Airflow.
* Write checks as Python dictionaries and in SQL.
* Use any SQL statement that returns a single row of booleans as a data quality check.
* Implement many different downstream dependencies depending on the outcome of different checks.
* Have full observability of which checks failed from within Airflow task logs, including the full SQL statements of failed checks.

Astronomer recommends using Dag-level checks with a data validation framework such as Great Expectations or Soda in the following circumstances:

* You want to collect the results of your data quality checks in a central place.
* You prefer to write checks in JSON (Great Expectations) or YAML (Soda).
* Most or all of your checks can be implemented by the predefined checks in the solution of your choice.
* You want to abstract your data quality checks from the Dag code.

In both cases, Astronomer recommends adding pre-load checks with the `AnalyticsOperator` to validate incoming files before they are loaded.

## Dag-level checks with Airflow

### When to use each operator

| Use case | Recommended operator |
| - | - |
| Primary key validation | `SQLColumnCheckOperator` |
| Null/duplicate checks | `SQLColumnCheckOperator` |
| Value range validation | `SQLColumnCheckOperator` |
| Row count thresholds | `SQLTableCheckOperator` |
| Cross-column business rules | `SQLTableCheckOperator` |
| Aggregate validations | `SQLTableCheckOperator` |
| Compare to expected value | `SQLValueCheckOperator` |
| Compare to historical data | `SQLIntervalCheckOperator` |
| Min/max threshold validation | `SQLThresholdCheckOperator` |
| Complex multi-table queries | `SQLCheckOperator` |

### SQL check operators

To access the SQL check operators, install the [Common SQL provider](https://airflow.apache.org/registry/providers/common-sql/):

```text wrap theme={null}
apache-airflow-providers-common-sql
```

Import and use them within your Dag:

```python wrap theme={null}
from airflow.providers.common.sql.operators.sql import (
    SQLColumnCheckOperator,
    SQLTableCheckOperator,
    SQLCheckOperator,
    SQLValueCheckOperator,
    SQLIntervalCheckOperator,
    SQLThresholdCheckOperator,
)
```

<AccordionGroup>
  <Accordion title="SQLCheckOperator">
    The `SQLCheckOperator` is the most generic check operator. It runs any SQL query and evaluates the result, giving you enough freedom to cover complex business rules. The check fails if any returned value evaluates to `False` in Python (for example, `0`, `None`, empty string).

    ```python wrap theme={null}
    _check_no_orphaned_payments = SQLCheckOperator(
        task_id="check_no_orphaned_payments",
        conn_id=_DB_CONN_ID,
        sql="""
            SELECT COUNT(*) = 0
            FROM payments p
            LEFT JOIN bookings b ON p.booking_id = b.booking_id
            WHERE b.booking_id IS NULL
        """,
    )
    ```
  </Accordion>

  <Accordion title="SQLColumnCheckOperator" defaultOpen={true}>
    The `SQLColumnCheckOperator` validates individual columns using built-in check types. Define a `column_mapping` dictionary to run multiple checks in a single task.

    **Built-in check types:**

    | Check | Description |
    | - | - |
    | `null_check` | Count of NULL values |
    | `unique_check` | Count of duplicate values |
    | `distinct_check` | Count of unique values |
    | `min` | Minimum value in column |
    | `max` | Maximum value in column |

    **Comparison options:**

    * `equal_to`, `greater_than`, `geq_to` (`>=`)
    * `less_than`, `leq_to` (`<=`)
    * `tolerance` (percentage threshold, as a fraction: `0.1` = 10%)

    ```python wrap theme={null}
    _check_columns = SQLColumnCheckOperator(
        task_id="check_non_promo_bookings_columns",
        conn_id=_DB_CONN_ID,
        table="bookings",
        partition_clause="promo_code IS NOT NULL",
        column_mapping={
            "booking_id": {
                "null_check": {"equal_to": 0},
                "unique_check": {"equal_to": 0},
            },
            "passengers": {
                "min": {"geq_to": 1},
                "max": {"leq_to": 10, "tolerance": 0.05},  # 5% tolerance
            },
        },
    )
    ```

    <Info>
      **partition\_clause**

      The optional `partition_clause` is an additional WHERE filter applied before the checks. It can be added at the operator level (partitions all checks), at the column level in the column mapping (partitions all checks for that column), or at the check level (partitions just that check).
    </Info>
  </Accordion>

  <Accordion title="SQLTableCheckOperator">
    The `SQLTableCheckOperator` runs custom SQL expressions against a given table that must evaluate to `true`. It is suited for business rules spanning multiple columns or requiring aggregations.

    ```python wrap theme={null}
    _check_report_business_rules = SQLTableCheckOperator(
        task_id="check_report_business_rules",
        conn_id=_DB_CONN_ID,
        table="daily_planet_report",
        checks={
            "net_fare_not_negative": {
                "check_statement": "total_net_fare_usd >= 0",
            },
            "discounts_leq_gross": {
                "check_statement": "total_discounts_usd <= total_gross_fare_usd",
            },
            "has_rows_for_today": {
                "check_statement": "COUNT(*) >= 1",
                "partition_clause": "report_date = '{{ ds }}'",
            },
        },
    )
    ```

    The `SQLTableCheckOperator` also supports an optional `partition_clause` on check level for an additional WHERE filter applied before the check.
  </Accordion>

  <Accordion title="SQLValueCheckOperator">
    Performs a simple value check by comparing a SQL result to an expected value (`pass_value`). The value can be of any type. For numerical values, you can set an additional tolerance percentage.

    ```python wrap theme={null}
    _check_planet_count = SQLValueCheckOperator(
        task_id="check_planet_count",
        conn_id=_DB_CONN_ID,
        sql="SELECT COUNT(*) FROM planets",
        pass_value=3,
        tolerance=0.1,  # 10% tolerance
    )
    ```
  </Accordion>

  <Accordion title="SQLIntervalCheckOperator">
    Verify that metrics defined as SQL expressions remain within tolerance compared to those from previous days (`days_back`). This utility helps track how values change over time and identify potential outliers.

    ```python wrap theme={null}
    _check_bookings_vs_last_week = SQLIntervalCheckOperator(
        task_id="check_bookings_vs_last_week",
        conn_id=_DB_CONN_ID,
        table="bookings",
        date_filter_column="CAST(booked_at AS DATE)",
        days_back=-7,
        ratio_formula="max_over_min",
        metrics_thresholds={"COUNT(*)": 3},  # max 3x deviation
    )
    ```

    <Info>
      **Defaults**

      The default for `days_back` is `-7`, and `ds` for the `date_filter_column`. Always set `date_filter_column` explicitly to your table's actual date column.
    </Info>
  </Accordion>

  <Accordion title="SQLThresholdCheckOperator">
    Performs a value check against a minimum and maximum threshold.

    ```python wrap theme={null}
    _check_avg_fare_in_range = SQLThresholdCheckOperator(
        task_id="check_avg_fare_in_range",
        conn_id=_DB_CONN_ID,
        sql="SELECT AVG(amount_usd) FROM payments",
        min_threshold=40000,
        max_threshold=80000,
    )
    ```

    <Info>
      **SQL expression thresholds**

      Thresholds can also be SQL expressions, not just numeric values. For example: `min_threshold="SELECT MIN(target_avg) FROM benchmarks"`.
    </Info>
  </Accordion>
</AccordionGroup>

### Notifications

Detecting data quality issues is only part of the story. The other part is raising awareness. To be notified when a data quality check fails, combine the `on_failure_callback` task parameter with [Airflow notifiers](/docs/airflow/error-notifications-in-airflow).

**Slack example:**

To use the `SlackNotifier`, install the following package:

```text wrap theme={null}
apache-airflow-providers-slack
```

In the templated text field, you can access the table through `task.table` and the actual quality issue through `exception`.

````python wrap theme={null}
from airflow.providers.slack.notifications.slack import SlackNotifier

SQLTableCheckOperator(
    task_id="business_rule_checks",
    on_failure_callback=SlackNotifier(
        slack_conn_id="slack_default",
        text="""
            Data quality checks failed for table: `{{ task.table }}`!
            ```
            {{ exception }}
            ```
        """,
        channel="#data-alerts",
    ),
    ...
)
````

<Info>
  **SlackNotifier**

  The `SlackNotifier` requires a properly configured Slack connection. In this case, the connection ID is `slack_default`.
</Info>

<Info>
  **AppriseNotifier**

  The `AppriseNotifier` supports 100+ notification services (Slack, Email, PagerDuty, Teams, etc.) through a unified interface. Install it using `apache-airflow-providers-apprise`.
</Info>

### Check patterns

Astronomer recommends running **column checks first** (field-level validation), followed by **table checks** (business logic), and then proceed with downstream processing.

Additionally, use **task groups** to organize your data quality checks and use the `default_args` parameter to configure notifications for all checks at once.

```python wrap theme={null}
@task_group(
    default_args={
        "on_failure_callback": notify_dq_failure,
    },
)
def data_quality_checks():
    _column_checks = SQLColumnCheckOperator(...)
    _table_checks = SQLTableCheckOperator(...)
```

## Pre-load checks with the AnalyticsOperator

The `AnalyticsOperator` runs SQL directly against files in object storage using Apache DataFusion. Use it to validate files before they are loaded into a warehouse.

The operator reads from Amazon S3 and the local filesystem, supports the Parquet, CSV, and Avro formats, and queries Apache Iceberg tables through a catalog.

To access the `AnalyticsOperator`, install the Common SQL provider with the `datafusion` extra:

```text wrap theme={null}
apache-airflow-providers-common-sql[datafusion]
```

### Define a datasource

A `DataSourceConfig` defines the location and format of the files the operator queries:

```python wrap theme={null}
from airflow.providers.common.sql.config import DataSourceConfig
from airflow.providers.common.sql.operators.analytics import AnalyticsOperator

_landing_bookings = DataSourceConfig(
    conn_id="aws_default",
    table_name="raw_bookings",
    uri="s3://daily-planet-landing/bookings/",
    format="parquet",
)
```

`table_name` is the identifier you reference in your queries. It does not have to exist in a database.

<Info>
  **Partitioned data**

  For partitioned data, pass the partition columns through `options`. For example: `options={"table_partition_cols": [("year", "integer")]}`.
</Info>

### Run the checks

Pass your datasource configs and a list of queries to the operator:

```python wrap theme={null}
_check_landing_zone = AnalyticsOperator(
    task_id="check_landing_zone",
    datasource_configs=[_landing_bookings],
    queries=[
        "SELECT COUNT(*) AS row_count FROM raw_bookings",
        "SELECT COUNT(*) AS null_ids FROM raw_bookings WHERE booking_id IS NULL",
    ],
    result_output_format="json",
)
```

### Fail the task on the results

The `AnalyticsOperator` runs your queries and returns the results. Unlike the SQL check operators, it does not evaluate a condition or fail. If you want to fail your Dag based on the check results, add a downstream task that reads the results and raises an exception:

```python wrap theme={null}
import json

from airflow.exceptions import AirflowException
from airflow.sdk import task


@task
def validate_landing_zone(raw: str) -> None:
    results = json.loads(raw)
    row_count = results[0]["data"][0]["row_count"]
    null_ids = results[1]["data"][0]["null_ids"]

    if row_count == 0 or null_ids > 0:
        raise AirflowException(
            f"Check failed: {row_count} rows, {null_ids} null booking_id values"
        )


validate_landing_zone(_check_landing_zone.output)
```

<Info>
  **Aggregate in SQL**

  `max_rows_check` defaults to 100. Aggregate in SQL rather than pulling rows into XCom. A check should return a number, not a dataset.
</Info>

## Third-party frameworks

### Great Expectations

The [airflow-provider-great-expectations](https://great-expectations.github.io/airflow-provider-great-expectations/latest/getting-started/) package provides operators for running Great Expectations validations directly in your Dags.

When deciding which operator fits your use case, consider:

1. **Where is your data?** In memory as a DataFrame, or in an external data source?
2. **Do you need to trigger actions?** Such as sending notifications or updating external systems based on validation results.
3. **What Data Context do you need?** Ephemeral for stateless validations, or persistent to track results over time.

| Scenario | Recommended operator |
| - | - |
| Data already in memory as Pandas or Spark DataFrame | `GXValidateDataFrameOperator` |
| Data in a database, warehouse, or file system | `GXValidateBatchOperator` |
| Need to trigger Slack notifications, emails, or other actions | `GXValidateCheckpointOperator` |
| Want full GX Core features with ValidationDefinitions | `GXValidateCheckpointOperator` |

See [Orchestrate Great Expectations with Airflow](/docs/airflow/airflow-great-expectations) to learn how to use these operators in your Dag.

### dbt tests

If you use dbt with Airflow through [Cosmos](/docs/airflow/airflow-dbt), use dbt's built-in testing framework:

* **Schema tests**: `unique`, `not_null`, `accepted_values`, `relationships`.
* **Custom tests**: SQL-based assertions for business-specific validation.

See [Orchestrate dbt Core with Airflow](/docs/airflow/airflow-dbt) for integration patterns.

### Soda Core

Run [Soda](/docs/airflow/soda-data-quality) checks using `@task.bash`:

```python wrap theme={null}
@task.bash
def soda_scan():
    return "soda scan -d snowflake -c soda_config.yml checks.yml"
```

## Quick reference

| Approach | Best for | Trade-offs |
| - | - | - |
| `SQLColumnCheckOperator` | Column-level validation (null, unique, min/max) | Requires code changes |
| `SQLTableCheckOperator` | Business rules, aggregations | Requires code changes |
| `SQLValueCheckOperator` | Compare to expected value with tolerance | Simple comparisons only |
| `SQLIntervalCheckOperator` | Compare to historical data | Requires consistent history |
| `SQLThresholdCheckOperator` | Min/max bounds validation | Simple bounds only |
| `SQLCheckOperator` | Complex multi-table queries | Most flexible, most verbose |
| `AnalyticsOperator` | Validating files in object storage before load | Returns results instead of failing, so it requires a downstream task to be a gate |
| `GXValidateDataFrameOperator` | In-memory DataFrame validation | Requires data in memory |
| `GXValidateBatchOperator` | Database/file validation using BatchDefinition | More setup than DataFrame |
| `GXValidateCheckpointOperator` | Full GX features with actions | Most configuration required |
| dbt tests | Model-layer validation | Requires dbt |
| Soda Core | Declarative YAML-based checks | Additional tool |

## Conclusion

Pre-load checks with the `AnalyticsOperator` validate files before they are loaded, without warehouse compute. Dag-level SQL check operators validate tables after the load and fail the task when a check does not pass.

Many teams should use both, and set the failure behavior for each check based on whether it should stop the pipeline.
