การกำหนด DAG
ในแบบฝึกหัดก่อนหน้า คุณได้ทำตามสามขั้นตอนของกระบวนการ 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() ถูกกำหนดไว้ในเวิร์กสเปซของคุณแล้ว ในแบบฝึกหัดนี้ คุณจะเปลี่ยนฟังก์ชันเหล่านี้ให้เป็นไปป์ไลน์ 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()