Skip to content
Learning Platform

Ментальная модель: XCom — это записка, а не грузовик

Задачи в Airflow по умолчанию изолированы. Каждый task instance выполняется в своём процессе, нередко на другой машине, и не делит ни память, ни локальные переменные с соседями по DAG. Если задача extract вычислила run_id загрузки, а задача load должна его использовать — нужен явный канал передачи. Этот канал и есть XCom.

Главная ошибка новичков — думать про XCom как про «способ передать данные между задачами». XCom передаёт не данные, а маленькие значения координации: идентификаторы, пути к файлам, флаги, имена партиций, число обработанных строк. Правильная аналогия — это записка на стикере, а не грузовик с грузом. Записку можно прикрепить к холодильнику (metadata DB), и любой member команды её прочитает. Но если вы попытаетесь засунуть в записку весь груз, холодильник рухнет.

Запомните этот тезис, всё остальное в статье — его развёртка: где физически лежит записка, как её положить и забрать, и что делать, когда груз всё-таки нужно передать.

Где физически живёт XCom

XCom расшифровывается как cross-communication. Под капотом это обычная таблица в metadata database. Таблица xcom хранит по строке на каждое опубликованное значение, и ключевые её колонки выглядят так:

КолонкаЧто хранитЗачем
dag_idидентификатор DAGscope: к какому пайплайну относится значение
task_idзадача-публикаторкто положил записку
run_idконкретный запуск DAGизолирует значения разных DagRun
map_indexиндекс mapped-задачи (или -1)поддержка dynamic task mapping
keyимя значения (по умолчанию return_value)несколько записок от одной задачи
valueсериализованный payloadсобственно содержимое
timestampвремя записиаудит и отладка

Критичный нюанс — тип колонки value. В стандартном бэкенде это LargeBinary/JSONB, и Airflow сериализует значение JSON-ом (исторически — pickle, сейчас по умолчанию запрещён ради безопасности). Каждая публикация XCom — это INSERT в metadata DB, а каждое чтение — SELECT. Та же база, в которой scheduler ведёт critical section и держит состояние всего кластера. Отсюда прямое следствие: любой XCom создаёт нагрузку на самый нагруженный и самый критичный компонент Airflow. Положили туда мегабайт — раздули таблицу, замедлили scheduler, рискуете распухшей БД при тысячах запусков в день.

push и pull: ручной способ

Любая задача может явно положить значение через ti.xcom_push и забрать через ti.xcom_pull. Здесь ti — это task instance object, который Airflow прокидывает в контекст исполнения.

from airflow.decorators import dag
from airflow.operators.python import PythonOperator
from datetime import datetime

def extract(**context):
    ti = context["ti"]
    # кладём не данные, а ссылку на данные
    s3_path = "s3://lake/raw/orders/2026-06-01/part-0000.parquet"
    ti.xcom_push(key="orders_path", value=s3_path)

def load(**context):
    ti = context["ti"]
    path = ti.xcom_pull(task_ids="extract", key="orders_path")
    print(f"Гружу из {path}")  # реальные данные читаем напрямую из S3

with __import__("airflow").DAG(
    dag_id="manual_xcom",
    start_date=datetime(2026, 1, 1),
    schedule=None,
    catchup=False,
) as d:
    e = PythonOperator(task_id="extract", python_callable=extract)
    l = PythonOperator(task_id="load", python_callable=load)
    e >> l

Важная деталь: значение, которое функция-callable просто return-ит, Airflow автоматически кладёт в XCom с ключом return_value. Поэтому даже без явного xcom_push ваши задачи могут незаметно засорять таблицу xcom — это включается опцией do_xcom_push (по умолчанию True у многих операторов). Если задача возвращает большой объект, а вам это не нужно, выключайте do_xcom_push=False.

TaskFlow API: тот же XCom, но без церемоний

push/pull многословны и хрупки: вы вручную пишете task_ids строкой, легко опечататься. TaskFlow API (декораторы @task) делает то же самое, но превращает зависимости по данным в обычный Python: возвращаемое значение становится XCom-ом автоматически, а передача его как аргумента в следующую задачу автоматически создаёт и зависимость в графе, и xcom_pull под капотом.

from airflow.decorators import dag, task
from datetime import datetime

@dag(start_date=datetime(2026, 1, 1), schedule=None, catchup=False)
def taskflow_xcom():
    @task
    def extract() -> str:
        # снова возвращаем ССЫЛКУ, а не датафрейм
        return "s3://lake/raw/orders/2026-06-01/part-0000.parquet"

    @task
    def load(path: str) -> int:
        # читаем данные напрямую из объектного хранилища
        return 42  # например, число загруженных строк

    load(extract())

taskflow_xcom()

Здесь extract() возвращает строку — она едет через XCom как return_value. Когда вы передаёте extract() в load(...), Airflow видит зависимость и при исполнении подставит реальное значение. Графа extract >> load вы не писали — он выведен из потока данных. Это сахар: физика та же — INSERT и SELECT в metadata DB. TaskFlow не отменяет лимита размера, он лишь делает код чище.

Где предел: жёсткие границы размера

XCom не предназначен для больших payload, и это не рекомендация, а инженерное ограничение бэкенда. Ориентиры по умолчанию:

Бэкенд metadata DBЛимит на одно XCom-значениеЧто происходит при превышении
PostgreSQLдо 1 GB на поле (теоретически), но практический потолок — десятки KBраздувание таблицы, деградация scheduler
MySQL64 KB (BLOB) по умолчаниюошибка записи, задача падает
SQLite (dev)мизер, годен только для локалкимоментальные проблемы

Цифры из таблицы — про то, что физически влезет. Но «влезет» и «нужно» — разные вещи. Практическое правило: если значение больше нескольких килобайт — это уже не записка, и XCom тут неуместен. Каждый лишний килобайт умножается на число запусков и число mapped-задач: dynamic task mapping на 10000 элементов с XCom по 5 KB — это 50 MB в metadata DB за один DagRun, и так каждый день.

Антипаттерн: «передам DataFrame через XCom»

Самый частый и самый дорогой антипаттерн выглядит безобидно:

@task
def transform() -> "pd.DataFrame":
    df = read_huge_table()      # 2 GB в памяти
    return df                   # АНТИПАТТЕРН: едет через XCom

Что тут идёт не так, по уровням:

  • Сериализация. DataFrame не JSON-сериализуем штатно. Либо задача упадёт, либо (со старым pickle-бэкендом) вы сериализуете 2 GB в bytes и пытаетесь записать в metadata DB.
  • Metadata DB. Даже если влезет — вы пишете гигабайты в БД, которую scheduler опрашивает в критической секции. Это замедляет весь кластер для всех DAG, не только для вашего.
  • Двойная пересылка по сети. Воркер сериализует и шлёт payload в БД; следующий воркер читает обратно. Данные дважды проходят через сеть и через узкое место — БД, — вместо того чтобы лежать там, где их умеют обрабатывать (объектное хранилище, warehouse).
  • Очистка. XCom-строки живут до тех пор, пока их не вычистит retention. Распухшая таблица xcom — классическая причина того, что «Airflow тормозит, а почему — непонятно».

Правильно — передавать ссылку, а данные держать снаружи:

@task
def transform() -> str:
    df = read_huge_table()
    path = "s3://lake/staging/orders/{{ run_id }}/data.parquet"
    df.to_parquet(path)         # данные — в объектное хранилище
    return path                 # через XCom едет только строка-путь

Задача-потребитель получит путь и прочитает Parquet напрямую. XCom остаётся запиской: координирует, но не возит. Этот сдвиг мышления — «передавай ссылки, не груз» — отделяет игрушечный DAG от продакшен-пайплайна.

Custom XCom backend: когда payload всё-таки крупный

Иногда передавать ссылки руками неудобно: хочется писать return df и не думать о путях. Airflow даёт легальный выход — custom XCom backend. Вы подменяете класс BaseXCom своим, и переопределяете два метода:

from airflow.models.xcom import BaseXCom
import uuid, pandas as pd

class S3XComBackend(BaseXCom):
    PREFIX = "xcom_s3://"
    BUCKET = "my-airflow-xcom"

    @staticmethod
    def serialize_value(value, **kwargs):
        if isinstance(value, pd.DataFrame):
            key = f"data/{uuid.uuid4()}.parquet"
            value.to_parquet(f"s3://{S3XComBackend.BUCKET}/{key}")
            # в metadata DB уходит ТОЛЬКО маркер-ссылка
            reference = f"{S3XComBackend.PREFIX}{key}"
            return BaseXCom.serialize_value(reference)
        return BaseXCom.serialize_value(value)

    @staticmethod
    def deserialize_value(result):
        reference = BaseXCom.deserialize_value(result)
        if isinstance(reference, str) and reference.startswith(S3XComBackend.PREFIX):
            key = reference[len(S3XComBackend.PREFIX):]
            return pd.read_parquet(f"s3://{S3XComBackend.BUCKET}/{key}")
        return reference

Подключается через конфиг: [core] xcom_backend = my_pkg.S3XComBackend (или переменная AIRFLOW__CORE__XCOM_BACKEND). После этого код DAG остаётся наивным — return df работает, — но физически крупный объект уезжает в S3/GCS, а в metadata DB ложится лишь короткая ссылка-маркер. Это и есть «лучшее из двух миров»: чистый код задач плюс БД, не раздутая гигабайтами.

Подводные камни custom-бэкенда: вы отвечаете за жизненный цикл объектов в хранилище (Airflow не вычистит ваши Parquet за вас — нужен TTL/lifecycle policy на бакете), за права доступа воркеров к S3 и за то, чтобы сериализация была быстрой. Готовые реализации часто идут в провайдер-пакетах (apache-airflow-providers-amazon, -google) — посмотрите туда, прежде чем писать своё.

Жизненный цикл записки: когда XCom исчезает

Записка не висит на холодильнике вечно, но и не убирается сама в тот же миг. Понимание lifecycle важно, потому что распухшая таблица xcom — одна из самых частых причин медленного Airflow.

По умолчанию XCom-значения переживают свой DagRun: строка в таблице остаётся, пока её явно не удалят. Это сделано осознанно — иногда нужно посмотреть, что задача отдала во вчерашнем запуске, для отладки или backfill. Но в продакшене с тысячами запусков это означает линейный рост таблицы. Чистка происходит в нескольких местах:

  • При clear task instance через UI или CLI Airflow удаляет связанные XCom этой задачи — поэтому повторный запуск стартует с чистого листа.
  • Команда airflow db clean --table xcom --clean-before-timestamp ... физически удаляет старые строки. В продакшене её ставят на регулярный запуск (отдельным maintenance-DAG или cron).
  • При custom backend на S3/GCS удаление строки в metadata DB не удаляет объект в хранилище. За это отвечаете вы — через lifecycle policy бакета или явную очистку. Иначе «холодильник» в БД чистый, а склад S3 растёт бесконтрольно.

Отдельный важный момент — изоляция по run_id. Два разных DagRun одного DAG не видят XCom друг друга по умолчанию: xcom_pull читает значения текущего запуска. Это правильно — иначе backfill за прошлый месяц подхватывал бы сегодняшние записки. Если вам действительно нужно прочитать XCom из другого запуска, это делается явно через параметры xcom_pull, и сам факт такой потребности — повод задуматься, не лучше ли тут Dataset/Asset или внешнее хранилище состояния.

Тонкости, на которых спотыкаются

Несколько неочевидных моментов, которые экономят часы отладки.

  • Sensor в режиме reschedule и XCom. Дефферабельные операторы и сенсоры, которые «засыпают» и просыпаются, при возврате из обработчика тоже могут публиковать XCom. Следите, чтобы они не возвращали тяжёлые объекты на каждом poke.
  • Несколько ключей от одной задачи. Одна задача может опубликовать сколько угодно XCom с разными key. Это нормально для нескольких маленьких значений (row_count, partition, checksum), но не превращайте это в обход лимита размера — десять записок по 10 KB не лучше одной на 100 KB.
  • Mapped-задачи умножают строки. При dynamic task mapping каждый из N экземпляров пишет собственную строку XCom. 10000 mapped-экземпляров — это минимум 10000 INSERT-ов. Это законно и работает, но размер payload здесь критичен вдвойне.
  • render_template_as_native_obj. По умолчанию шаблоны Jinja ({{ ti.xcom_pull(...) }}) возвращают строку. Если вам нужен реальный тип (число, список), включайте этот флаг на уровне DAG — иначе 42 приедет как "42".

Шпаргалка по решениям

Что передаётеРазмерЧем передавать
id, путь, флаг, имя партициибайтыXCom напрямую (push/pull или TaskFlow)
небольшой dict с метаданнымидо пары KBXCom, но следите за mapping
DataFrame, файл, выгрузкаMB–GBданные в S3/GCS/warehouse, через XCom — только путь
то же, но хочется return dfMB–GBcustom XCom backend (S3/GCS)

Что запомнить

XCom — это механизм координации поверх metadata DB, а не транспорт данных. Каждое значение — INSERT/SELECT в самую критичную базу кластера, поэтому передавайте записки (id, пути, флаги), а груз (DataFrame, файлы) держите в объектном хранилище и пересылайте лишь ссылку. TaskFlow API делает push/pull читаемым, но не меняет физику и лимиты. Когда крупный payload всё же нужно «передавать как объект» — это работа для custom XCom backend на S3/GCS, а не для дефолтной таблицы xcom.

Если хотите разобрать это руками — с метаданными в реальной БД, dynamic task mapping и собственным XCom-бэкендом — на нашем бесплатном курсе по Apache Airflow это отдельный модуль с практикой в песочнице: /landing/airflow-course/. А чтобы увидеть, как XCom встаёт в общую картину инженерии данных — от хранилищ до оркестрации, — смотрите направление /directions/data/.

Готовы потрогать XCom вживую?

Курс по Apache Airflow бесплатный: интерактивная песочница, лабораторные с настоящей metadata DB и разбор внутренностей до уровня транзакций. Начните здесь: /landing/airflow-course/.

Ещё в направлении · Data Engineering

Все материалы направления →