Ментальная модель: XCom — это записка, а не грузовик
Задачи в Airflow по умолчанию изолированы. Каждый extract вычислила run_id загрузки, а задача load должна его использовать — нужен явный канал передачи. Этот канал и есть XCom.
Главная ошибка новичков — думать про XCom как про «способ передать данные между задачами». XCom передаёт не данные, а маленькие значения координации: идентификаторы, пути к файлам, флаги, имена партиций, число обработанных строк. Правильная аналогия — это записка на стикере, а не грузовик с грузом. Записку можно прикрепить к холодильнику (metadata DB), и любой member команды её прочитает. Но если вы попытаетесь засунуть в записку весь груз, холодильник рухнет.
Запомните этот тезис, всё остальное в статье — его развёртка: где физически лежит записка, как её положить и забрать, и что делать, когда груз всё-таки нужно передать.
Где физически живёт XCom
XCom расшифровывается как cross-communication. Под капотом это обычная таблица в xcom хранит по строке на каждое опубликованное значение, и ключевые её колонки выглядят так:
| Колонка | Что хранит | Зачем |
|---|---|---|
dag_id | идентификатор DAG | scope: к какому пайплайну относится значение |
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 — это
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 |
| MySQL | 64 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 даёт легальный выход — 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. Но в продакшене с тысячами запусков это означает линейный рост таблицы. Чистка происходит в нескольких местах:
- При
cleartask 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механизм Airflow, разворачивающий одну задачу в N экземпляров по числу элементов на входе; каждый экземпляр пишет свой XCom с уникальным map_index каждый из 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 с метаданными | до пары KB | XCom, но следите за mapping |
| DataFrame, файл, выгрузка | MB–GB | данные в S3/GCS/warehouse, через XCom — только путь |
то же, но хочется return df | MB–GB | custom 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/.