TaskFlow API

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

Volker Janz

Senior Developer Advocate at Astronomer

従来のアプローチ

def extract_data():
    return {"users": 150, "events": 4200}

with DAG("etl_pipeline") as dag: t1 = PythonOperator(task_id="extract", python_callable=extract_data) t2 = PythonOperator(task_id="summary", python_callable=print_summary)
t1 >> t2
Airflow によるデータパイプラインの構築

TaskFlow のアプローチ

from airflow.sdk import dag, task

@dag def etl_pipeline(): @task def extract_data(): return {"users": 150}
data = extract_data() print_summary(data)
Airflow によるデータパイプラインの構築

TaskFlow を使う理由

 

  • 定型コードが少ない:デコレータがオペレーターインスタンスを置き換えます
  • 依存関係が自動化:戻り値でタスクが自動的に連結されます
  • 可読性が高い:DAG が Python スクリプトのように読めます
  • クラシックオペレーターはプロバイダー統合でも引き続き利用可能です

TaskFlow API - ビジュアル

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

Airflow UI:グリッドビュー

Airflow グリッドビュー

 

  • 各列は DAG の実行を表します
  • 各行はタスクを表します
  • 色でステータスを確認できます: = 成功、 = 失敗
  • マスをクリックするとログと詳細を確認できます
Airflow によるデータパイプラインの構築

Airflow UI:グラフビュー

Airflow グリッドビュー

 

  • 1 回の DAG 実行に焦点を当てます
  • 依存関係と並列実行を視覚的に確認できます
  • 設定から表示の詳細を調整できます
Airflow によるデータパイプラインの構築

Airflow UI:DAG バージョン管理

Airflow UI の DAG バージョン表示

 

  • 構造の変更を自動的に追跡します
  • 各実行はその時点で有効だったバージョンに紐づきます
  • バージョンは DAG 詳細パネルで確認できます
Airflow によるデータパイプラインの構築

練習しましょう!

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

Preparing Video For Download...