Aan de slagBegin gratis

Een DAG definiëren

In de vorige oefeningen heb je de drie stappen in het ETL-proces toegepast:

  • Extract: Haal de PostgreSQL-tabel film op in pandas.
  • Transform: Splits de kolom rental_rate van de DataFrame film.
  • Load: Laad de DataFrame film in 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

Bekijk cursus

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 etl moet wachten tot wait_for_table klaar 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()
Code bewerken en uitvoeren