Rularea workload-urilor SQL

Construirea pipeline-urilor de date cu Airflow

Volker Janz

Senior Developer Advocate at Astronomer

De ce SQL în Airflow?

 

  • Workload-urile SQL sunt cel mai frecvent caz de utilizare Airflow
  • Lasă baza de date să facă munca grea
  • Airflow orchestrează când și unde
  • Baza de date gestionează cum

Orchestrarea SQL în Airflow

Construirea pipeline-urilor de date cu 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", )

 

  • Funcționează cu orice bază de date care are un provider Airflow compatibil
  • PostgreSQL, Snowflake, BigQuery, DuckDB și altele
Construirea pipeline-urilor de date cu Airflow

Conexiuni

  • Stochează credențialele în afara codului
  • Fiecare are un ID, tip, host, port, login
  • Se creează în mai multe moduri: UI, CLI, API sau variabile de mediu
  • Referite prin conn_id în operatori

Interfața Connections din Airflow

Construirea pipeline-urilor de date cu Airflow

Construirea unui 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
        """,
    )
Construirea pipeline-urilor de date cu Airflow

Fișiere SQL externe

@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",
    )
  • Setează template_searchpath pe @dag
  • Referă numele fișierului în sql
  • Păstrează scripturile și logica de business în afara folderului dags/ 💡

Structura fișierelor proiectului

Construirea pipeline-urilor de date cu Airflow

Șabloane Jinja în fișiere 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 }} redă data logică (YYYY-MM-DD)
  • DELETE urmat de INSERT: tiparului idempotent din Capitolul 2
  • Rulând aceeași dată se obține același rezultat
Construirea pipeline-urilor de date cu Airflow

params vs parameters

params (redare Jinja)

SQLExecuteQueryOperator(
    sql="""SELECT * FROM orders
           WHERE product = '{{ params.product }}'""",
    params={"product": user_input},
)
  • Valoarea este interpolată în șirul SQL
  • Vulnerabil la injecție SQL

parameters (binding la nivel de bază de date)

SQLExecuteQueryOperator(
    sql="SELECT * FROM orders
         WHERE product = $product",
    parameters={"product": user_input},
)
  • Valoarea este transmisă driverului bazei de date
  • Sigur față de injecție: driverul gestionează escaparea
Construirea pipeline-urilor de date cu Airflow

De la local la producție cu Astro

$$

  • Astro de la Astronomer: platformă Airflow gestionată

$$

  • Astro CLI: Airflow local dintr-o singură comandă

$$

  • Dezvoltare locală, deploy în producție fără probleme

Astro CLI și platforma

Construirea pipeline-urilor de date cu Airflow

Hai să exersăm!

Construirea pipeline-urilor de date cu Airflow

Preparing Video For Download...