Визначення DAG
У попередніх вправах ви застосували три кроки процесу ETL:
- Extract: Екстрагуйте таблицю PostgreSQL
filmдоpandas. - Transform: Розбийте стовпець
rental_rateу датафрейміfilm. - Load: Завантажте датафрейм
filmдо сховища даних PostgreSQL.
Функції extract_film_to_pandas(), transform_rental_rate() і load_dataframe_to_film() уже визначено у вашому робочому середовищі. У цій вправі ви додасте ETL‑завдання до наявного DAG. DAG, який потрібно розширити, і завдання, на яке слід чекати, визначено у вашому середовищі як dag і wait_for_table відповідно.
Ця вправа є частиною курсу
Вступ до Data Engineering
Інструкції до вправи
- Доповніть функцію
etl()з використанням функцій, наведених в описі вправи. - Переконайтеся, що
etl_taskвикористовує викличний об'єктetl. - Налаштуйте правильну залежність upstream. Зверніть увагу:
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()