Zacznij terazZacznij za darmo

Definiowanie DAG-a

W poprzednich ćwiczeniach zastosowałeś trzy etapy procesu ETL:

  • Extract: wyodrębnienie tabeli film z PostgreSQL do pandas.
  • Transform: podzielenie kolumny rental_rate w ramce danych film.
  • Load: załadowanie ramki danych film do hurtowni danych PostgreSQL.

Funkcje extract_film_to_pandas(), transform_rental_rate() i load_dataframe_to_film() są już zdefiniowane w twoim środowisku pracy. W tym ćwiczeniu zamienisz je w zaplanowany potok Airflow: zadanie etl uruchomi się dopiero po tym, jak zadanie wait_for_table da sygnał, że tabela źródłowa jest gotowa.

To ćwiczenie jest częścią kursu

Wprowadzenie do inżynierii danych

Zobacz kurs

Instrukcje do ćwiczenia

  • Uzupełnij zadanie etl(), wykorzystując funkcje opisane w treści ćwiczenia.
  • Ustaw odpowiednią zależność. Pamiętaj, że zadanie etl powinno czekać na zakończenie wait_for_table.
  • Ostatnia linijka to przykładowe uruchomienie: etl.function() wywołuje zwykłą funkcję Pythona stojącą za tym zadaniem, dzięki czemu potok uruchomi się raz w konsoli.

Interaktywne ćwiczenie praktyczne

Spróbuj tego ćwiczenia, uzupełniając ten przykładowy kod.

# 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()
Edytuj i uruchom kod