태스크 간 데이터 공유

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

Volker Janz

Senior Developer Advocate at Astronomer

XCom이란?

XCom 푸시-풀 흐름

 

  • XCom = 태스크 간 교차 통신
  • 값은 메타데이터 데이터베이스에 저장됩니다
  • 기본 방식: xcom_push()xcom_pull()을 명시적으로 사용
Airflow로 데이터 파이프라인 구축하기

XComArgs: TaskFlow 방식

 

  • @task의 반환값은 XComArg가 됩니다
  • XComArg는 실제 데이터가 아닌 지연 참조입니다
  • 다른 함수에 전달하면 의존성이 생성됩니다
  • 런타임 시 Airflow가 XCom에서 값을 가져와 해석합니다

TaskFlow의 XComArgs

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

코드로 보는 XComArgs

@task
def get_config():
    return {"name": "daily_etl"}

@task
def run_pipeline(config):
    print(config["name"])


config = get_config() # XComArg
run_pipeline(config) # dependency + data
Airflow로 데이터 파이프라인 구축하기

기본 방식과 TaskFlow 혼합 사용

get_timestamp = BashOperator(
    task_id="get_timestamp",
    bash_command="date +%Y-%m-%d")

@task
def log_timestamp(ts_value):
    print(f"Timestamp: {ts_value}")


log_timestamp(get_timestamp.output)
Airflow로 데이터 파이프라인 구축하기

XCom의 제한 사항

백엔드별 XCom 크기 제한

 

  • 최대 크기는 데이터베이스 백엔드에 따라 다릅니다
  • JSON 직렬화 가능해야 합니다 (딕셔너리, 리스트, 문자열, 숫자)
  • 모든 값은 데이터베이스에 부하를 줍니다
Airflow로 데이터 파이프라인 구축하기

커스텀 XCom 백엔드

커스텀 XCom 백엔드 라우팅

  • 값을 오브젝트 스토리지로 라우팅합니다
  • AIRFLOW__COMMON_IO__XCOM_OBJECTSTORAGE_THRESHOLD: 임계값에 도달했을 때만 XCom을 오브젝트 스토리지에 저장합니다
Airflow로 데이터 파이프라인 구축하기

연습해 봅시다!

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

Preparing Video For Download...