НачатьНачать бесплатно

Определение 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={"____": ____},
)
Редактировать и запускать код