定义一个 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 流水线:etl 任务会在 wait_for_table 任务发出源表就绪的信号后再运行。
本练习是课程的一部分
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()