Definiowanie DAG-a
W poprzednich ćwiczeniach zrealizowałeś trzy etapy procesu ETL:
- Ekstrakcja: Wyodrębnienie tabeli
filmz bazy PostgreSQL dopandas. - Transformacja: Podział kolumny
rental_ratew DataFramefilm. - Ładowanie: Załadowanie DataFrame
filmdo hurtowni danych PostgreSQL.
Funkcje extract_film_to_pandas(), transform_rental_rate() i load_dataframe_to_film() są już zdefiniowane w twoim środowisku roboczym. W tym ćwiczeniu dodasz zadanie ETL do istniejącego DAG-a. DAG do rozszerzenia oraz zadanie, na które należy czekać, są zdefiniowane w środowisku jako dag i wait_for_table.
To ćwiczenie jest częścią kursu
Wprowadzenie do inżynierii danych
Instrukcje do ćwiczenia
- Uzupełnij funkcję
etl(), korzystając z funkcji opisanych w treści ćwiczenia. - Upewnij się, że
etl_taskużywa funkcji wywoływalnejetl. - Ustaw właściwą zależność nadrzędną. Pamiętaj, że
etl_taskpowinno czekać na zakończeniewait_for_table. - Przykładowy kod zawiera próbne uruchomienie – oznacza to, że potok ETL zostanie wykonany po uruchomieniu kodu.
Interaktywne ćwiczenie praktyczne
Spróbuj tego ćwiczenia, uzupełniając ten przykładowy kod.
# Define the ETL function
def etl():
film_df = ____()
film_df = ____(____)
____(____)
# Define the ETL task using PythonOperator
etl_task = PythonOperator(task_id='etl_film',
python_callable=____,
dag=dag)
# Set the upstream to wait_for_table and sample run etl()
etl_task.____(wait_for_table)
etl()