Định nghĩa một DAG
Ở các bài trước, bạn đã áp dụng ba bước trong quy trình ETL:
- Extract: Trích xuất bảng PostgreSQL
filmvàopandas. - Transform: Tách cột
rental_ratecủa DataFramefilm. - Load: Nạp DataFrame
filmvào một kho dữ liệu PostgreSQL.
Các hàm extract_film_to_pandas(), transform_rental_rate() và load_dataframe_to_film() đã được định nghĩa trong không gian làm việc của bạn. Trong bài này, bạn sẽ thêm một tác vụ ETL vào một DAG hiện có. DAG cần mở rộng và tác vụ cần chờ đều đã được định nghĩa trong không gian làm việc lần lượt là dag và wait_for_table.
Bài tập này là một phần của khóa học
Introduction to Data Engineering
Hướng dẫn bài tập
- Hoàn thiện hàm
etl()bằng cách sử dụng các hàm đã nêu trong mô tả bài tập. - Đảm bảo
etl_taskdùng callableetl. - Thiết lập quan hệ phụ thuộc upstream chính xác. Lưu ý
etl_taskphải chờwait_for_tablehoàn thành. - Mã mẫu có kèm một lần chạy mẫu. Điều này nghĩa là pipeline ETL sẽ chạy khi bạn chạy mã.
Bài tập tương tác thực hành trực tiếp
Hãy thử làm bài tập này bằng cách hoàn thành đoạn mã mẫu này.
# 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()