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
filmz PostgreSQL dopandas. - Transform (transformace): Rozdělení sloupce
rental_ratev DataFramefilm. - Load (načtení): Uložení DataFrame
filmdo 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
Pokyny k cvičení
- Doplň funkci
etl()s využitím funkcí popsaných v zadání cvičení. - Ujisti se, že
etl_taskpoužívá callableetl. - Nastav správnou závislost –
etl_taskmusí č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()