Esecuzione di carichi di lavoro SQL

Creare data pipeline con Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Perché SQL in Airflow?

 

  • I carichi di lavoro SQL sono il caso d'uso più comune in Airflow
  • Lascia che sia il database a fare il grosso
  • Airflow orchestra quando e dove
  • Il database gestisce il come

Orchestrazione SQL in Airflow

Creare data pipeline con Airflow

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", )

 

  • Funziona con qualsiasi database che abbia un provider Airflow compatibile
  • PostgreSQL, Snowflake, BigQuery, DuckDB e altri
Creare data pipeline con Airflow

Connessioni

  • Conserva le credenziali fuori dal codice
  • Ognuna ha un ID, tipo, host, porta, login
  • Create in più modi: UI, CLI, API o variabili d'ambiente
  • Riferisciti con conn_id negli operatori

Interfaccia Connections di Airflow

Creare data pipeline con Airflow

Creare una pipeline 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
        """,
    )
Creare data pipeline con Airflow

File SQL esterni

@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",
    )
  • Imposta template_searchpath su @dag
  • Indica il nome file in sql
  • Tieni script e logica di business fuori dalla cartella dags/ 💡

Struttura dei file di progetto

Creare data pipeline con Airflow

Template Jinja nei file SQL

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 }} restituisce la data logica (YYYY-MM-DD)
  • Prima DELETE poi INSERT: il pattern idempotente del Capitolo 2
  • Rieseguire la stessa data produce lo stesso risultato
Creare data pipeline con Airflow

params vs parameters

params (rendering Jinja)

SQLExecuteQueryOperator(
    sql="""SELECT * FROM orders
           WHERE product = '{{ params.product }}'""",
    params={"product": user_input},
)
  • Valore interpolato nella stringa SQL
  • Vulnerabile a SQL injection

parameters (binding a livello DB)

SQLExecuteQueryOperator(
    sql="SELECT * FROM orders
         WHERE product = $product",
    parameters={"product": user_input},
)
  • Valore passato al driver del database
  • Sicuro contro injection: il driver gestisce l'escaping
Creare data pipeline con Airflow

Dal locale alla produzione con Astro

$$

  • Astronomer's Astro: piattaforma Airflow gestita

$$

  • Astro CLI: Airflow locale con un solo comando

$$

  • Sviluppa in locale e distribuisci in produzione senza attriti

Astro CLI e piattaforma

Creare data pipeline con Airflow

Esercitiamoci!

Creare data pipeline con Airflow

Preparing Video For Download...