開始使用免費開始

定義 DAG

在前面的練習中,你已經完成 ETL 流程的三個步驟:

  • Extract:將 PostgreSQL 的 film 資料表匯入 pandas
  • Transform:將 film DataFrame 的 rental_rate 欄位拆分。
  • Load:將 film DataFrame 載入 PostgreSQL 的資料倉儲。

extract_film_to_pandas()transform_rental_rate()load_dataframe_to_film() 這三個函式已在你的工作區中定義。這個練習中,你會把它們串成一個排程的 Airflow 管線:在來源資料表就緒的 wait_for_table 任務發出訊號後,etl 任務才會執行。

本練習屬於課程

Data Engineering 入門

檢視課程

練習說明

  • 使用練習說明中提供的函式完成 etl() 任務。
  • 設定正確的相依關係。注意 etl 任務應該要等 wait_for_table 完成後才執行。
  • 最後一行是範例執行:etl.function() 會呼叫這個任務背後的純 Python 函式,因此在主控台中會讓管線執行一次。

動手互動練習

試著完成這個範例程式碼,體驗一下這個練習。

# Define the ETL task
@task(task_id="etl_film")
def etl():
    film_df = ____()
    film_df = ____(____)
    ____(____)

# Add the task to the DAG
@dag(dag_id="etl",
     start_date=datetime(2024, 1, 1),
     schedule="0 0 * * *")
def etl_dag():
    wait_for_table = EmptyOperator(task_id="wait_for_table")
    wait_for_table >> ____

etl_dag()

# Sample run
etl.function()
編輯並執行程式碼