Определение DAG
В предыдущих упражнениях вы выполнили три этапа процесса ETL:
- Извлечение (Extract): извлечение таблицы
filmиз PostgreSQL вpandas. - Преобразование (Transform): разделение столбца
rental_rateв DataFramefilm. - Загрузка (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()