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