Визначення DAG
У попередніх вправах ви окремо виконали етапи extract, transform і load. Тепер усе це зібрано в охайну функцію etl(), з якою ви можете ознайомитися в консолі.
Функція etl() витягує сирі дані про курси та рейтинги з відповідних баз даних, очищує пошкоджені дані й заповнює пропущені значення, обчислює середній рейтинг для кожного курсу, створює рекомендації на основі правил прийняття рішень для формування рекомендацій і зрештою завантажує ці рекомендації до бази даних.
Як ви, мабуть, пам'ятаєте з відео, etl() приймає один аргумент: db_engines. Ви можете передати його до завдання через op_kwargs у PythonOperator. Передайте туди словник, який буде підставлено як kwargs у виклику функції.
Ця вправа є частиною курсу
Вступ до Data Engineering
Інструкції до вправи
- Завершіть визначення DAG так, щоб він запускався щодня. Обов'язково використайте нотацію cron.
- Завершіть створення
PythonOperator(), передавши потрібні аргументи. Окрімetl, у вашому робочому середовищі також доступнийdb_engines.
Інтерактивна практична вправа
Спробуйте виконати цю вправу, доповнивши цей зразок коду.
# Define the DAG so it runs on a daily basis
dag = DAG(dag_id="recommendations",
schedule_interval="____")
# Make sure `etl()` is called in the operator. Pass the correct kwargs.
task_recommendations = PythonOperator(
task_id="recommendations_task",
python_callable=____,
op_kwargs={"____": ____},
)