Definirea unui DAG
În exercițiile anterioare ai aplicat cei trei pași ai procesului ETL:
- Extract: Extragi tabelul PostgreSQL
filmînpandas. - Transform: Împarți coloana
rental_ratedin DataFrame-ulfilm. - Load: Încarci DataFrame-ul
filmîntr-un data warehouse 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, le vei transforma într-un pipeline Airflow programat: task-ul etl rulează după ce un task wait_for_table semnalează că tabelul sursă este pregătit.
Acest exercițiu face parte din cursul
Introducere în Data Engineering
Instrucțiuni pentru exercițiu
- Completează task-ul
etl()folosind funcțiile definite în descrierea exercițiului. - Configurează dependența corectă. Reține că task-ul
etltrebuie să aștepte finalizarea task-uluiwait_for_table. - Ultima linie este o rulare exemplu:
etl.function()apelează funcția Python simplă din spatele task-ului, astfel încât pipeline-ul rulează o dată aici, în consolă.
Exercițiu interactiv practic
Încearcă acest exercițiu completând acest cod de exemplu.
# 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()