Köra SQL-arbetsflöden

Bygg datapipelines med Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Varför SQL i Airflow?

 

  • SQL-arbetsflöden är det vanligaste användningsfallet för Airflow
  • Låt databasen göra det tunga arbetet
  • Airflow styr när och var
  • Databasen hanterar hur

SQL-orkestrering i Airflow

Bygg datapipelines med 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", )

 

  • Fungerar med valfri databas som har en kompatibel Airflow-provider
  • PostgreSQL, Snowflake, BigQuery, DuckDB med flera
Bygg datapipelines med Airflow

Anslutningar

  • Lagra autentiseringsuppgifter utanför din kod
  • Var och en har ett ID, typ, värd, port och inloggning
  • Skapas på flera sätt: via gränssnittet, CLI, API eller miljövariabler
  • Refereras med conn_id i operatorer

Airflow Connections-gränssnitt

Bygg datapipelines med Airflow

Bygga en 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
        """,
    )
Bygg datapipelines med Airflow

Externa SQL-filer

@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",
    )
  • Ange template_searchpath@dag
  • Referera till filnamnet i sql
  • Håll skript och affärslogik utanför mappen dags/ 💡

Projektets filstruktur

Bygg datapipelines med Airflow

Jinja-mallar i SQL-filer

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 }} renderar det logiska datumet (YYYY-MM-DD)
  • DELETE sedan INSERT: det idempotenta mönstret från kapitel 2
  • Att köra om samma datum ger samma resultat
Bygg datapipelines med Airflow

params vs parameters

params (Jinja-rendering)

SQLExecuteQueryOperator(
    sql="""SELECT * FROM orders
           WHERE product = '{{ params.product }}'""",
    params={"product": user_input},
)
  • Värdet interpoleras i SQL-strängen
  • Sårbart för SQL-injektion

parameters (bindning på databasnivå)

SQLExecuteQueryOperator(
    sql="SELECT * FROM orders
         WHERE product = $product",
    parameters={"product": user_input},
)
  • Värdet skickas till databasdrivrutinen
  • Injektionssäkert: drivrutinen hanterar escapning
Bygg datapipelines med Airflow

Från lokal miljö till produktion med Astro

$$

  • Astronomer's Astro: hanterad Airflow-plattform

$$

  • Astro CLI: lokal Airflow med ett kommando

$$

  • Utveckla lokalt och driftsätt i produktion sömlöst

Astro CLI och plattform

Bygg datapipelines med Airflow

Nu kör vi en övning!

Bygg datapipelines med Airflow

Preparing Video For Download...