เริ่มต้นใช้งานเริ่มต้นใช้งานได้ฟรี

การกำหนด DAG

ในแบบฝึกหัดก่อนหน้า คุณได้ทำตามสามขั้นตอนของกระบวนการ ETL:

  • Extract: ดึงข้อมูลตาราง film จาก PostgreSQL มาไว้ใน pandas
  • Transform: แยกคอลัมน์ rental_rate ของ DataFrame film
  • Load: โหลด DataFrame film เข้าสู่ data warehouse ของ PostgreSQL

ฟังก์ชัน extract_film_to_pandas(), transform_rental_rate() และ load_dataframe_to_film() ถูกกำหนดไว้ในเวิร์กสเปซของคุณแล้ว ในแบบฝึกหัดนี้ คุณจะเปลี่ยนฟังก์ชันเหล่านี้ให้เป็นไปป์ไลน์ Airflow ที่มีการตั้งเวลา โดย task etl จะรันหลังจาก task wait_for_table แจ้งว่าตารางต้นทางพร้อมใช้งาน

แบบฝึกหัดนี้เป็นส่วนหนึ่งของหลักสูตร

Data Engineering เบื้องต้น

ดูคอร์ส

คำแนะนำการฝึกหัด

  • ทำให้ task etl() สมบูรณ์ โดยใช้ฟังก์ชันที่กำหนดไว้ในคำอธิบายของแบบฝึกหัด
  • กำหนดความสัมพันธ์ระหว่าง task ให้ถูกต้อง โดย task etl ต้องรอให้ wait_for_table เสร็จสิ้นก่อน
  • บรรทัดสุดท้ายเป็นการรันตัวอย่าง: etl.function() จะเรียกใช้ฟังก์ชัน Python ธรรมดาที่อยู่เบื้องหลัง task นี้ ทำให้ไปป์ไลน์รันหนึ่งครั้งในคอนโซล

แบบฝึกหัดเชิงโต้ตอบแบบลงมือทำ

ลองทำแบบฝึกหัดนี้โดยเติมโค้ดตัวอย่างนี้ให้สมบูรณ์

# 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()
แก้ไขและรันโค้ด