Definire un DAG
Negli esercizi precedenti hai applicato i tre passaggi del processo ETL:
- Extract: estrai la tabella PostgreSQL
filminpandas. - Transform: suddividi la colonna
rental_ratedel DataFramefilm. - Load: carica il DataFrame
filmin un data warehouse PostgreSQL.
Le funzioni extract_film_to_pandas(), transform_rental_rate() e load_dataframe_to_film() sono già definite nel tuo workspace. In questo esercizio le trasformerai in una pipeline pianificata con Airflow: il task etl viene eseguito dopo che un task wait_for_table segnala che la tabella sorgente è pronta.
Questo esercizio fa parte del corso
Introduzione al Data Engineering
Istruzioni dell'esercizio
- Completa il task
etl()utilizzando le funzioni descritte nel testo dell'esercizio. - Imposta la dipendenza corretta. Nota che il task
etldeve attendere il completamento diwait_for_table. - L'ultima riga esegue un run di esempio:
etl.function()richiama la normale funzione Python dietro al task, così la pipeline viene eseguita una volta qui nella console.
esercizio interattivo pratico
Prova questo esercizio completando questo codice di esempio.
# 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()