Začněte nyníZačněte zdarma

Definice DAGu

Předchozí cvičení tě provedla třemi kroky procesu ETL:

  • Extrakce: Extrahování tabulky film z PostgreSQL do pandas.
  • Transformace: Rozdělení sloupce rental_rate v DataFrame film.
  • Načtení: Načtení DataFrame film do 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

Zobrazit kurz

Pokyny k cvičení

  • Dokonči task etl() pomocí funkcí definovaných v zadání cvičení.
  • Nastav správnou závislost. Task etl by 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()
Upravit a spustit kód