SQL-Workloads ausführen

Data-Pipelines mit Airflow aufbauen

Volker Janz

Senior Developer Advocate at Astronomer

Warum SQL in Airflow?

 

  • SQL-Workloads sind der häufigste Airflow-Einsatzfall
  • Lass die Datenbank die schwere Arbeit machen
  • Airflow steuert wann und wo
  • Die Datenbank kümmert sich um das wie

SQL-Orchestrierung in Airflow

Data-Pipelines mit Airflow aufbauen

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

 

  • Funktioniert mit jeder Datenbank, die einen kompatiblen Airflow-Provider hat
  • PostgreSQL, Snowflake, BigQuery, DuckDB und mehr
Data-Pipelines mit Airflow aufbauen

Verbindungen

  • Zugangsdaten außerhalb deines Codes speichern
  • Jede Verbindung hat ID, Typ, Host, Port, Login
  • Auf viele Arten anlegen: UI, CLI, API oder Umgebungsvariablen
  • In Operatoren per conn_id referenzieren

Airflow Connections UI

Data-Pipelines mit Airflow aufbauen

Eine SQL-Pipeline bauen

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
        """,
    )
Data-Pipelines mit Airflow aufbauen

Externe SQL-Dateien

@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",
    )
  • Setze template_searchpath auf @dag
  • Referenziere den Dateinamen in sql
  • Halte Skripte und Business-Logik außerhalb des Ordners dags/ 💡

Projektverzeichnisstruktur

Data-Pipelines mit Airflow aufbauen

Jinja-Templates in SQL-Dateien

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 }} rendert das logische Datum (YYYY-MM-DD)
  • DELETE dann INSERT: das idempotente Muster aus Kapitel 2
  • Erneutes Ausführen für dasselbe Datum liefert dasselbe Ergebnis
Data-Pipelines mit Airflow aufbauen

params vs. parameters

params (Jinja-Rendering)

SQLExecuteQueryOperator(
    sql="""SELECT * FROM orders
           WHERE product = '{{ params.product }}'""",
    params={"product": user_input},
)
  • Wert wird in den SQL-String eingesetzt
  • Anfällig für SQL-Injection

parameters (DB-seitiges Binding)

SQLExecuteQueryOperator(
    sql="SELECT * FROM orders
         WHERE product = $product",
    parameters={"product": user_input},
)
  • Wert wird an den Datenbanktreiber übergeben
  • Injection-sicher: Treiber übernimmt Escaping
Data-Pipelines mit Airflow aufbauen

Von lokal bis Produktion mit Astro

$$

  • Astronomer's Astro: verwaltete Airflow-Plattform

$$

  • Astro CLI: lokales Airflow mit einem Befehl

$$

  • Lokal entwickeln, nahtlos in Produktion deployen

Astro CLI und Plattform

Data-Pipelines mit Airflow aufbauen

Lass uns üben!

Data-Pipelines mit Airflow aufbauen

Preparing Video For Download...