Kom igångKom igång gratis

Definiera en DAG

I de föregående övningarna tillämpade du de tre stegen i ETL-processen:

  • Extract: Extrahera PostgreSQL-tabellen film till pandas.
  • Transform: Dela upp kolumnen rental_rate i DataFrame:n film.
  • Load: Ladda DataFrame:n film till ett PostgreSQL-datalager.

Funktionerna extract_film_to_pandas(), transform_rental_rate() och load_dataframe_to_film() är definierade i din arbetsyta. I den här övningen ska du omvandla dem till en schemalagd Airflow-pipeline: uppgiften etl körs efter att uppgiften wait_for_table signalerar att källtabellen är redo.

Den här övningen är en del av kursen

Introduktion till datatekniker

Visa kurs

Övningsinstruktioner

  • Slutför uppgiften etl() genom att använda funktionerna som beskrivs i övningsbeskrivningen.
  • Ställ in rätt beroende. Observera att uppgiften etl ska vänta på att wait_for_table är klar.
  • Den sista raden är en exempelkörning: etl.function() anropar den vanliga Python-funktionen bakom uppgiften, så pipelinen körs en gång här i konsolen.

Interaktiv övning med praktiskt arbete

Testa den här övningen genom att slutföra den här exempelkoden.

# 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()
Redigera och kör kod