Bir DAG Tanımlama
Önceki egzersizlerde ETL sürecindeki üç adımı uyguladın:
- Extract:
filmPostgreSQL tablosunupandas'a aktar. - Transform:
filmDataFrame'indekirental_ratesütununu böl. - Load:
filmDataFrame'ini bir PostgreSQL veri ambarına yükle.
extract_film_to_pandas(), transform_rental_rate() ve load_dataframe_to_film() fonksiyonları çalışma alanında tanımlı. Bu egzersizde, bunları zamanlanmış bir Airflow veri hattına dönüştüreceksin: kaynak tablonun hazır olduğunu bildiren wait_for_table görevinin ardından etl görevi çalışacak.
Bu egzersiz, kursun bir parçasıdır
Data Engineering'e Giriş
Egzersiz talimatları
- Egzersiz açıklamasında verilen fonksiyonları kullanarak
etl()görevini tamamla. - Doğru bağımlılığı ayarla.
etlgörevi,wait_for_tabletamamlanana kadar beklemelidir. - Son satır örnek bir çalıştırmadır:
etl.function()görev arkasındaki düz Python fonksiyonunu çağırır; böylece boru hattı burada konsolda bir kez çalışır.
Uygulamalı etkileşimli egzersiz
Bu egzersizi bu örnek kodu tamamlayarak deneyin.
# 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()