skip to content

How do you configure a custom XCom backend in Airflow, and what must the class implement?

level: seniorimportance: should knowfreq 42%

answer

  1. you change where, not how, values are stored
  2. one config key names a class
  3. serialize out, deserialize back
  4. every component must import that class
  5. deleting the row leaves the object behind

basics

~20 s

Point the core xcom_backend setting at your class, which subclasses BaseXCom and overrides serialize_value and deserialize_value. Serialize writes the payload to external storage and returns a small pointer; deserialize resolves the pointer back into the value.

solid answer

~50 s

You set `[core] xcom_backend` (or the matching `AIRFLOW__CORE__XCOM_BACKEND` environment variable) to the dotted import path of a class that subclasses `airflow.models.xcom.BaseXCom`. That class overrides `serialize_value`, which writes the real payload somewhere external — an S3 or GCS object keyed by dag, task and run — and returns a small pointer string that becomes the row in the metadata database; and `deserialize_value`, which takes the row back and fetches the payload. Practical caveats matter more than the code. The module must be importable by **every** component — scheduler, workers, triggerer and webserver — or you get a backend that works on workers and breaks the UI. Override `orm_deserialize_value` so the grid view can render a pointer without downloading the object. And own the lifecycle: Airflow deleting the XCom row does not delete your S3 object, so you need bucket lifecycle rules or an explicit cleanup path.

code

python · 24 lines
python
from airflow.models.xcom import BaseXCom

class S3XComBackend(BaseXCom):
    PREFIX = "s3xcom://"
    BUCKET = "company-airflow-xcom"

    @staticmethod
    def serialize_value(value, **kwargs):
        key = build_key(**kwargs)
        write_object(S3XComBackend.BUCKET, key, value)
        ref = f"{S3XComBackend.PREFIX}{S3XComBackend.BUCKET}/{key}"
        return BaseXCom.serialize_value(ref)

    @staticmethod
    def deserialize_value(result):
        ref = BaseXCom.deserialize_value(result)
        if isinstance(ref, str) and ref.startswith(S3XComBackend.PREFIX):
            return read_object(ref)
        return ref

    @staticmethod
    def orm_deserialize_value(result):
        # keep the UI from downloading the payload just to render it
        return BaseXCom.deserialize_value(result)

go deeper

for a junior

Know that the default storage is the metadata database and that Airflow lets you swap in a different backend by configuration; the details are not expected at this level.

for a middle

Explain the mechanism: a configured class subclassing BaseXCom serializes the payload to external storage and stores a pointer, and deserialization resolves that pointer back.

for a senior

Bring the deployment realities — importable by every component, override orm_deserialize_value for the UI, and own the object lifecycle that Airflow's cleanup will not touch.

for a principal

Frame it as a platform-wide policy choice: transparent offloading buys ergonomics at the cost of hidden data movement, a shared bucket's blast radius, and a retention bill someone must own.

## Why a custom backend exists By default, an XCom value is serialized and stored as a row in Airflow's metadata database, which caps the useful payload at kilobytes and puts any bloat on the scheduler's own database. A custom XCom backend changes *where* the bytes live without changing the API DAG authors use: `ti.xcom_push`, `ti.xcom_pull` and TaskFlow argument passing all look identical, while the payload is transparently offloaded to object storage and only a pointer is written to the database. ## Configuring it One setting, deployment-wide: ```ini [core] xcom_backend = my_company.airflow.xcom.S3XComBackend ``` or equivalently `AIRFLOW__CORE__XCOM_BACKEND=my_company.airflow.xcom.S3XComBackend`. There is no per-DAG or per-task override; the choice applies to every XCom in the installation, which is the first thing to say when someone proposes one "just for our team's DAG." ## The class contract The class subclasses `airflow.models.xcom.BaseXCom` and implements the serialize/deserialize pair: ```python from airflow.models.xcom import BaseXCom class S3XComBackend(BaseXCom): PREFIX = "s3xcom://" BUCKET = "company-airflow-xcom" @staticmethod def serialize_value(value, **kwargs): key = build_key(**kwargs) # from dag_id / task_id / run_id / key write_object(S3XComBackend.BUCKET, key, value) reference = f"{S3XComBackend.PREFIX}{S3XComBackend.BUCKET}/{key}" return BaseXCom.serialize_value(reference) @staticmethod def deserialize_value(result): reference = BaseXCom.deserialize_value(result) if isinstance(reference, str) and reference.startswith(S3XComBackend.PREFIX): return read_object(reference) return reference ``` The shape is always the same: `serialize_value` puts the payload somewhere durable and returns something small enough for the column; `deserialize_value` recognises its own pointers and resolves them, passing anything else straight through so existing small XComs still work. The exact signatures of these methods have changed across Airflow releases — check them against the version you are running rather than copying a blog post, because a mismatched signature fails at runtime on every task, not at import. ## `orm_deserialize_value`: the one people miss `BaseXCom` also has `orm_deserialize_value`, used when Airflow needs a representation of the value for display and bookkeeping rather than for a task to consume — most visibly, rendering the XCom tab in the UI. If you do not override it, the webserver will happily download every offloaded object just to draw a table. Overriding it to return the pointer string (or a short summary) keeps the UI fast and keeps webserver processes out of your data path. ## Deployment is where this actually goes wrong The class is imported by the scheduler, every worker, the triggerer and the webserver. If it lives in a package installed only in the worker image, tasks will run fine and the UI will raise import errors when it tries to render XComs — a confusing failure that shows up days later. Ship the backend in a package installed into the same base image every component uses, and make sure the credentials or IAM role for the storage bucket are available to all of them too. Putting the class in the DAGs folder is tempting and fragile: DAG-folder modules are not reliably importable by every component. ## Lifecycle and cost, which the backend does not solve for you A custom backend creates an object per XCom, and Airflow's own cleanup — clearing a task instance, or `airflow db clean` — deletes only the database row. Nothing calls back into your backend to remove the object. You therefore need one of: bucket lifecycle rules that expire the prefix after N days, a scheduled cleanup DAG, or a key layout that makes reruns overwrite rather than accumulate. Otherwise the bucket grows monotonically and the storage bill becomes someone's surprise. There is also a latency cost: every push and every pull now performs a network round trip to object storage, on top of the database write. For thousands of small XComs — a dynamically mapped task with a large fan-out, say — that overhead is real. Some implementations offload only above a size threshold and let small values take the normal path. ## When it is actually worth it Use a custom backend when many teams legitimately pass medium-sized values, when you want the metadata database protected from payload growth without asking every DAG author to change code, or when you need XCom contents to sit inside a specific compliance boundary rather than in the Airflow database. Do not use it as a way to make gigabyte hops between tasks acceptable. The more common and often better answer is the convention version of the same idea: have tasks write to storage explicitly and push a URI, so the DAG shows a reader where the data went and the storage lifecycle is owned by the team that created it.

  • The backend works on workers but the UI raises an import error on the XCom tab. What went wrong?
    The backend class is only installed in the worker image. Scheduler, triggerer and webserver all import the configured `xcom_backend`, so the package must be present in every component's environment, with credentials for the storage bucket available to each. Ship it in the shared base image rather than dropping it in the DAGs folder.
  • What is `orm_deserialize_value` for, and why should you override it?
    It supplies the representation Airflow uses for display and bookkeeping rather than for task consumption, most visibly the UI's XCom view. Left at the default, the webserver downloads every offloaded object just to render a table. Overriding it to return the pointer or a short summary keeps the UI responsive and keeps webserver processes out of the data path.
  • What happens to the S3 objects when `airflow db clean` removes old XCom rows?
    Nothing — they stay. Airflow's cleanup only touches its own tables and never calls back into the backend, so an offloading backend hands you a retention problem. Solve it with bucket lifecycle rules on the prefix, a scheduled cleanup DAG, or a deterministic key layout where reruns overwrite rather than accumulate.

saying these in an interview costs you the question

  • Thinks the backend can be set per DAG or per task
  • Installs the class only in the worker image
  • Assumes deleting the XCom row deletes the S3 object
  • Believes a backend makes gigabyte task-to-task hops fine
  • Cannot name serialize_value and deserialize_value

context