Definirea unui DAG
În exercițiile anterioare ai aplicat cei trei pași din procesul ETL:
- Extragere (Extract): Extragerea tabelului PostgreSQL
filmînpandas. - Transformare (Transform): Divizarea coloanei
rental_ratedin DataFrame-ulfilm. - Încărcare (Load): Încărcarea DataFrame-ului
filmîntr-un depozit de date PostgreSQL.
Funcțiile extract_film_to_pandas(), transform_rental_rate() și load_dataframe_to_film() sunt definite în spațiul tău de lucru. În acest exercițiu, vei adăuga o sarcină ETL la un DAG existent. DAG-ul pe care îl vei extinde și sarcina pentru care trebuie să aștepți sunt definite în spațiul de lucru ca dag, respectiv wait_for_table.
Acest exercițiu face parte din cursul
Introducere în Data Engineering
Instrucțiuni pentru exercițiu
- Completează funcția
etl()folosind funcțiile descrise în enunțul exercițiului. - Asigură-te că
etl_taskutilizează funcția apelabilăetl. - Configurează dependența upstream corectă. Reține că
etl_tasktrebuie să aștepte finalizareawait_for_table. - Codul exemplu conține o rulare demonstrativă. Asta înseamnă că pipeline-ul ETL rulează în momentul în care execuți codul.
Exercițiu interactiv practic
Încearcă acest exercițiu completând acest cod de exemplu.
# 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()