एक DAG परिभाषित करना
पिछले अभ्यासों में आपने ETL प्रक्रिया के तीन चरण लागू किए थे:
- Extract:
filmPostgreSQL टेबल कोpandasमें एक्सट्रैक्ट किया. - Transform:
filmDataFrame केrental_rateकॉलम को स्प्लिट किया. - Load:
filmDataFrame को 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()