用 XCom 在任務間傳遞資料

Python 中的 Apache Airflow 入門

Mike Metzger

Data Engineer

什麼是 XCom?

  • 「Cross Communication」
    • 讓任務彼此溝通
  • 儲存在 Airflow 中繼資料資料庫
    • 傳遞少量資料
    • 檔名、URI、列數

示意圖:XCom 在 Airflow 任務間傳遞小型資料

Python 中的 Apache Airflow 入門

不該用 XCom 傳什麼

  • 巨大檔案
  • DataFrame
  • 完整資料庫
  • 大型影像

示意圖:不應透過 XCom 傳送的型別,如大型檔案與 DataFrame

Python 中的 Apache Airflow 入門

實作 XCom

  • 使用 XCom 的方式很多
  • 本章聚焦 TaskFlow API
  • 擴充我們已用過的 @task 做法
Python 中的 Apache Airflow 入門

XCom 範例

@dag(dag_id='Example_XCom')
def example_xcom():

@task def get_data(): return data
@task(multiple_outputs=True) def clean_data(sourcedata): return clean(sourcedata) # Example, not implemented
clean_data(get_data()) example_xcom()
Python 中的 Apache Airflow 入門

XCom 的相依關係

  • XCom 會自動定義相依順序
  • 範例
    clean_data(get_data())
    
  • 觀念上等同 get_data() >> clean_data()
  • 更多範例
     result = clean_data(get_data())
     result >> alert_when_complete()
    
Python 中的 Apache Airflow 入門

檢視 XCom 資料

Airflow XCom 頁面,依 key、Dag 與任務列出儲存的值

Python 中的 Apache Airflow 入門

一起來練習吧!

Python 中的 Apache Airflow 入門

Preparing Video For Download...