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