डेटा क्वालिटी चेक्स

Airflow के साथ Data Pipelines बनाना

Volker Janz

Senior Developer Advocate at Astronomer

क्वालिटी गेट्स क्यों ज़रूरी हैं

 

  • बिना चेक वाला पाइपलाइन ऐसा पाइपलाइन है जिस पर आप भरोसा नहीं कर सकते
  • एक अपस्ट्रीम बदलाव से nulls या निगेटिव वैल्यूज़ आ जाती हैं
  • हर डाउनस्ट्रीम डैशबोर्ड पर गलत नंबर दिखते हैं
  • क्वालिटी गेट्स समस्याएँ सोर्स पर ही पकड़ लेते हैं

खराब डेटा नीचे की ओर बह रहा है

Airflow के साथ Data Pipelines बनाना

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 के साथ Data Pipelines बनाना

कॉलम चेक के प्रकार

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 के साथ Data Pipelines बनाना

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 अभिव्यक्ति जो true या false देती है
Airflow के साथ Data Pipelines बनाना

क्वालिटी चेक्स की चेइनिंग

@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 के साथ Data Pipelines बनाना

कॉलम बनाम टेबल चेक्स

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}},
    },
)
  • क्या nulls हैं?
  • क्या वैल्यूज़ रेंज में हैं?

SQLTableCheckOperator

check_table = SQLTableCheckOperator(
    task_id="check_table",
    conn_id="duckdb_analytics",
    table="daily_summary",
    checks={
        "row_count_check": {
            "check_statement": "COUNT(*) > 0",
        },
    },
)
  • क्या टेबल में rows हैं?
  • क्या टोटल expected bounds में है?
Airflow के साथ Data Pipelines बनाना

Astro के साथ स्केल पर क्वालिटी मॉनिटरिंग

$$

  • Astronomer का Astro Observe पाइपलाइन हेल्थ मॉनिटर करने में मदद करता है
  • डेडलाइन से पहले अलर्ट्स के साथ SLA ट्रैकिंग
  • Dags और टेबल्स में पाइपलाइन लिनिएज
  • बड़े पैमाने पर डेटा क्वालिटी मॉनिटरिंग
  • AI-आधारित root cause analysis

Astro Observe लिनिएज

Airflow के साथ Data Pipelines बनाना

अभ्यास करते हैं!

Airflow के साथ Data Pipelines बनाना

Preparing Video For Download...