Mendefinisikan DAG
Pada latihan sebelumnya Anda menerapkan tiga langkah dalam proses ETL:
- Extract: Mengekstrak tabel PostgreSQL
filmkepandas. - Transform: Memecah kolom
rental_ratedari DataFramefilm. - Load: Memuat DataFrame
filmke 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
Instruksi latihan
- Lengkapi task
etl()dengan menggunakan fungsi-fungsi yang didefinisikan dalam deskripsi latihan. - Atur dependensi yang benar. Perhatikan bahwa task
etlharus menunggu hinggawait_for_tableselesai. - 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()