상위: 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
| 항목 | XCom | Variable |
|---|---|---|
| 범위 | 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 사용 패턴
파일 경로 전달
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)