एक DAG परिभाषित करना
पिछले अभ्यासों में आपने ETL प्रक्रिया के तीन चरण लागू किए:
- Extract:
filmPostgreSQL टेबल कोpandasमें एक्सट्रैक्ट किया. - Transform:
filmDataFrame केrental_rateकॉलम को स्प्लिट किया. - Load:
filmDataFrame को PostgreSQL डेटा वेयरहाउस में लोड किया.
extract_film_to_pandas(), transform_rental_rate() और load_dataframe_to_film() फंक्शन आपके वर्कस्पेस में परिभाषित हैं. इस अभ्यास में, आप एक मौजूदा DAG में एक ETL टास्क जोड़ेंगे. जिसे बढ़ाना है वह DAG और जिसका इंतज़ार करना है वह टास्क, आपके वर्कस्पेस में क्रमशः dag और wait_for_table के रूप में परिभाषित हैं.
यह अभ्यास पाठ्यक्रम का हिस्सा है
Introduction to Data Engineering
अभ्यास निर्देश
- अभ्यास विवरण में दिए गए फंक्शनों का उपयोग करके
etl()फंक्शन को पूरा करें. - सुनिश्चित करें कि
etl_taskमेंetlcallable उपयोग हो. - सही upstream dependency सेट करें. ध्यान दें,
etl_taskको तब तक इंतज़ार करना चाहिए जब तकwait_for_tableपूरा न हो जाए. - नमूना कोड में एक sample run शामिल है. यानी जब आप कोड चलाएँगे तो ETL पाइपलाइन चलेगी.
इंटरैक्टिव व्यावहारिक अभ्यास
इस अभ्यास को इस नमूना कोड को पूरा करके आज़माएँ।
# 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()