定義 DAG
在前面的練習中,你已經完成 ETL 流程的三個步驟:
- Extract:將 PostgreSQL 的
film資料表匯入pandas。 - Transform:將
filmDataFrame 的rental_rate欄位拆分。 - Load:將
filmDataFrame 載入 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()