开始使用免费开始使用

定义 DAG

在之前的练习中,您分别完成了抽取、转换和加载三个阶段。现在,这些步骤已整合到一个简洁的 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()
编辑并运行代码