НачатьНачать бесплатно

Определение DAG

В предыдущих упражнениях вы выполнили три этапа процесса ETL:

  • Извлечение (Extract): извлечение таблицы film из PostgreSQL в pandas.
  • Преобразование (Transform): разделение столбца rental_rate в DataFrame film.
  • Загрузка (Load): загрузка DataFrame film в хранилище данных PostgreSQL.

Функции extract_film_to_pandas(), transform_rental_rate() и load_dataframe_to_film() уже определены в вашем рабочем пространстве. В этом упражнении вы добавите задачу ETL в существующий DAG. DAG, который нужно расширить, и задача, которую необходимо дождаться, определены в рабочем пространстве как dag и wait_for_table соответственно.

Это упражнение является частью курса

Введение в дата-инжиниринг

Посмотреть курс

Инструкции к упражнению

  • Завершите функцию 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()
Редактировать и запускать код