การกำหนด DAG
ในแบบฝึกหัดก่อนหน้า คุณได้ลองทำทั้ง 3 ขั้นตอนของกระบวนการ ETL แล้ว ได้แก่
- Extract: ดึงตาราง
filmจาก PostgreSQL เข้าสู่pandas - Transform: แบ่งคอลัมน์
rental_rateของ DataFramefilm - Load: โหลด DataFrame
filmลงใน data warehouse ของ PostgreSQL
ฟังก์ชัน extract_film_to_pandas(), transform_rental_rate() และ load_dataframe_to_film() ถูกกำหนดไว้ใน workspace แล้ว ในแบบฝึกหัดนี้ คุณจะเพิ่ม ETL task เข้าไปใน DAG ที่มีอยู่ โดย DAG ที่ต้องการขยาย และ task ที่ต้องรอถูกกำหนดไว้ใน workspace เป็น dag และ wait_for_table ตามลำดับ
แบบฝึกหัดนี้เป็นส่วนหนึ่งของหลักสูตร
Data Engineering เบื้องต้น
คำแนะนำการฝึกหัด
- เติมฟังก์ชัน
etl()ให้สมบูรณ์โดยใช้ฟังก์ชันที่ระบุไว้ในคำอธิบายแบบฝึกหัด - ตรวจสอบให้แน่ใจว่า
etl_taskใช้ callable ชื่อetl - กำหนด upstream dependency ให้ถูกต้อง โดย
etl_taskต้องรอให้wait_for_tableทำงานเสร็จก่อน - โค้ดตัวอย่างมีการรันทดสอบไว้แล้ว หมายความว่า ETL pipeline จะทำงานทันทีเมื่อรันโค้ด
แบบฝึกหัดเชิงโต้ตอบแบบลงมือทำ
ลองทำแบบฝึกหัดนี้โดยเติมโค้ดตัวอย่างนี้ให้สมบูรณ์
# 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()