Визначення 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()