データ品質チェック

Airflow によるデータパイプラインの構築

Volker Janz

Senior Developer Advocate at Astronomer

品質ゲートが重要な理由

 

  • チェックのないパイプラインは信頼できません
  • 上流の変更がNULL値負の値を生じさせます
  • 下流のダッシュボードはすべて誤った数値を表示します
  • 品質ゲートは問題を発生源で検出します

不良データが下流に流れる様子

Airflow によるデータパイプラインの構築

SQLColumnCheckOperator

from airflow.providers.common.sql.operators.sql import
    (SQLExecuteQueryOperator, SQLColumnCheckOperator)

check_columns = SQLColumnCheckOperator( task_id="check_columns", conn_id="duckdb_analytics", table="daily_summary", column_mapping={ "total_revenue": {"min": {"greater_than": 0}}, "total_orders": {"null_check": {"equal_to": 0}}, }, )
Airflow によるデータパイプラインの構築

カラムチェックの種類

column_mapping={
    "total_revenue": {"min": {"greater_than": 0}},
    "total_orders": {"null_check": {"equal_to": 0}},
    "order_date": {"distinct_check": {"greater_than": 0}},
}
  • min: カラムの最小値を検証します
  • null_check: NULL値をカウントします(NULL なしの場合は equal_to: 0 を指定)
  • distinct_check: 個別値の数をカウントします
  • 利用可能なチェック: null_check, unique_check, distinct_check, min, max
  • 利用可能な比較演算子: equal_to, greater_than, geq_to, less_than, leq_to
Airflow によるデータパイプラインの構築

SQLTableCheckOperator

from airflow.providers.common.sql.operators.sql import SQLTableCheckOperator

check_table = SQLTableCheckOperator( task_id="check_table", conn_id="duckdb_analytics", table="daily_summary", checks={ "row_count_check": { "check_statement": "COUNT(*) > 0", }, }, )
  • テーブル全体のプロパティを検証します
  • check_statement: 真または偽に評価される任意のSQL式を指定します
Airflow によるデータパイプラインの構築

品質チェックの連結

@dag(schedule="@daily", template_searchpath="/path/to/sql")
def sales_pipeline():

    load = SQLExecuteQueryOperator(
        task_id="load_daily_sales",
        conn_id="duckdb_analytics",
        sql="aggregate_sales.sql",
    )

    check_columns = SQLColumnCheckOperator(...)
    check_table = SQLTableCheckOperator(...)

    load >> check_columns >> check_table
  • まずロードし、次に検証します
  • いずれかのチェックが失敗すると、パイプラインは停止します
  • 下流の利用者に不良データが渡ることはありません
Airflow によるデータパイプラインの構築

カラムチェックとテーブルチェックの比較

SQLColumnCheckOperator

check_columns = SQLColumnCheckOperator(
    task_id="check_columns",
    conn_id="duckdb_analytics",
    table="daily_summary",
    column_mapping={
        "total_revenue": {"min": {"greater_than": 0}},
        "total_orders": {"null_check": {"equal_to": 0}},
    },
)
  • NULL値はありますか?
  • 値は範囲内に収まっていますか?

SQLTableCheckOperator

check_table = SQLTableCheckOperator(
    task_id="check_table",
    conn_id="duckdb_analytics",
    table="daily_summary",
    checks={
        "row_count_check": {
            "check_statement": "COUNT(*) > 0",
        },
    },
)
  • テーブルに行が存在しますか?
  • 合計値は想定範囲内ですか?
Airflow によるデータパイプラインの構築

Astro を使った大規模な品質監視

$$

  • Astronomer の Astro Observe でパイプラインの状態を監視できます
  • 期限前にアラートを通知する SLA トラッキング
  • DAG とテーブルをまたぐパイプラインリネージ
  • 大規模なデータ品質モニタリング
  • AI による根本原因分析

Astro Observe のリネージ画面

Airflow によるデータパイプラインの構築

練習しましょう!

Airflow によるデータパイプラインの構築

Preparing Video For Download...