在任務之間共享資料

使用 Airflow 建置資料管線

Volker Janz

Senior Developer Advocate at Astronomer

什麼是 XCom?

XCom push pull flow

 

  • XCom = 任務間的交互通訊
  • 值儲存在 中繼資料資料庫
  • 傳統作法:明確呼叫 xcom_push()xcom_pull()
使用 Airflow 建置資料管線

XComArg:TaskFlow 的做法

 

  • @task 的回傳值會變成一個 XComArg
  • XComArg 是延遲參照,不是實際資料
  • 把它傳給另一個函式會建立相依關係
  • 執行時,Airflow 會從 XCom 解析並取回

XComArgs in TaskFlow

使用 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 size limits by backend

 

  • 最大大小取決於你的 資料庫後端
  • 必須能 JSON 序列化(dict、list、字串、數字)
  • 每個值都會為資料庫帶來負載
使用 Airflow 建置資料管線

自訂 XCom 後端

Custom XCom backend routing

  • 將值導向 物件儲存
  • AIRFLOW__COMMON_IO__XCOM_OBJECTSTORAGE_THRESHOLD:僅在達到臨界值時,才把 XCom 存到物件儲存
使用 Airflow 建置資料管線

一起來練習吧!

使用 Airflow 建置資料管線

Preparing Video For Download...