Định nghĩa một DAG
Trong các bài trước, bạn đã áp dụng ba bước của 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 sẵn trong không gian làm việc của bạn. Trong bài này, bạn sẽ biến chúng thành một pipeline Airflow chạy theo lịch: task etl chạy sau khi task wait_for_table báo hiệu bảng nguồn đã sẵn sàng.
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 task
etl()bằng cách sử dụng các hàm đã nêu trong mô tả bài tập. - Thiết lập phụ thuộc đúng. Lưu ý task
etlphải đợiwait_for_tablehoàn thành. - Dòng cuối là một lần chạy mẫu:
etl.function()gọi hàm Python thuần phía sau task, nên pipeline sẽ chạy một lần tại đây trong console.
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 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()