Определение DAG
В предыдущих упражнениях вы выполнили фазы извлечения, преобразования и загрузки по отдельности. Теперь всё это объединено в одну аккуратную функцию etl(), с которой вы можете ознакомиться в консоли.
Функция etl() извлекает необработанные данные о курсах и оценках из соответствующих баз данных, очищает повреждённые данные и заполняет пропущенные значения, вычисляет средний рейтинг по каждому курсу и формирует рекомендации на основе заданных правил, а затем загружает эти рекомендации в базу данных.
Как вы помните из видео, функция etl() принимает единственный аргумент: db_engines. Передать его в задачу можно с помощью op_kwargs в PythonOperator. Для этого передайте словарь, значения которого будут подставлены как kwargs в вызываемую функцию.
Это упражнение является частью курса
Введение в дата-инжиниринг
Инструкции к упражнению
- Завершите определение 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={"____": ____},
)