DAG 정의하기
이전 연습 문제에서 ETL 프로세스의 세 단계에 대해 실습해 보셨죠:
- Extract: PostgreSQL의
film테이블을pandas로 추출합니다. - Transform:
filmDataFrame의rental_rate열을 분리합니다. - Load:
filmDataFrame을 PostgreSQL 데이터 웨어하우스에 적재합니다.
extract_film_to_pandas(), transform_rental_rate(), load_dataframe_to_film() 함수는 워크스페이스에 정의되어 있어요. 이번 연습에서는 기존 DAG에 ETL 작업을 추가하겠습니다. 확장할 DAG와 대기해야 하는 작업은 각각 dag와 wait_for_table로 워크스페이스에 정의되어 있어요.
이 연습은 강의의 일부입니다
데이터 엔지니어링 입문
연습 안내
- 연습 설명에 나온 함수들을 사용해
etl()함수를 완성하세요. etl_task가etl호출 가능 객체를 사용하도록 하세요.- 올바른 상위(Upstream) 의존성을 설정하세요.
etl_task는wait_for_table이 완료될 때까지 기다려야 합니다. - 샘플 코드는 예시 실행을 포함합니다. 즉, 코드를 실행하면 ETL 파이프라인이 함께 실행됩니다.
실습형 인터랙티브 연습
이 예제를 이 샘플 코드를 완성하여 풀어보세요.
# 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()