SQL वर्कलोड चलाना

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

Volker Janz

Senior Developer Advocate at Astronomer

Airflow में SQL क्यों?

 

  • SQL वर्कलोड Airflow का सबसे आम उपयोग केस है
  • भारी काम डेटाबेस को करने दें
  • Airflow तय करता है कब और कहाँ
  • कैसे का काम डेटाबेस संभालता है

Airflow में SQL ऑर्केस्ट्रेशन

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

SQLExecuteQueryOperator

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

SQLExecuteQueryOperator( task_id="load_sales", conn_id="duckdb_analytics", sql="INSERT INTO sales SELECT * FROM staging", )

 

  • ऐसे किसी भी डेटाबेस के साथ काम करता है जिसके लिए संगत Airflow प्रोवाइडर हो
  • PostgreSQL, Snowflake, BigQuery, DuckDB, आदि
Airflow के साथ Data Pipelines बनाना

Connections

  • क्रेडेंशियल्स को कोड से बाहर रखें
  • हर कनेक्शन का ID, type, host, port, login होता है
  • कई तरीकों से बनाएँ: UI, CLI, API, या environment variables
  • ऑपरेटर्स में conn_id से रेफ़र करें

Airflow Connections UI

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

SQL पाइपलाइन बनाना

from airflow.sdk import dag, task
from airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperator

@dag(schedule="@daily")
def sales_pipeline():

    aggregate = SQLExecuteQueryOperator(
        task_id="aggregate_daily_sales",
        conn_id="duckdb_analytics",
        sql="""
            INSERT INTO daily_summary (order_date, total_orders, total_revenue)
            SELECT order_date, COUNT(*), SUM(revenue)
            FROM raw_orders GROUP BY order_date
        """,
    )
Airflow के साथ Data Pipelines बनाना

External SQL फ़ाइलें

@dag(
    schedule="@daily",
    template_searchpath="/path/to/include/sql",
)
def sales_pipeline():
    aggregate = SQLExecuteQueryOperator(
        task_id="aggregate_daily_sales",
        conn_id="duckdb_analytics",
        sql="aggregate_sales.sql",
    )
  • @dag पर template_searchpath सेट करें
  • sql में फ़ाइलनाम रेफ़र करें
  • स्क्रिप्ट्स और बिज़नेस लॉजिक को dags/ फ़ोल्डर से बाहर रखें 💡

प्रोजेक्ट फ़ाइल स्ट्रक्चर

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

SQL फ़ाइलों में Jinja टेम्पलेट्स

DELETE FROM daily_summary WHERE order_date = '{{ ds }}';

INSERT INTO daily_summary (order_date, total_orders, total_revenue)
SELECT order_date, COUNT(*), SUM(revenue)
FROM raw_orders
WHERE order_date = '{{ ds }}'
GROUP BY order_date;

 

  • {{ ds }} logical date (YYYY-MM-DD) रेंडर करता है
  • पहले DELETE फिर INSERT: Chapter 2 का idempotent पैटर्न
  • उसी तारीख को दोबारा चलाने पर एक जैसा परिणाम मिलता है
Airflow के साथ Data Pipelines बनाना

params बनाम parameters

params (Jinja rendering)

SQLExecuteQueryOperator(
    sql="""SELECT * FROM orders
           WHERE product = '{{ params.product }}'""",
    params={"product": user_input},
)
  • वैल्यू SQL स्ट्रिंग में interpolate होती है
  • SQL injection का जोखिम

parameters (DB-स्तर binding)

SQLExecuteQueryOperator(
    sql="SELECT * FROM orders
         WHERE product = $product",
    parameters={"product": user_input},
)
  • वैल्यू database driver को पास होती है
  • Injection-safe: ड्राइवर escaping संभालता है
Airflow के साथ Data Pipelines बनाना

Astro के साथ लोकल से प्रोडक्शन तक

$$

  • Astronomer का Astro: managed Airflow प्लेटफ़ॉर्म

$$

  • Astro CLI: एक कमांड में लोकल Airflow

$$

  • लोकली डेवलप करें, बिना रुकावट production में deploy करें

Astro CLI और प्लेटफ़ॉर्म

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

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

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

Preparing Video For Download...