Spouštění SQL úloh

Tvorba datových pipeline s Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Proč SQL v Airflow?

 

  • SQL úlohy jsou nejčastějším případem použití Airflow
  • Nech databázi udělat těžkou práci
  • Airflow řídí kdy a kde
  • Databáze řeší jak

Orchestrace SQL v Airflow

Tvorba datových pipeline s 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", )

 

  • Funguje s jakoukoli databází, která má kompatibilního Airflow providera
  • PostgreSQL, Snowflake, BigQuery, DuckDB a další
Tvorba datových pipeline s Airflow

Connections

  • Ukládej přihlašovací údaje mimo kód
  • Každé připojení má ID, typ, host, port, login
  • Vytváření více způsoby: UI, CLI, API nebo proměnné prostředí
  • Odkazuj se na ně pomocí conn_id v operátorech

Airflow Connections UI

Tvorba datových pipeline s Airflow

Tvorba SQL pipeline

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
        """,
    )
Tvorba datových pipeline s Airflow

Externí SQL soubory

@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",
    )
  • Nastav template_searchpath na @dag
  • V sql uveď název souboru
  • Skripty a business logiku uchovávej mimo složku dags/ 💡

Struktura souborů projektu

Tvorba datových pipeline s Airflow

Jinja šablony v SQL souborech

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 }} vykreslí logické datum (YYYY-MM-DD)
  • DELETE a INSERT: idempotentní vzor z kapitoly 2
  • Opakované spuštění se stejným datem dává stejný výsledek
Tvorba datových pipeline s Airflow

params vs parameters

params (Jinja rendering)

SQLExecuteQueryOperator(
    sql="""SELECT * FROM orders
           WHERE product = '{{ params.product }}'""",
    params={"product": user_input},
)
  • Hodnota interpolována do SQL řetězce
  • Zranitelné vůči SQL injection

parameters (vazba na úrovni DB)

SQLExecuteQueryOperator(
    sql="SELECT * FROM orders
         WHERE product = $product",
    parameters={"product": user_input},
)
  • Hodnota předána databázovému driveru
  • Bezpečné vůči injection: driver řeší escapování
Tvorba datových pipeline s Airflow

Z lokálního prostředí do produkce s Astro

$$

  • Astronomer's Astro: spravovaná platforma pro Airflow

$$

  • Astro CLI: lokální Airflow jedním příkazem

$$

  • Vyvíjej lokálně, nasazuj do produkce bez komplikací

Astro CLI a platforma

Tvorba datových pipeline s Airflow

Pojďme cvičit!

Tvorba datových pipeline s Airflow

Preparing Video For Download...