SQL-workloads uitvoeren

Data-pijplijnen bouwen met Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Waarom SQL in Airflow?

 

  • SQL-workloads zijn de meest voorkomende Airflow-use case
  • Laat de database het zware werk doen
  • Airflow orkestreert wanneer en waar
  • De database bepaalt het hoe

SQL-orkestratie in Airflow

Data-pijplijnen bouwen met 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", )

 

  • Werkt met elke database met een compatibele Airflow-provider
  • PostgreSQL, Snowflake, BigQuery, DuckDB en meer
Data-pijplijnen bouwen met Airflow

Connections

  • Sla referenties buiten je code op
  • Elk heeft een ID, type, host, poort, login
  • Aan te maken via meerdere manieren: UI, CLI, API of omgevingsvariabelen
  • Verwijs ernaar met conn_id in operators

Airflow Connections-UI

Data-pijplijnen bouwen met Airflow

Een SQL-pijplijn bouwen

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-pijplijnen bouwen met Airflow

Externe SQL-bestanden

@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",
    )
  • Zet template_searchpath op @dag
  • Verwijs in sql naar de bestandsnaam
  • Houd scripts en businesslogica buiten de map dags/ 💡

Projectmappenstructuur

Data-pijplijnen bouwen met Airflow

Jinja-templates in SQL-bestanden

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 de logische datum (YYYY-MM-DD)
  • Eerst DELETE, dan INSERT: het idempotente patroon uit hoofdstuk 2
  • Opnieuw draaien voor dezelfde datum geeft hetzelfde resultaat
Data-pijplijnen bouwen met Airflow

params vs parameters

params (Jinja-rendering)

SQLExecuteQueryOperator(
    sql="""SELECT * FROM orders
           WHERE product = '{{ params.product }}'""",
    params={"product": user_input},
)
  • Waarde geïnterpoleerd in de SQL-string
  • Kwetsbaar voor SQL-injectie

parameters (DB-level binding)

SQLExecuteQueryOperator(
    sql="SELECT * FROM orders
         WHERE product = $product",
    parameters={"product": user_input},
)
  • Waarde doorgegeven aan de databasedriver
  • Injectie-veilig: driver verzorgt escaping
Data-pijplijnen bouwen met Airflow

Van lokaal naar productie met Astro

$$

  • Astronomer's Astro: beheerd Airflow-platform

$$

  • Astro CLI: lokale Airflow met één commando

$$

  • Ontwikkel lokaal, deploy naar productie zonder gedoe

Astro CLI en platform

Data-pijplijnen bouwen met Airflow

Laten we oefenen!

Data-pijplijnen bouwen met Airflow

Preparing Video For Download...