开始使用免费开始使用

定义一个 DAG

在前面的练习中,您已经实践了 ETL 过程中的三步:

  • Extract:将 PostgreSQL 中的 film 表提取到 pandas
  • Transform:拆分 film DataFrame 的 rental_rate 列。
  • Load:将 film DataFrame 加载到 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()
编辑并运行代码