전체 그래프
Airflow

XCom

data-engineeringtool/airflowxcom

상위: Airflow

요약

XCom(Cross-Communication)은 Airflow에서 Task 간 데이터 전달을 위한 메커니즘입니다. DAG Run 내에서만 유효하며, 소량의 데이터(문자열, 숫자 등) 전달에 적합합니다. Push로 데이터를 저장하고 Pull로 조회하며, PythonOperator의 return 값은 자동으로 저장됩니다.

XCom이란?

  • Cross-Communication의 약자
  • Airflow에서 Task 간 데이터 전달을 위해 사용
  • 각 Task는 독립적으로 실행되므로 XCom으로 공유
  • DAG Run 내에서만 존재 (다른 DAG Run과 공유 불가)
  • 대용량 데이터 지원 불가
    • 문자열, 숫자 등 작은 크기의 데이터만 공유 가능

XCom vs Variable

항목XComVariable
범위DAG Run 내부전역 (모든 DAG)
유지 기간DAG Run 완료 후 삭제 가능영구 저장
사용 목적Task 간 데이터 전달설정값 공유
저장 위치메타데이터 DB메타데이터 DB

데이터 저장 (Push)

명시적 Push

def push_function(**context):
    context['ti'].xcom_push(key='my_key', value='Hello')

자동 Push (PythonOperator)

def return_value():
    return "Automatically pushed"

task = PythonOperator(
    task_id='auto_push',
    python_callable=return_value
)
# return 값이 자동으로 XCom에 저장됨

데이터 조회 (Pull)

Python에서 Pull

def pull_function(**context):
    value = context['ti'].xcom_pull(
        task_ids='push_task',
        key='my_key'
    )
    print(f"Received: {value}")

기본 키로 Pull

# key 지정 없이 Pull (return  가져오기)
value = context['ti'].xcom_pull(task_ids='task_id')

PythonOperator 예제

Push Task

def push_xcom_value(**kwargs):
    kwargs['ti'].xcom_push(key='message', value='Hello from push_task')

push_task = PythonOperator(
    task_id='push_task',
    python_callable=push_xcom_value
)

Pull Task

def pull_xcom_value(**kwargs):
    message = kwargs['ti'].xcom_pull(task_ids='push_task', key='message')
    print(f"XCom에서 받은 값: {message}")

pull_task = PythonOperator(
    task_id='pull_task',
    python_callable=pull_xcom_value
)

push_task >> pull_task

BashOperator와 XCom

Push

bash_push = BashOperator(
    task_id='bash_push',
    bash_command="echo 'Hello from Bash'",
    do_xcom_push=True  # 출력값을 XCom에 저장
)

Pull (Jinja Template)

bash_pull = BashOperator(
    task_id='bash_pull',
    bash_command="echo '{{ ti.xcom_pull(task_ids=\"bash_push\") }}'"
)

XCom

XCom 사용 패턴

파일 경로 전달

def save_to_s3(**context):
    s3_path = "s3://bucket/data.csv"
    # 파일 저장 로직
    return s3_path  # 경로만 XCom에 저장

def process_s3_file(**context):
    path = context['ti'].xcom_pull(task_ids='save_to_s3')
    # path를 사용하여 처리

처리 결과 전달

def count_records(**context):
    count = 1000  # 실제 카운트 로직
    return count

def check_threshold(**context):
    count = context['ti'].xcom_pull(task_ids='count_records')
    if count > 500:
        print("Threshold exceeded")

주의사항

대용량 데이터 금지

#  잘못된 사용
def wrong_usage():
    df = pd.read_csv('large_file.csv')
    return df  # DataFrame은 너무 

#  올바른 사용
def correct_usage():
    df = pd.read_csv('large_file.csv')
    df.to_parquet('s3://bucket/file.parquet')
    return 's3://bucket/file.parquet'  # 경로만 전달

JSON 직렬화 가능한 데이터만

#  가능
return "string"
return 123
return {"key": "value"}
return [1, 2, 3]

#  불가능 (직렬화 안됨)
return SomeCustomObject()

민감한 정보 전달 피하기

#  비밀번호를 XCom에 저장하지  
return {"password": "secret123"}

#  Connection 사용
from airflow.hooks.base import BaseHook
conn = BaseHook.get_connection('my_conn_id')

XCom 고급 활용

여러 Task에서 Pull

def aggregate_results(**context):
    result_1 = context['ti'].xcom_pull(task_ids='task_1')
    result_2 = context['ti'].xcom_pull(task_ids='task_2')
    result_3 = context['ti'].xcom_pull(task_ids='task_3')
    
    total = result_1 + result_2 + result_3
    return total

[task_1, task_2, task_3] >> aggregate

여러 값 Push

def push_multiple(**context):
    ti = context['ti']
    ti.xcom_push(key='value_1', value=100)
    ti.xcom_push(key='value_2', value=200)
    ti.xcom_push(key='value_3', value=300)

다음 단계

참고 자료