शुरू करेंमुफ़्त में शुरू करें

एक DAG परिभाषित करना

पिछले अभ्यासों में आपने ETL प्रक्रिया के तीन चरण लागू किए थे:

  • Extract: film PostgreSQL टेबल को pandas में एक्सट्रैक्ट किया.
  • Transform: film DataFrame के rental_rate कॉलम को स्प्लिट किया.
  • Load: film DataFrame को PostgreSQL डेटा वेयरहाउस में लोड किया.

extract_film_to_pandas(), transform_rental_rate() और load_dataframe_to_film() फंक्शन्स आपके वर्कस्पेस में परिभाषित हैं. इस अभ्यास में, आप इन्हें एक शेड्यूल्ड Airflow पाइपलाइन में बदलेंगे: etl टास्क, सोर्स टेबल के तैयार होने का संकेत देने वाले wait_for_table टास्क के बाद चलेगा.

यह अभ्यास पाठ्यक्रम का हिस्सा है

Introduction to Data Engineering

पाठ्यक्रम देखें

अभ्यास निर्देश

  • अभ्यास विवरण में दिए गए फंक्शन्स का उपयोग करके etl() टास्क पूरा करें.
  • सही डिपेंडेंसी सेट करें. ध्यान दें कि etl टास्क को wait_for_table के खत्म होने का इंतज़ार करना चाहिए.
  • आखिरी लाइन एक सैंपल रन है: etl.function() टास्क के पीछे वाले साधारण Python फंक्शन को कॉल करता है, इसलिए पाइपलाइन यहाँ कंसोल में एक बार चलती है.

इंटरैक्टिव व्यावहारिक अभ्यास

इस अभ्यास को इस नमूना कोड को पूरा करके आज़माएँ।

# 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()
कोड संपादित करें और चलाएँ