시작하기무료로 시작하기

DAG 정의하기

이전 연습 문제에서 ETL 프로세스의 세 단계에 대해 실습해 보셨죠:

  • Extract: PostgreSQL의 film 테이블을 pandas로 추출합니다.
  • Transform: film DataFrame의 rental_rate 열을 분리합니다.
  • Load: film DataFrame을 PostgreSQL 데이터 웨어하우스에 적재합니다.

extract_film_to_pandas(), transform_rental_rate(), load_dataframe_to_film() 함수는 워크스페이스에 정의되어 있어요. 이번 연습에서는 기존 DAG에 ETL 작업을 추가하겠습니다. 확장할 DAG와 대기해야 하는 작업은 각각 dagwait_for_table로 워크스페이스에 정의되어 있어요.

이 연습은 강의의 일부입니다

데이터 엔지니어링 입문

강의 보기

연습 안내

  • 연습 설명에 나온 함수들을 사용해 etl() 함수를 완성하세요.
  • etl_tasketl 호출 가능 객체를 사용하도록 하세요.
  • 올바른 상위(Upstream) 의존성을 설정하세요. etl_taskwait_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()
코드 편집 및 실행