ПочатиПочніть безкоштовно

Визначення DAG

У попередніх вправах ви виконали три кроки процесу ETL:

  • Extract: Екстрагуйте таблицю PostgreSQL film у pandas.
  • Transform: Розбийте стовпчик rental_rate у датафреймі film.
  • Load: Завантажте датафрейм film до сховища даних 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()
Редагувати та запускати код