Data Quality Tasks

Data Quality Checks as Pipeline Tasks

A check task runs a query and fails if the answer is wrong, so the orchestrator's normal machinery (retries, downstream skipping, alerts) applies to data problems too. Airflow 129 's common-SQL provider (2.1.1 here) has check operators that work with any SQL connection; dags/booknest_quality.py checks the published mart and the fact table before a downstream finance extract:

dags/booknest_quality.py (excerpt): two check operatorsPython
    mart_day = SQLTableCheckOperator(
        task_id="mart_day_complete",
        conn_id="booknest_pg",
        table="mart.daily_genre_sales",
        partition_clause="sale_date = '{{ ds }}'",
        checks={"six_genres": {"check_statement": "COUNT(*) = 6"},
                "gross_positive": {"check_statement": "MIN(gross) > 0"}},
    )
    fact_columns = SQLColumnCheckOperator(
        task_id="fct_sales_columns",
        conn_id="booknest_pg",
        table="analytics.fct_sales",
        column_mapping={"gross_amount": {"null_check": {"equal_to": 0}, "min": {"geq_to": 0}},
                        "order_id": {"null_check": {"equal_to": 0}}},
    )

airflow/alerts/q1.sh ran the DAG for 28 June (published, all checks passed) and for 29 June, which the pipeline had not published yet. The operator turned both checks into one UNION ALL query and failed the task with each result:

Output of 51
Results:
[('six_genres', 0), ('gross_positive', 0)]
The following tests have failed:
   Check: six_genres,
   Check Values: {'check_statement': 'COUNT(*) = 6', 'result': '0', 'success': False}
,    Check: gross_positive,
   Check Values: {'check_statement': 'MIN(gross) > 0', 'result': '0', 'success': False}

finance_extract was marked upstream_failed and never ran. Keep such checks cheap (aggregates over one partition, not table scans), name each one after the rule it enforces, and decide per check whether it should block or only warn. For richer rule sets, run GX Core or Soda Core 2,433 as a task (Data Quality); note that Soda Core v4 is licensed under the Elastic License 2.0, not an open-source licence.