Asset Partitions によるパーティション対応スケジューリング

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

Volker Janz

Senior Developer Advocate at Astronomer

ds だけでは不十分な場合

 

  • {{ ds }} は各実行を日付に紐付けますが、これは時刻スケジュールの DAG にのみ有効です
  • アセットトリガーの DAG では dsNone に設定され、利用できません
  • ダウンストリームの DAG はどの日付が更新されたかを知る手段がありません
  • Asset Partitionspartition_key をアップストリームからダウンストリームへ伝播させます

アップストリームのスケジュール実行が、ds が None のダウンストリームのアセットスケジュール実行をトリガーする

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

CronPartitionTimetable

from airflow.sdk import dag, CronPartitionTimetable

@dag(schedule=CronPartitionTimetable("0 0 * * *", timezone="UTC"))
def sales_pipeline():
    ...

 

  • スケジュール実行ごとに partition_key が自動的に付与されます
  • 手動実行はキーを指定しない限りパーティション化されません
Airflow によるデータパイプラインの構築

パーティションキーへのアクセス

Python タスクの場合:

@task
def process_data(dag_run=None):
    partition_key = dag_run.partition_key
    print(f"Processing: {partition_key}")

SQL テンプレートの場合:

DELETE FROM daily_summary
WHERE order_date = '{{ dag_run.partition_key[:10] }}';

 

  • パーティション化された実行内のすべてのタスクで利用できます
Airflow によるデータパイプラインの構築

パーティション化されたアセットイベント

from airflow.sdk import dag, task, Asset, CronPartitionTimetable

daily_sales = Asset("daily_sales")

@dag(schedule=CronPartitionTimetable("0 0 * * *", timezone="UTC"))
def sales_pipeline():
    @task(outlets=[daily_sales])
    def load_data(**context):
        ...

 

  • タスクの outlet + CronPartitionTimetable = パーティション化されたアセットイベント
  • イベントは DAG 実行と同じパーティションキーを持ちます
  • データが更新されたことだけでなく、どのパーティションが準備完了かを通知します
Airflow によるデータパイプラインの構築

PartitionedAssetTimetable

from airflow.sdk import dag, task, Asset, PartitionedAssetTimetable

daily_sales = Asset("daily_sales")

@dag(schedule=PartitionedAssetTimetable(assets=daily_sales))
def sales_report():
    @task
    def generate_report(dag_run=None):
        partition_key = dag_run.partition_key
        print(f"Report for: {partition_key}")

 

  • パーティション化されたアセットイベントに対してのみトリガーされます
  • ダウンストリームの実行はパーティションキーを引き継ぎます
Airflow によるデータパイプラインの構築

StartOfDayMapper

from airflow.sdk import (
    dag, PartitionedAssetTimetable,
    StartOfDayMapper,
)

@dag(schedule=PartitionedAssetTimetable(
    assets=daily_sales,
    partition_mapper_config={
        daily_sales: StartOfDayMapper()
    },
))
def sales_report():
    ...

 

マッピング前: 2026-04-23T00:00:00

StartOfDayMapper 適用後: 2026-04-23

$$

  • 時間系キーStartOfHourMapperStartOfWeekMapperStartOfMonthMapper
  • 非時間系キー:地域や部門には AllowedKeyMapper
Airflow によるデータパイプラインの構築

全体フロー

パーティションの全体フロー

  • アップストリーム:CronPartitionTimetable が各実行にパーティションキーを付与
  • outlet タスクがそのキーを持つパーティション化されたアセットイベントを発行
  • StartOfDayMapper がタイムスタンプを日付に正規化
  • ダウンストリーム:PartitionedAssetTimetable がトリガーされ、マッピング済みキーを引き継ぐ
Airflow によるデータパイプラインの構築

練習しましょう!

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

Preparing Video For Download...