Визначення DAG
У попередніх вправах ви окремо виконали етапи extract, transform і load. Тепер усе це зібрано в одну акуратну функцію etl(), з якою ви можете ознайомитися в консолі.
Функція etl() отримує сирі дані про курси та рейтинги з відповідних баз даних, очищає пошкоджені дані й заповнює пропуски, обчислює середній рейтинг для кожного курсу та створює рекомендації на основі правил ухвалення рішення для формування рекомендацій, а наприкінці завантажує рекомендації до бази даних.
Як ви, мабуть, пам'ятаєте з відео, etl() приймає єдиний аргумент: db_engines. І etl, і db_engines доступні у вашому робочому середовищі, тож завдання, яке ви визначаєте, має лише викликати одну з цих сутностей з іншою як аргументом.
Ця вправа є частиною курсу
Вступ до Data Engineering
Інструкції до вправи
- Доповніть визначення DAG так, щоб він запускався щодня. Обов'язково використайте нотацію cron.
- Завершіть завдання так, щоб воно викликало функцію
etl()із рушіями баз даних.
Інтерактивна практична вправа
Спробуйте виконати цю вправу, доповнивши цей зразок коду.
# Define the DAG so it runs on a daily basis
@dag(dag_id="recommendations",
start_date=datetime(2024, 1, 1),
schedule="____")
def recommendations():
# Make sure the task calls etl() with the database engines
@task(task_id="recommendations_task")
def recommendations_task():
____(____)
recommendations_task()
# Run the DAG
recommendations()