Airflow コアコンセプト

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

Volker Janz

Senior Developer Advocate at Astronomer

インストラクター紹介

  Volker Janz のプロフィール写真

$$

Volker Janz

$$

  • Senior Developer Advocate、Astronomer
  • データエンジニアとしてゲーム業界で 14年以上 の経験
  • Airflow バージョン 1.x から使用
  • Data Engineer Things のスピーカー、メンター、ニュースレター担当
Airflow によるデータパイプラインの構築

学習内容

 

  • TaskFlow API で DAG を作成する
  • タスクマッピングとアセットベースのスケジューリングで動的なワークフローを構築する
  • リトライとコールバックで失敗を処理する
  • Airflow を通じて SQL ワークロードを実行する

各チャプターの内容を示す図

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

前提知識

 

$$

  • DAG、タスク、オペレーターを理解している
  • スケジューリングの基本を知っている

Introduction to Airflow のコースページ

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

復習

from airflow.sdk import dag, task

@dag
def star_wars_dag():


@task def get_star_wars_person(): import requests return requests.get("https://swapi.dev/api/people/1/").json()
@task.bash def print_name(person): return f"echo '{person['name']}'"
person = get_star_wars_person() print_name(person) star_wars_dag()
  • DAG は依存関係を持つタスクの集まりです
  • タスクは個々の作業単位です
  • オペレーター/デコレーターは各タスクの処理内容を定義します
  • 依存関係は実行順序を定義します

$$

シンプルな DAG

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

Airflow のアーキテクチャ

$$

Airflow 3 のアーキテクチャ

 

$$

  • オーケストレーション: Scheduler、Dag Processor
  • 実行: Worker、Triggerer
  • インターフェースとストレージ: API Server、Metadata DB
Airflow によるデータパイプラインの構築

スケジューリングの方法

# No automatic runs: trigger manually (default)
@dag(schedule=None)
def my_pipeline(): ...

# Time-based: runs every day at 6 AM @dag(schedule="0 6 * * *") def daily_pipeline(): ...
# Data-aware: runs when an Asset updates @dag(schedule=[Asset("my_asset")]) def downstream_pipeline(): ...
Airflow によるデータパイプラインの構築

DAG の 2 つの書き方

クラシックオペレーター

extract = PythonOperator(
    task_id="extract",
    python_callable=extract_fn)
extract >> transform
  • デコレーターが使えない場合に選択する

TaskFlow API

@task
def extract():
    return {"users": 150}

data = extract()
transform(data)
  • シンプルな Python デコレーター
  • ボイラープレートコードが少ない
  • クラシックオペレーターと組み合わせ可能
Airflow によるデータパイプラインの構築

このコースの演習について

$$

IDE 演習のスクリーンショット

 

  • IDE 演習: 実際の .py ファイルを編集する
  • 「Run this file」 をクリックするか、python3 filename.py を使用する
  • dag.test() で DAG 全体を 1 つのプロセスとして実行できる
Airflow によるデータパイプラインの構築

練習しましょう!

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

Preparing Video For Download...