← 학습 카테고리

Learn

Airflow

151개 모듈 · 현재 21번째

Airflow 모듈 21/151 airflow-learn-21

XCom을 통한 태스크 간 통신

Apache Airflow Official Documentation (in-repo snapshot) — Apache Software Foundation core-concepts/xcoms.rst

이 모듈을 다 읽으면

  • XCom을 식별하는 세 요소(key, task_id, dag_id)와 return_value 기본 키의 의미를 설명할 수 있다
  • Airflow 3에서 xcom_pull()이 task_ids 없이 호출될 때의 동작이 Airflow 2와 어떻게 달라졌는지 설명할 수 있다
  • 재시도 시 XCom이 클리어되어 상태를 지속시키는 용도로 쓸 수 없는 이유를 설명할 수 있다

XCom('cross-communications')은 기본적으로 서로 격리되어 있고 다른 머신에서 실행될 수도 있는 Task들이 서로 통신할 수 있게 해주는 메커니즘이다. 이 모듈은 XCom의 식별 방식, push/pull API, multiple_outputs를 통한 다중 XCom 처리, Airflow 2와 달라진 xcom_pull() 기본 동작, 그리고 대용량 데이터를 위한 오브젝트 스토리지 백엔드와 커스텀 백엔드를 다룬다.

XCom의 식별과 push/pull

XCom은 ``key``(사실상 이름)와, 그것이 어디서 왔는지를 나타내는 ``task_id``·``dag_id``로 식별된다. 직렬화 가능한 어떤 값이든(``@dataclass``나 ``@attr.define``으로 데코레이트된 객체 포함) XCom으로 가질 수 있지만, 작은 양의 데이터만을 위해 설계되었으므로 데이터프레임 같은 큰 값을 옮기는 용도로 쓰면 안 된다. XCom 조작은 ``get_current_context()``로 얻은 Task Context를 통해 이루어져야 하며, XCom 데이터베이스 모델을 직접 업데이트하는 것은 불가능하다.

XCom은 Task Instance의 ``xcom_push``와 ``xcom_pull`` 메서드로 명시적으로 '푸시'되고 '풀'된다. "task-1"이라는 태스크 안에서 값을 푸시하려면 ``task_instance.xcom_push(key="식별자", value=값)``을 쓰고, 다른 태스크에서 이를 풀하려면 ``task_instance.xcom_pull(key="식별자", task_ids="task-1")``을 쓴다. 많은 오퍼레이터가 ``do_xcom_push`` 인자가 True(기본값)이면 결과를 ``return_value``라는 키의 XCom으로 자동 push하며, ``@task`` 함수도 마찬가지다. ``xcom_pull``은 키를 넘기지 않으면 기본적으로 ``return_value``를 쓰므로, ``task_instance.xcom_pull(task_ids='pushing_task')``처럼 간단히 쓸 수 있다. ``TriggerDagRunOperator``로 트리거된 자식 Dag Run처럼 특정 Dag Run에서 값을 풀해야 한다면, ``dag_id``와 ``run_id``를 명시적으로 지정해야 한다. return_value 키는 ``BaseXCom`` 클래스의 ``XCOM_RETURN_KEY`` 상수로 정의되어 있고 ``BaseXCom.XCOM_RETURN_KEY``로 접근할 수 있다. Jinja 템플릿 안에서도 ``{{ task_instance.xcom_pull(task_ids='foo', key='table_name') }}``처럼 XCom을 참조할 수 있다.

핵심 포인트

  • XCom은 key, task_id, dag_id 세 요소로 식별되며, 값은 직렬화 가능해야 하고 작은 데이터 전용으로 설계되었다
  • do_xcom_push=True(기본값)인 오퍼레이터와 @task 함수는 결과를 return_value 키로 자동 push하며, xcom_pull()도 키를 생략하면 기본적으로 return_value를 조회한다
  • 다른 Dag Run(예: TriggerDagRunOperator로 트리거된 자식 Dag)의 XCom을 풀하려면 dag_id와 run_id를 명시적으로 지정해야 한다

Airflow 3에서 달라진 xcom_pull() 기본 동작

``xcom_pull()``을 ``task_ids`` 인자 없이 호출하면 현재 태스크로부터만 pull한다. 이는 Airflow 2와 다른 지점이다 — Airflow 2에서는 같은 호출이 모든 태스크를 검색해 가장 최근 값을 반환했다. 따라서 다른 태스크로부터 값을 풀할 때는 항상 ``task_ids``를 명시적으로 지정해야 한다.

또한 XCom은 :doc:`Variables <variables>`와 친척 관계이지만, 핵심 차이는 XCom이 태스크 인스턴스 단위이고 하나의 Dag Run 안에서의 통신을 위해 설계된 반면, Variable은 전역적이며 전체적인 설정과 값 공유를 위해 설계되었다는 점이다.

만약 한 번에 여러 XCom을 push하고 싶다면 ``do_xcom_push``와 ``multiple_outputs`` 인자를 모두 True로 설정하고 딕셔너리를 반환하면 된다. 예를 들어 ``@task(do_xcom_push=True, multiple_outputs=True)``로 데코레이트된 함수가 ``{"key1": "value1", "key2": "value2"}``를 반환하면, 다운스트림 태스크에서 ``ti.xcom_pull(task_ids="push_multiple", key="key1")``처럼 개별 키를 풀할 수도 있고, ``key="return_value"``로 전체 딕셔너리를 한 번에 풀할 수도 있다.

핵심 포인트

  • Airflow 3에서 xcom_pull()을 task_ids 없이 호출하면 현재 태스크로부터만 pull한다 — Airflow 2에서는 전체 태스크를 검색해 최신 값을 반환했다는 점에서 동작이 달라졌다
  • XCom은 태스크 인스턴스 단위·Dag Run 내 통신용이고, Variable은 전역·전체 설정용이라는 점이 두 개념의 핵심 차이다
  • do_xcom_push=True와 multiple_outputs=True를 함께 설정하고 딕셔너리를 반환하면, 다운스트림에서 개별 키로 또는 return_value로 전체를 풀할 수 있다

재시도 시 XCom 클리어, 오브젝트 스토리지·커스텀 백엔드

첫 시도가 실패하면 매 재시도마다 그 태스크의 XCom은 클리어되어 태스크 실행을 멱등하게(idempotent) 만든다. 따라서 XCom은 태스크 재시도나 :doc:`Sensor <sensors>`의 poke 사이에 상태를 지속시키는 용도로는 쓸 수 없다.

기본 XCom 백엔드인 BaseXCom은 XCom을 Airflow 데이터베이스에 저장하는데, 이는 작은 값에는 잘 동작하지만 큰 값이나 대량의 XCom이 오갈 때 문제를 일으킬 수 있다. 이를 극복하기 위해 더 큰 데이터를 효율적으로 다루려면 오브젝트 스토리지를 XCom 백엔드로 쓰는 것이 권장된다.

XCom 시스템은 백엔드를 교체할 수 있게 설계되어 있으며, ``xcom_backend`` 설정 옵션으로 사용할 백엔드를 지정한다. 커스텀 백엔드를 구현하려면 ``BaseXCom``을 서브클래싱하고 ``serialize_value``와 ``deserialize_value`` 메서드를 오버라이드하면 된다. 커스텀 백엔드에서 XCom 데이터를 어떻게 정리(purge)할지 제어하고 싶다면 ``BaseXCom``의 ``purge`` 메서드를 오버라이드할 수 있으며, 이는 ``delete``의 일부로 호출된다. 컨테이너 환경(로컬, Docker, K8s 등)에서 커스텀 XCom 백엔드가 실제로 초기화되었는지 확인하려면, 컨테이너 터미널에 접속해 ``from airflow.sdk.execution_time.xcom import XCom; print(XCom.__name__)``으로 실제 사용 중인 XCom 클래스를 출력해볼 수 있다.

핵심 포인트

  • 첫 시도가 실패하면 재시도마다 그 태스크의 XCom이 클리어되므로, XCom으로 재시도나 센서 poke 사이의 상태를 지속시킬 수 없다
  • 대용량/고빈도 XCom에는 기본 DB 백엔드 대신 오브젝트 스토리지 백엔드가 권장되며, BaseXCom을 서브클래싱해 serialize_value/deserialize_value(선택적으로 purge)를 오버라이드하면 커스텀 백엔드를 만들 수 있다