開始使用免費開始

定義 DAG

在前面的練習中,你已分別完成了擷取(extract)、轉換(transform)和載入(load)各階段。現在這些步驟已整理成一個乾淨的 etl() 函式,你可以在主控台中查看。

etl() 會從相關資料庫擷取原始課程與評分資料,清理損毀資料並填補遺漏值,計算各課程的平均評分,並依照制定的決策規則產生推薦,最後將推薦結果載入資料庫。

如同影片中所提到,etl() 只接受一個引數:db_enginesetldb_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()
編輯並執行程式碼