SQL ワークロードの実行

Airflow によるデータパイプラインの構築

Volker Janz

Senior Developer Advocate at Astronomer

Airflow で SQL を使う理由

 

  • SQL ワークロードは Airflow の最も一般的なユースケースです
  • 重い処理はデータベースに任せましょう
  • Airflow はいつどこで実行するかを管理します
  • どのように処理するかはデータベースが担います

AirflowにおけるSQLオーケストレーション

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

 

  • 互換性のある Airflow プロバイダーがあれば任意のデータベースで動作します
  • PostgreSQL、Snowflake、BigQuery、DuckDB などに対応
Airflow によるデータパイプラインの構築

コネクション

  • 認証情報をコードの外部に保管します
  • それぞれに ID、タイプ、ホスト、ポート、ログイン情報があります
  • UI、CLI、API、環境変数など複数の方法で作成できます
  • オペレーターから conn_id で参照します

Airflow Connections UI

Airflow によるデータパイプラインの構築

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
        """,
    )
Airflow によるデータパイプラインの構築

外部 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",
    )
  • @dagtemplate_searchpath を設定します
  • sql にはファイル名を指定します
  • スクリプトやビジネスロジックは dags/ フォルダーの外に置きましょう 💡

プロジェクトのファイル構成

Airflow によるデータパイプラインの構築

SQL ファイルでの Jinja テンプレート

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 }}論理日付YYYY-MM-DD)に展開されます
  • DELETE してから INSERT する:第2章で紹介した冪等パターンです
  • 同じ日付を再実行しても同じ結果が得られます
Airflow によるデータパイプラインの構築

params と parameters の違い

params(Jinja レンダリング)

SQLExecuteQueryOperator(
    sql="""SELECT * FROM orders
           WHERE product = '{{ params.product }}'""",
    params={"product": user_input},
)
  • 値が SQL 文字列に展開されます
  • SQL インジェクションの危険があります

parameters(DB レベルのバインド)

SQLExecuteQueryOperator(
    sql="SELECT * FROM orders
         WHERE product = $product",
    parameters={"product": user_input},
)
  • 値はデータベースドライバーに渡されます
  • インジェクション安全:エスケープはドライバーが処理します
Airflow によるデータパイプラインの構築

Astro でローカルから本番環境へ

$$

  • Astronomer の Astro:マネージド Airflow プラットフォーム

$$

  • Astro CLI:1つのコマンドでローカル Airflow を起動

$$

  • ローカルで開発し、本番環境へシームレスにデプロイ

Astro CLI とプラットフォーム

Airflow によるデータパイプラインの構築

練習しましょう!

Airflow によるデータパイプラインの構築

Preparing Video For Download...