Mulai sekarangMulai gratis

Mendefinisikan DAG

Pada latihan sebelumnya Anda menerapkan tiga langkah dalam proses ETL:

  • Extract: Mengekstrak tabel PostgreSQL film ke pandas.
  • Transform: Memecah kolom rental_rate dari DataFrame film.
  • Load: Memuat DataFrame film ke gudang data PostgreSQL.

Fungsi extract_film_to_pandas(), transform_rental_rate() dan load_dataframe_to_film() sudah didefinisikan di workspace Anda. Pada latihan ini, Anda akan mengubahnya menjadi pipeline terjadwal di Airflow: task etl berjalan setelah task wait_for_table memberi sinyal bahwa tabel sumber siap.

Latihan ini merupakan bagian dari kursus

Pengantar Data Engineering

Lihat Kursus

Instruksi latihan

  • Lengkapi task etl() dengan menggunakan fungsi-fungsi yang didefinisikan dalam deskripsi latihan.
  • Atur dependensi yang benar. Perhatikan bahwa task etl harus menunggu hingga wait_for_table selesai.
  • Baris terakhir adalah contoh run: etl.function() memanggil fungsi Python biasa di balik task tersebut, sehingga pipeline dijalankan sekali di konsol ini.

Latihan interaktif langsung praktik

Cobalah latihan ini dengan melengkapi kode contoh ini.

# 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()
Edit dan Jalankan Kode