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

Definování DAG

V předchozích cvičeních jsi prošel/prošla třemi kroky ETL procesu:

  • Extract (extrakce): Načtení tabulky film z PostgreSQL do pandas.
  • Transform (transformace): Rozdělení sloupce rental_rate v DataFrame film.
  • Load (načtení): Uložení DataFrame film do datového skladu v PostgreSQL.

Funkce extract_film_to_pandas(), transform_rental_rate() a load_dataframe_to_film() jsou definované v tvém pracovním prostředí. V tomto cvičení přidáš ETL úlohu do existujícího DAG. DAG, který budeš rozšiřovat, a úloha, na kterou je třeba čekat, jsou v pracovním prostředí definovány jako dag a wait_for_table.

Toto cvičení je součástí kurzu

Introduction to Data Engineering

Zobrazit kurz

Pokyny k cvičení

  • Doplň funkci etl() s využitím funkcí popsaných v zadání cvičení.
  • Ujisti se, že etl_task používá callable etl.
  • Nastav správnou závislost – etl_task musí čekat na dokončení wait_for_table.
  • Ukázkový kód obsahuje testovací spuštění, takže ETL pipeline se spustí hned po spuštění kódu.

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 function
def etl():
    film_df = ____()
    film_df = ____(____)
    ____(____)

# Define the ETL task using PythonOperator
etl_task = PythonOperator(task_id='etl_film',
                          python_callable=____,
                          dag=dag)

# Set the upstream to wait_for_table and sample run etl()
etl_task.____(wait_for_table)
etl()
Upravit a spustit kód