Definice DAGu
Předchozí cvičení tě provedla třemi kroky procesu ETL:
- Extrakce: Extrahování tabulky
filmz PostgreSQL dopandas. - Transformace: Rozdělení sloupce
rental_ratevDataFramefilm. - Načtení: Načtení
DataFramefilmdo datového skladu PostgreSQL.
Funkce extract_film_to_pandas(), transform_rental_rate() a load_dataframe_to_film() jsou definované ve tvém pracovním prostoru. V tomto cvičení z nich vytvoříš naplánovaný pipeline v Airflow: task etl se spustí poté, co task wait_for_table signalizuje, že zdrojová tabulka je připravená.
Toto cvičení je součástí kurzu
Introduction to Data Engineering
Pokyny k cvičení
- Dokonči task
etl()pomocí funkcí definovaných v zadání cvičení. - Nastav správnou závislost. Task
etlby měl počkat, až se dokončíwait_for_table. - Poslední řádek je ukázkové spuštění:
etl.function()zavolá obyčejnou Python funkci za taskem, takže se pipeline jednou spustí přímo v konzoli.
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 ETL task
@task(task_id="etl_film")
def etl():
film_df = ____()
film_df = ____(____)
____(____)
# Add the task to the DAG
@dag(dag_id="etl",
start_date=datetime(2024, 1, 1),
schedule="0 0 * * *")
def etl_dag():
wait_for_table = EmptyOperator(task_id="wait_for_table")
wait_for_table >> ____
etl_dag()
# Sample run
etl.function()