タスク間のデータ共有

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

Volker Janz

Senior Developer Advocate at Astronomer

XCom とは?

XComのプッシュ・プルフロー

 

  • XCom = タスク間のクロスコミュニケーション
  • 値はメタデータデータベースに保存される
  • 従来の方法:明示的な xcom_push()xcom_pull()
Airflow によるデータパイプラインの構築

XComArg:TaskFlow の方法

 

  • @task の戻り値は XComArg になる
  • XComArg は実際のデータではなく遅延参照
  • 別の関数に渡すことで依存関係が生成される
  • 実行時に Airflow が XCom から値を取得して解決する

TaskFlow における XComArg

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

コードでの XComArg

@task
def get_config():
    return {"name": "daily_etl"}

@task
def run_pipeline(config):
    print(config["name"])


config = get_config() # XComArg
run_pipeline(config) # dependency + data
Airflow によるデータパイプラインの構築

従来の方法と TaskFlow の組み合わせ

get_timestamp = BashOperator(
    task_id="get_timestamp",
    bash_command="date +%Y-%m-%d")

@task
def log_timestamp(ts_value):
    print(f"Timestamp: {ts_value}")


log_timestamp(get_timestamp.output)
Airflow によるデータパイプラインの構築

XCom の制限

バックエンド別の XCom サイズ制限

 

  • 最大サイズはデータベースバックエンドによって異なる
  • JSON シリアライズ可能な値(辞書、リスト、文字列、数値)のみ対応
  • 値を保存するたびにデータベースに負荷がかかる
Airflow によるデータパイプラインの構築

カスタム XCom バックエンド

カスタム XCom バックエンドのルーティング

  • 値をオブジェクトストレージにルーティングする
  • AIRFLOW__COMMON_IO__XCOM_OBJECTSTORAGE_THRESHOLD:しきい値に達した場合にのみ XCom をオブジェクトストレージに保存する
Airflow によるデータパイプラインの構築

練習しましょう!

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

Preparing Video For Download...