Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing

Orchestration & Reliable Pipelines

Progress0/14
x

What is a Pipeline

  • What is a data pipeline?12m
  • Scheduling and cron12m
  • Monitoring and alerting12m

DAG Mental Model

  • Why a scheduler12m
  • Tasks, edges, topological order14m
  • Operators, sensors, and XCom12m

Retries, Idempotency, Dedup

  • Retries in the DAG runner14m
  • Idempotent loads14m
  • Deduplicate before load12m

Backfills, Catchup, Intervals

  • Backfills and catchup12m
  • Data interval vs execution clock12m

CDC Concepts

  • CDC: inserts, updates, deletes12m
  • Apply CDC idempotently14m

Capstone

  • Capstone: a reliable night job16m
Back to track
  1. Learn
  2. Orchestration & Reliable Pipelines
  3. DAG Mental Model
  4. Operators, sensors, and XCom

Lesson 6 of 14 · Theory first, then run it

Operators, sensors, and XCom

orchestrationpythonbeginner12 min

Overview

Airflow's words for 'run this,' 'wait until that exists,' and 'pass a small note downstream.' You are not running Airflow.

On this page7 sections›
  1. 1The idea
  2. 2Why this exists
  3. 3Picture this
  4. 4A small example
  5. 5Common beginner questions
  6. 6What comes next
  7. 7Practice

The idea

Apache Airflow is the most popular open-source orchestrator for data pipelines. It uses a specific vocabulary to describe the pieces of a DAG. Understanding these terms prepares you for any Airflow job posting or codebase, even though other orchestrators use different names for the same ideas.

The three core Airflow concepts are: Operators (the tasks that do work), Sensors (tasks that wait for a condition), and XCom (a small communication channel between tasks). Every orchestrator has equivalents: Azure Data Factory calls them Activities, Prefect calls them Tasks, Dagster calls them Ops.

This lesson simulates Airflow's vocabulary in Python dictionaries. No actual Airflow scheduler parses a dag.py file here. You practice the classification (what is an operator, what is a sensor, what is XCom) so you can read real Airflow code confidently.

Why this exists

Imagine a daily pipeline that loads orders from a cloud storage bucket into the warehouse. The bucket is populated by an upstream team. Your pipeline should not start loading until the file actually appears. You could poll the bucket in a loop inside your load script, but that mixes 'waiting' logic with 'loading' logic. Airflow separates these concerns: a Sensor handles the waiting, an Operator handles the loading, and XCom passes small metadata between them.

Misusing these concepts causes real outages. Storing a 200 MB DataFrame in XCom overloads the metadata database. A sensor without a timeout blocks the DAG forever. Understanding what each piece is for prevents these production incidents.

Picture this

Operators do the actual work. A PythonOperator runs a Python function. A BashOperator runs a shell command. A BigQueryOperator runs a SQL query. Each operator is one node in the DAG.

Sensors wait for an external condition to become true. A FileSensor checks if a file exists. An ExternalTaskSensor checks if another DAG's task has succeeded. Sensors 'poke' (check the condition) on an interval, like checking if the mail has arrived every 5 minutes.

XCom (short for cross-communication) is a small key-value store. One task can push a value (like a filename or a row count), and a downstream task can pull it. XCom is designed for tiny metadata, not for data transfer.

Operators, Sensors, and XCom
WorkWaitXCom: pass a small noteOperator: do the workPythonOperatorBashOperatorBigQueryOperatorSensor: wait for a signalFileSensorExternalTaskSensorHttpSensor

Three building blocks with different purposes. Mixing them up causes outages.

Apache Airflow pipeline orchestration architecture showing scheduler, worker execution, and metadata coordination
Production Airflow pipeline architecture: the scheduler coordinates task execution across worker nodes based on DAG dependency graphs with metadata tracking for retries and backfills.
Source: Wikimedia Commons / WikitechCC BY-SA 4.0

A small example

Run the example below in this tab. Read the input, follow the code, then check the output matches what you expect.

Python
# Simulating Airflow concepts as Python dicts
tasks = [
    {"id": "wait_for_file", "kind": "sensor", "target": "gs://bucket/orders/"},
    {"id": "load_orders", "kind": "operator", "action": "COPY INTO orders"},
    {"id": "transform", "kind": "operator", "action": "INSERT INTO gold SELECT ..."},
    {"id": "notify", "kind": "operator", "action": "send Slack message"},
]

# XCom: a tiny value passed between tasks
xcom_store = {}

def push_xcom(task_id, key, value):
    xcom_store[(task_id, key)] = value

def pull_xcom(task_id, key):
    return xcom_store.get((task_id, key))

# After load_orders runs, push the row count
push_xcom("load_orders", "row_count", 1847)

# The notify task can pull that value
count = pull_xcom("load_orders", "row_count")
print(f"Loaded {count} rows")  # "Loaded 1847 rows"

The sensor's job is to wait. It does not transform data. When it senses the condition is true (the file exists), it marks itself as successful and the downstream operator can begin.

ConceptPurposeCommon trap
OperatorRun an extract, load, or transformRetrying a non-idempotent write duplicates data
SensorWait for a file, partition, or external signalNo timeout means the DAG hangs forever
XComPass a filename, row count, or status flagStoring a DataFrame (too large) kills the metadata DB
Python
# Classifying tasks by kind
def get_sensors(tasks):
    """Return IDs of tasks whose kind is 'sensor'."""
    return [t["id"] for t in tasks if t["kind"] == "sensor"]

print(get_sensors(tasks))  # ["wait_for_file"]

Common beginner questions

How big can XCom be?

The practical limit depends on your metadata database. Airflow stores XCom in a database table, so anything over a few kilobytes slows things down. Never push a DataFrame or a large JSON blob. Push a file path or a row count instead.

What happens if a sensor never succeeds?

Without a timeout, the sensor task stays in a 'running' state indefinitely, blocking the entire DAG. Always set a timeout (like timeout=3600 for one hour) and a poke_interval (like poke_interval=300 for every 5 minutes).

Why not just put the wait logic inside the operator?

Separation of concerns. A sensor failing means 'the data is not ready yet.' An operator failing means 'the processing broke.' Different problems need different responses.

Sensors need a timeout

A sensor without a timeout is a DAG that might hang forever. Always set timeout and poke_interval. If the condition never becomes true, the task should fail, not wait eternally.

Same pattern, different names

Azure Data Factory calls operators 'Activities.' Prefect calls them 'Tasks.' Dagster calls them 'Ops' and 'Assets.' The three-way split (do work, wait, pass metadata) is universal.

What comes next

The next lesson covers retries: when a task fails, how do you decide whether to try again or fail immediately? You will implement retry logic with backoff and learn which tasks should retry and which should not.

Practice

Given the tasks list, return the IDs of all tasks whose kind is 'sensor'.

Practicals · load into the editor

After you read the theory, run these in the pane on the right. They execute in this tab, no cluster.

Practice this

Same ideas as interview drills. These challenges open in the studio with a dataset and tests already set up.

  • Order DAG tasks so every dependency runs firstProduction ticket: A scheduler can't just run tasks in the order they were declared - it has to respect the dependency graph.Studiointermediatepython15 minPro
Rate:
Was this useful?
Tasks, edges, topological orderRetries in the DAG runner