TaskFlow API

Airflow로 데이터 파이프라인 구축하기

Volker Janz

Senior Developer Advocate at Astronomer

기존 방식

def extract_data():
    return {"users": 150, "events": 4200}

with DAG("etl_pipeline") as dag: t1 = PythonOperator(task_id="extract", python_callable=extract_data) t2 = PythonOperator(task_id="summary", python_callable=print_summary)
t1 >> t2
Airflow로 데이터 파이프라인 구축하기

TaskFlow 방식

from airflow.sdk import dag, task

@dag def etl_pipeline(): @task def extract_data(): return {"users": 150}
data = extract_data() print_summary(data)
Airflow로 데이터 파이프라인 구축하기

TaskFlow를 사용하는 이유

 

  • 보일러플레이트 감소, 데코레이터가 오퍼레이터 인스턴스를 대체
  • 의존성 자동 연결, 반환값으로 태스크가 자동으로 연결
  • 가독성 향상, DAG가 Python 스크립트처럼 읽힘
  • 프로바이더 통합을 위한 기존 오퍼레이터도 계속 사용 가능

TaskFlow API - 시각적 표현

Airflow로 데이터 파이프라인 구축하기

Airflow UI: 그리드 뷰

Airflow 그리드 뷰

 

  • 각 열은 DAG 실행을 나타냄
  • 각 행은 태스크를 나타냄
  • 색상으로 상태 표시: 초록 = 성공, 빨강 = 실패
  • 사각형을 클릭하면 로그와 상세 정보 확인 가능
Airflow로 데이터 파이프라인 구축하기

Airflow UI: 그래프 뷰

Airflow 그리드 뷰

 

  • 하나의 DAG 실행에 집중
  • 의존성 및 병렬 실행 구조를 한눈에 확인
  • 설정에서 세부 사항 조정 가능
Airflow로 데이터 파이프라인 구축하기

Airflow UI: DAG 버전 관리

Airflow UI의 DAG 버전 표시기

 

  • 구조적 변경 사항을 자동으로 추적
  • 각 실행은 해당 시점에 활성화된 버전과 연결
  • DAG 상세 패널에서 버전 확인 가능
Airflow로 데이터 파이프라인 구축하기

연습해 봅시다!

Airflow로 데이터 파이프라인 구축하기

Preparing Video For Download...