タスクグループで複雑な DAG を整理する

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

Volker Janz

Senior Developer Advocate at Astronomer

複雑化の問題

複雑な DAG

  • 大規模な DAG はグラフビューで見づらくなります
  • タスク名が混在して判別しにくくなります
  • 新しいメンバーが担当箇所を見つけにくくなります
Airflow によるデータパイプラインの構築

@task_group

from airflow.sdk import dag, task, task_group

@task_group(
    group_id="ingest_orders",
    default_args={"retries": 3},
)

def process_orders(): @task def extract_orders(): return [{"id": 1, "amount": 99.99}] @task def transform_orders(orders): return [{"id": o["id"], "total": o["amount"] * 1.08} for o in orders] return transform_orders(extract_orders())
  • group_idカスタム識別子を設定します。省略するとメソッド名が使われます
  • default_args はグループ内のすべてのタスクに適用され、設定の重複を防ぎます
  • タスクグループは UI で折りたたみ可能なブロックとして表示されます
Airflow によるデータパイプラインの構築

Airflow UI での表示

タスクグループなし

extract_orders、transform_orders、load の 3 つのタスクが同じ階層で並ぶシンプルなグラフ

  • すべてのタスクが同じ階層に配置されます

タスクグループあり

extract と transform を含む process_orders という折りたたまれたブロックと、グループ外の load タスクを示すグラフ

  • UI で折りたたみ可能なブロックとして表示されます
Airflow によるデータパイプラインの構築

ネストとカスタム表示名

@task_group(group_display_name="Process All")
def process_all():

    @task_group(group_display_name="Ingest Orders")
    def orders():
        return transform(extract())

    @task_group(group_display_name="Process Returns")
    def returns():
        return transform(extract())

    return {
        "orders": orders(),
        "returns": returns(),
    }
  • より複雑なパイプラインではタスクグループをネストできます
  • group_display_name は UI に表示されるわかりやすいラベルを設定します(絵文字にも対応しています)
Airflow によるデータパイプラインの構築

ファクトリーパターン

@task_group
def process_source(source_name, source_path):

    @task
    def extract():
        return read_data(source_path)

    @task
    def transform(data):
        return clean_data(data)

    return transform(extract())

# Reuse the same pattern for different sources orders = process_source("orders", "/data/orders.csv") returns = process_source("returns", "/data/returns.csv") events = process_source("events", "/data/events.csv")
  • タスクグループはデコレータを付けた Python 関数です
  • 異なるパラメーターで複数回呼び出すことができます
  • このファクトリーパターンにより、パイプラインのロジックを再利用できます
Airflow によるデータパイプラインの構築

グループ化のガイドライン

ドメインや関心事によるグループ化

  • オペレーターの種類ではなく、ドメインや関心事でグループ化しましょう
  • プログラム上の名前には group_id、UI ラベルには group_display_name を使いましょう
  • default_args を使ってグループ内の全タスクにリトライなどの設定を共有しましょう
  • マジカルナンバー 7 の法則を意識し、トップレベルの項目が 7 つを超える場合はタスクグループの導入を検討しましょう
Airflow によるデータパイプラインの構築

練習しましょう!

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

Preparing Video For Download...