定義 DAG
在前面的練習中,你已分別完成了擷取(extract)、轉換(transform)和載入(load)各階段。現在這些步驟已整理成一個乾淨的 etl() 函式,你可以在主控台中查看。
etl() 會從相關資料庫擷取原始課程與評分資料,清理損毀資料並填補遺漏值,計算各課程的平均評分,並依照制定的決策規則產生推薦,最後將推薦結果載入資料庫。
如同影片中所提到,etl() 只接受一個引數:db_engines。etl 和 db_engines 都已提供在你的工作環境中,所以你定義的任務只需要用其中一個去呼叫另一個即可。
本練習屬於課程
Data Engineering 入門
練習說明
- 完成 DAG 的定義,讓它每天執行。請使用 cron 表示法。
- 完成任務,使其以資料庫引擎呼叫
etl()函式。
動手互動練習
試著完成這個範例程式碼,體驗一下這個練習。
# Define the DAG so it runs on a daily basis
@dag(dag_id="recommendations",
start_date=datetime(2024, 1, 1),
schedule="____")
def recommendations():
# Make sure the task calls etl() with the database engines
@task(task_id="recommendations_task")
def recommendations_task():
____(____)
recommendations_task()
# Run the DAG
recommendations()