SQL iş yüklerini çalıştırma

Airflow ile Veri İş Hatları Oluşturma

Volker Janz

Senior Developer Advocate at Astronomer

Neden Airflow'da SQL?

 

  • SQL iş yükleri en yaygın Airflow kullanım alanıdır
  • Veritabanı ağır işi yapsın
  • Airflow ne zaman ve neredeyi orkestre eder
  • Nasılını veritabanı halleder

Airflow'da SQL orkestrasyonu

Airflow ile Veri İş Hatları Oluşturma

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

 

  • Uyumlu Airflow provider'ı olan her veritabanıyla çalışır
  • PostgreSQL, Snowflake, BigQuery, DuckDB ve daha fazlası
Airflow ile Veri İş Hatları Oluşturma

Connections

  • Kimlik bilgilerini kodunun dışında tut
  • Her birinin ID, tür, host, port, giriş bilgisi vardır
  • Birden çok yolla oluşturulur: UI, CLI, API veya ortam değişkenleri
  • Operatörlerde conn_id ile referans ver

Airflow Connections arayüzü

Airflow ile Veri İş Hatları Oluşturma

Bir SQL veri hattı oluşturma

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
        """,
    )
Airflow ile Veri İş Hatları Oluşturma

Harici SQL dosyaları

@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",
    )
  • @dag üzerinde template_searchpath ayarla
  • sql içinde dosya adını referans ver
  • Betikleri ve iş mantığını dags/ klasörü dışında tut 💡

Proje dosya yapısı

Airflow ile Veri İş Hatları Oluşturma

SQL dosyalarında Jinja şablonları

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 }} mantıksal tarihi (YYYY-MM-DD) üretir
  • Önce DELETE sonra INSERT: Bölüm 2'den idempotent desen
  • Aynı tarihi tekrar çalıştırmak aynı sonucu verir
Airflow ile Veri İş Hatları Oluşturma

params ve parameters

params (Jinja rendering)

SQLExecuteQueryOperator(
    sql="""SELECT * FROM orders
           WHERE product = '{{ params.product }}'""",
    params={"product": user_input},
)
  • Değer SQL dizesine eklenir
  • SQL injection'a karşı hassastır

parameters (VT düzeyinde bağlama)

SQLExecuteQueryOperator(
    sql="SELECT * FROM orders
         WHERE product = $product",
    parameters={"product": user_input},
)
  • Değer veritabanı sürücüsüne iletilir
  • Enjeksiyona dayanıklı: kaçışlamayı sürücü yapar
Airflow ile Veri İş Hatları Oluşturma

Astro ile yerelden production'a

$$

  • Astronomer's Astro: yönetilen Airflow platformu

$$

  • Astro CLI: tek komutla yerelde Airflow

$$

  • Yerelde geliştir, sorunsuzca production'a dağıt

Astro CLI ve platform

Airflow ile Veri İş Hatları Oluşturma

Hadi pratik yapalım!

Airflow ile Veri İş Hatları Oluşturma

Preparing Video For Download...