定义一个 DAG
在前面的练习中,您已经实践了 ETL 流程的三个步骤:
- Extract:将 PostgreSQL 中的
film表提取到pandas。 - Transform:拆分
filmDataFrame 的rental_rate列。 - Load:将
filmDataFrame 加载到 PostgreSQL 数仓中。
extract_film_to_pandas()、transform_rental_rate() 和 load_dataframe_to_film() 这三个函数已在您的工作区中定义。在本练习中,您将把一个 ETL 任务添加到现有的 DAG 中。需要扩展的 DAG 以及需要等待的任务已在您的工作区中分别定义为 dag 和 wait_for_table。
本练习是课程的一部分
Data Engineering 入门
练习说明
- 使用练习描述中给出的函数完善
etl()函数。 - 确保
etl_task使用可调用对象etl。 - 设置正确的上游依赖。注意,
etl_task应等待wait_for_table完成后再运行。 - 示例代码包含一次示例运行。这意味着在您运行代码时,ETL 流水线会随之运行。
交互式实操练习
通过完成这段示例代码来试试这个练习。
# Define the ETL function
def etl():
film_df = ____()
film_df = ____(____)
____(____)
# Define the ETL task using PythonOperator
etl_task = PythonOperator(task_id='etl_film',
python_callable=____,
dag=dag)
# Set the upstream to wait_for_table and sample run etl()
etl_task.____(wait_for_table)
etl()