Definiowanie DAG-a
W poprzednich ćwiczeniach zastosowałeś trzy etapy procesu ETL:
- Extract: wyodrębnienie tabeli
filmz PostgreSQL dopandas. - Transform: podzielenie kolumny
rental_ratew ramce danychfilm. - Load: załadowanie ramki danych
filmdo 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
Instrukcje do ćwiczenia
- Uzupełnij zadanie
etl(), wykorzystując funkcje opisane w treści ćwiczenia. - Ustaw odpowiednią zależność. Pamiętaj, że zadanie
etlpowinno czekać na zakończeniewait_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()