Definiera en DAG
I de föregående övningarna tillämpade du de tre stegen i ETL-processen:
- Extract: Extrahera PostgreSQL-tabellen
filmtillpandas. - Transform: Dela upp kolumnen
rental_ratei DataFrame:nfilm. - Load: Ladda DataFrame:n
filmtill 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
Övningsinstruktioner
- Slutför uppgiften
etl()genom att använda funktionerna som beskrivs i övningsbeskrivningen. - Ställ in rätt beroende. Observera att uppgiften
etlska vänta på attwait_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()