Definování DAGu
V předchozích cvičeních jsi postupně dokončil/a fáze extrakce, transformace a nahrání dat. Teď je to všechno spojené do jedné přehledné funkce etl(), kterou najdeš v konzoli.
Funkce etl() extrahuje surová data o kurzech a hodnoceních z příslušných databází, vyčistí poškozená data a doplní chybějící hodnoty, spočítá průměrné hodnocení pro každý kurz a na základě rozhodovacích pravidel vytvoří doporučení, která nakonec nahraje do databáze.
Jak si možná pamatuješ z videa, funkce etl() přijímá jediný argument: db_engines. Jak etl, tak db_engines máš k dispozici ve svém pracovním prostředí, takže úloha, kterou definuješ, stačí, aby jedno zavolala s druhým.
Toto cvičení je součástí kurzu
Introduction to Data Engineering
Pokyny k cvičení
- Doplň definici DAGu tak, aby se spouštěl denně. Nezapomeň použít cronový zápis.
- Doplň úlohu tak, aby volala funkci
etl()s databázovými enginy.
Interaktivní cvičení na vyzkoušení si v praxi
Vyzkoušejte si toto cvičení dokončením tohoto ukázkového kódu.
# 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()