XCom (XComs)

XCom (XComs)

XCom은 "cross-communications"의 줄임말로, Task끼리 서로 데이터를 주고받게 해 주는 메커니즘입니다. 기본적으로 Task는 서로 완전히 격리되어 있고, 전혀 다른 머신에서 실행될 수도 있기 때문에 이 통로가 필요한 거예요.

XCom은 key(본질적으로 이름)와, 그것이 어디서 왔는지를 나타내는 task_id, dag_id로 식별됩니다. 직렬화 가능한 어떤 값이든 담을 수 있는데, @dataclass@attr.define으로 데코레이션된 객체도 포함됩니다(TaskFlow arguments 참고). 다만 XCom은 애초에 작은 데이터를 주고받도록 설계되어 있으니, 데이터프레임 같은 큰 값을 넘기는 용도로는 쓰지 마세요.

XCom 연산은 Task Context를 통해 get_current_context()로 수행해야 합니다. XCom 데이터베이스 모델을 직접 업데이트하는 방식은 불가능해요.

XCom은 Task Instance의 xcom_pushxcom_pull 메서드를 이용해 저장소에 명시적으로 push(밀어 넣기) 하고 pull(끌어 오기) 합니다.

"task-1"이라는 태스크 안에서 값을 push해서 다른 태스크가 쓰게 하려면 이렇게 해요:

# any_serializable_value를 "identifier as string"이라는 key로 XCom에 push
task_instance.xcom_push(key="identifier as a string", value=any_serializable_value)

그 위 코드에서 push한 값을 다른 태스크에서 pull 하려면:

# task-1에서 push한 "identifier as string" key의 XCom 변수를 pull
task_instance.xcom_pull(key="identifier as string", task_ids="task-1")

많은 오퍼레이터는 do_xcom_push 인자가 True일 때(기본값이 True) 결과를 return_value라는 XCom key에 자동으로 push하며, @task 함수도 마찬가지로 동작합니다. 그리고 xcom_pull은 key를 넘기지 않으면 기본적으로 return_value를 key로 사용하므로, 다음과 같이 간단히 쓸 수 있어요:

# "pushing_task"에서 return_value XCom을 pull
value = task_instance.xcom_pull(task_ids='pushing_task')

특정 DAG run에서 값을 가져와야 할 때(예: TriggerDagRunOperator가 트리거한 DAG run)는 dag_idrun_id를 둘 다 명시해 주세요:

trigger_run_id = task_instance.xcom_pull(task_ids="trigger_child", key="trigger_run_id")
child_value = task_instance.xcom_pull(
    task_ids="child_task",
    dag_id="child_dag",
    run_id=trigger_run_id,
)

return_value key(기본적으로 XCom이 push되는 key)는 BaseXCom 클래스의 상수 XCOM_RETURN_KEY로 정의되어 있고, BaseXCom.XCOM_RETURN_KEY로 접근할 수 있습니다.

XCom은 템플릿에서도 사용할 수 있어요:

SELECT * FROM {{ task_instance.xcom_pull(task_ids='foo', key='table_name') }}

참고

xcom_pull()task_ids 인자를 주지 않으면 현재 태스크에서만 pull합니다. Airflow 2에서는 같은 호출이 모든 태스크를 검색해 가장 최근 값을 반환했는데요, 다른 태스크에서 pull할 때는 항상 task_ids를 명시해 주세요.

XCom은 Variables와 가깝지만 차이가 있어요. XCom은 task-instance 단위로 존재하고 하나의 DAG run 안에서 서로 통신하기 위한 용도인 반면, Variables는 전역적으로 존재하며 전반적인 설정과 값 공유를 위해 설계되었습니다.

한 번에 여러 XCom을 push하고 싶다면 do_xcom_pushmultiple_outputs 인자를 True로 설정하고, 값의 딕셔너리를 반환하면 됩니다.

여러 XCom을 push하고 개별적으로 pull하는 예시를 볼게요:

# 딕셔너리를 반환하는 태스크
@task(do_xcom_push=True, multiple_outputs=True)
def push_multiple(**context):
    return {"key1": "value1", "key2": "value2"}


@task
def xcom_pull_with_multiple_outputs(**context):
    # 여러 출력에서 특정 key만 pull
    key1 = context["ti"].xcom_pull(task_ids="push_multiple", key="key1")  # key1을 pull
    key2 = context["ti"].xcom_pull(task_ids="push_multiple", key="key2")  # key2를 pull

    # push_multiple 태스크의 전체 XCom 데이터를 pull
    data = context["ti"].xcom_pull(task_ids="push_multiple", key="return_value")

참고

첫 번째 태스크가 성공하지 못하면, 매 재시도(retry)마다 태스크 XCom이 지워져서 태스크 실행이 멱등(idempotent)하게 유지됩니다. 그래서 XCom으로는 태스크 재시도나 Sensor poke 사이의 상태를 유지할 수 없어요.

Object Storage XCom Backend

기본 XCom 백엔드인 BaseXCom은 XCom을 Airflow 데이터베이스에 저장하는데, 작은 값에는 잘 동작하지만 큰 값이나 XCom이 많아지면 문제가 생길 수 있습니다. 이 한계를 극복하기 위해 더 큰 데이터를 효율적으로 다룰 땐 object storage를 권장합니다. 자세한 내용은 여기 문서를 참고하세요.

Custom XCom Backends

XCom 시스템은 백엔드를 교체할 수 있습니다. 어떤 백엔드를 사용할지는 xcom_backend 설정 옵션으로 정합니다.

자체 백엔드를 구현하고 싶다면 BaseXCom을 서브클래싱하고 serialize_valuedeserialize_value 메서드를 오버라이드하면 됩니다.

BaseXCom 클래스의 purge 메서드도 오버라이드할 수 있는데, 이는 custom backend에서 XCom 데이터를 삭제하는 시점을 제어하기 위한 것입니다. 이 메서드는 delete의 일부로 호출됩니다.

Verifying Custom XCom Backend usage in Containers

Airflow가 어디에 배포되느냐(로컬, Docker, K8s 등)에 따라, custom XCom 백엔드가 실제로 초기화되고 있는지 확인하고 싶을 때가 있어요. 예를 들어 컨테이너 환경의 복잡성 때문에, 컨테이너 배포에서 백엔드가 올바르게 로드되는지 판단하기 어려울 수 있죠. 다행히 아래 방법으로 custom XCom 구현에 대한 신뢰를 쌓을 수 있습니다.

Airflow 컨테이너에서 터미널에 exec로 들어갈 수 있다면, 실제로 사용 중인 XCom 클래스를 출력해 볼 수 있어요:

from airflow.sdk.execution_time.xcom import XCom

print(XCom.__name__)