시작하기무료로 시작하기

DAG 정의하기

이전 연습 문제에서 ETL 프로세스의 세 단계를 적용해 보았습니다:

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

extract_film_to_pandas(), transform_rental_rate(), load_dataframe_to_film() 함수는 이미 작업 공간에 정의되어 있습니다. 이번 연습에서는 이 함수들을 예약된 Airflow 파이프라인으로 만들어 볼 거예요. etl 작업은 wait_for_table 작업이 원본 테이블 준비 완료 신호를 보낸 후에 실행됩니다.

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

데이터 엔지니어링 입문

강의 보기

연습 안내

  • 연습 문제 설명에 정의된 함수들을 사용해 etl() 작업을 완성하세요.
  • 올바른 의존성을 설정하세요. etl 작업은 wait_for_table이 완료될 때까지 기다려야 한다는 점에 유의하세요.
  • 마지막 줄은 샘플 실행입니다. etl.function()은 작업 뒤에 있는 순수 Python 함수를 호출하므로, 여기서는 콘솔에서 파이프라인이 한 번 실행됩니다.

실습형 인터랙티브 연습

이 예제를 이 샘플 코드를 완성하여 풀어보세요.

# 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()
코드 편집 및 실행