Een DAG definiëren
In de vorige oefeningen heb je de drie stappen in het ETL-proces toegepast:
- Extract: Haal de PostgreSQL-tabel
filmop inpandas. - Transform: Splits de kolom
rental_ratevan de DataFramefilm. - Load: Laad de DataFrame
filmin een PostgreSQL-datawarehouse.
De functies extract_film_to_pandas(), transform_rental_rate() en load_dataframe_to_film() staan klaar in je werkruimte. In deze oefening maak je er een geplande Airflow-pijplijn van: de taak etl draait nadat een taak wait_for_table heeft aangegeven dat de brontabel klaar is.
Deze oefening maakt deel uit van de cursus
Introductie tot Data Engineering
Oefeninstructies
- Maak de taak
etl()af door de functies te gebruiken die in de oefeningsbeschrijving staan. - Stel de juiste afhankelijkheid in. Let op: de taak
etlmoet wachten totwait_for_tableklaar is. - De laatste regel is een voorbeeldrun:
etl.function()roept de gewone Python-functie achter de taak aan, zodat de pijplijn hier één keer in de console draait.
Interactieve oefening met praktijkervaring
Probeer deze oefening door deze voorbeeldcode aan te vullen.
# 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()