Uruchamianie zadań SQL

Budowanie potoków danych z Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Po co SQL w Airflow?

 

  • Zadania SQL to najczęstszy przypadek użycia Airflow
  • Pozwól, żeby baza danych wykonała ciężką pracę
  • Airflow zarządza kiedy i gdzie
  • Baza danych odpowiada za jak

Orkiestracja SQL w Airflow

Budowanie potoków danych z 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", )

 

  • Działa z każdą bazą danych, która ma zgodny provider Airflow
  • PostgreSQL, Snowflake, BigQuery, DuckDB i inne
Budowanie potoków danych z Airflow

Connections

  • Przechowuj dane uwierzytelniające poza kodem
  • Każde połączenie ma ID, typ, host, port i login
  • Można je tworzyć na różne sposoby: przez UI, CLI, API lub zmienne środowiskowe
  • Odwołuj się przez conn_id w operatorach

Interfejs Connections w Airflow

Budowanie potoków danych z Airflow

Budowanie pipeline'u 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
        """,
    )
Budowanie potoków danych z Airflow

Zewnętrzne pliki 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",
    )
  • Ustaw template_searchpath w @dag
  • Podaj nazwę pliku w sql
  • Trzymaj skrypty i logikę biznesową poza folderem dags/ 💡

Struktura plików projektu

Budowanie potoków danych z Airflow

Szablony Jinja w plikach 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 }} wstawia datę logiczną (YYYY-MM-DD)
  • DELETE + INSERT: wzorzec idempotentny z rozdziału 2
  • Ponowne uruchomienie dla tej samej daty daje ten sam wynik
Budowanie potoków danych z Airflow

params vs parameters

params (rendering Jinja)

SQLExecuteQueryOperator(
    sql="""SELECT * FROM orders
           WHERE product = '{{ params.product }}'""",
    params={"product": user_input},
)
  • Wartość wstawiana bezpośrednio do SQL
  • Podatne na SQL injection

parameters (bindowanie na poziomie bazy)

SQLExecuteQueryOperator(
    sql="SELECT * FROM orders
         WHERE product = $product",
    parameters={"product": user_input},
)
  • Wartość przekazywana do sterownika bazy danych
  • Bezpieczne: sterownik obsługuje escaping
Budowanie potoków danych z Airflow

Od środowiska lokalnego do produkcji z Astro

$$

  • Astro od Astronomer: zarządzana platforma Airflow

$$

  • Astro CLI: lokalny Airflow jednym poleceniem

$$

  • Pracuj lokalnie, wdrażaj na produkcję bez komplikacji

Astro CLI i platforma

Budowanie potoków danych z Airflow

Czas na praktykę!

Budowanie potoków danych z Airflow

Preparing Video For Download...