Определение DAG
В предыдущих упражнениях вы применили три этапа процесса ETL:
- Извлечение: извлечение таблицы PostgreSQL
filmвpandas. - Преобразование: разделение столбца
rental_rateв DataFramefilm. - Загрузка: загрузка DataFrame
filmв хранилище данных PostgreSQL.
Функции extract_film_to_pandas(), transform_rental_rate() и load_dataframe_to_film() уже определены в вашем рабочем пространстве. В этом упражнении вы превратите их в запланированный пайплайн Airflow: задача etl будет запускаться после того, как задача wait_for_table подтвердит, что исходная таблица готова.
Это упражнение является частью курса
Введение в дата-инжиниринг
Инструкции к упражнению
- Завершите задачу
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()