skip to content

Executors and Deployment

How and where tasks actually run — Local, Celery or Kubernetes — plus the scheduler, webserver, triggerer and workers around them. Interviewers ask you to pick one for a given scale and to name the concurrency limits that throttle throughput.

on this pageshow

questions

7

In Airflow, which settings limit how many task instances can run at the same time?

level: middleimportance: must knowfreq 62%

answer

  1. throughput is a minimum, not a sum
  2. limits exist at four different scopes
  3. one of them is shared across DAGs
  4. a task can be worth more than one unit
  5. parallelism, max_active_tasks, pools

basics

~10 s

Airflow throttles at several levels: [core] parallelism per scheduler, max_active_runs and max_active_tasks per DAG, max_active_tis_per_dag per task, pool slots across DAGs, and the worker capacity your executor provides. The tightest one wins.

solid answer

~40 s

Throughput in Airflow is the minimum of several independent caps. **Deployment-wide**: `[core] parallelism` limits concurrently running task instances per scheduler, and real worker capacity — Celery workers × `worker_concurrency`, or what the Kubernetes cluster will schedule — sits underneath it. **Per DAG**: `max_active_runs` caps simultaneous DAG runs and `max_active_tasks` caps running task instances from that DAG (the config defaults are `[core] max_active_runs_per_dag` and `max_active_tasks_per_dag`). **Per task**: `max_active_tis_per_dag` caps copies of one task across all runs, and `max_active_tis_per_dagrun` within a run. **Across DAGs**: a **pool** with N slots limits everything assigned to it, and a task can consume more than one slot with `pool_slots`. When a DAG runs slower than expected, the diagnosis is finding which of these is binding — the UI's pool page and the task instance's queued state tell you which.

code

python · 18 lines
python
from airflow.decorators import dag, task


@dag(
    schedule="@hourly",
    catchup=False,
    max_active_runs=1,     # one run of this DAG at a time
    max_active_tasks=8,    # at most 8 of its tasks running
)
def warehouse_load():
    @task(pool="warehouse", pool_slots=2, max_active_tis_per_dag=1)
    def load():
        ...

    load()


warehouse_load()

go deeper

for a junior

Know that Airflow limits concurrency in more than one place, and that a DAG can set max_active_runs and max_active_tasks.

for a middle

Explain each scope — deployment, DAG, task, pool — and be able to say which one is binding for a given symptom.

for a senior

Demonstrate diagnosis: reading the pools page, spotting queued tasks against idle workers, and deciding between raising a cap and adding capacity.

for a principal

Own the policy: which shared resources get pools, what the deployment-wide defaults should be so one team cannot starve another, and how limits interact with autoscaling.

## Why there are so many limits Airflow deliberately layers throttles so that one runaway DAG cannot exhaust a cluster and one shared resource — a warehouse, an API with a rate limit, a database — cannot be hammered by unrelated pipelines. The consequence is that observed throughput is the **minimum** of every applicable cap, and tuning means finding the binding one rather than raising all of them. ## Deployment-wide `[core] parallelism` is the maximum number of task instances a scheduler will have running at once across all DAGs (32 in a stock configuration). Raising it does nothing if you have no workers to absorb the work, and lowering it is a blunt but effective circuit breaker. Underneath sits actual capacity. With `CeleryExecutor` that is the number of worker processes multiplied by `[celery] worker_concurrency`. With `LocalExecutor` it is the scheduler machine's CPU and memory. With `KubernetesExecutor` there is no fixed worker pool at all — capacity is whatever the cluster's quotas and nodes allow, so an Airflow-side limit and a Kubernetes ResourceQuota can each be the binding constraint. ## Per DAG Two DAG-level arguments matter most: - `max_active_runs` — how many runs of *this* DAG may be in flight. Setting it to 1 serialises the DAG, which is the standard defence when a catchup or backfill would otherwise launch dozens of runs against the same target table. - `max_active_tasks` — how many task instances of this DAG may run concurrently, across all of its active runs. The deployment-level defaults for these are `[core] max_active_runs_per_dag` and `[core] max_active_tasks_per_dag`. Note the historical naming: the DAG argument used to be `concurrency` and the config `dag_concurrency`; both were renamed in Airflow 2.2, and older blog posts still use the old spellings. ## Per task `max_active_tis_per_dag` on an operator limits how many instances of that *one* task run at once across all DAG runs — the right tool when a single task talks to a fragile system and you are backfilling many intervals. `max_active_tis_per_dagrun` does the same within a single run, which matters for mapped tasks. ## Pools — the cross-DAG lever A **pool** is a named bucket of slots defined in the UI or CLI. Any task assigned `pool="warehouse"` must acquire a slot before it can leave the `queued` state, and releases it when it finishes. Because pools are global, they are the only mechanism that limits concurrency *across* DAGs owned by different teams — exactly what you want in front of a shared warehouse or a vendor API. A task may declare `pool_slots=2` (or more) to represent that it is heavier than its neighbours, so a 10-slot pool admits five of them rather than ten. Tasks waiting on a pool show as `queued` (`scheduled` in some versions) with the pool named on the task instance detail page, and the Pools page shows used, queued and free slots — the fastest way to confirm a pool is the bottleneck. ## Queues are not limits A Celery `queue` routes a task to a subset of workers; it does not cap concurrency by itself, though the capacity of the workers subscribed to that queue effectively does. Confusing an Airflow **pool** (a concurrency semaphore) with a Celery **queue** (a routing key) is one of the most common interview mistakes on this topic. ## Diagnosing When a pipeline is slower than expected, walk the layers from the outside in: is the pool full? Is the DAG at `max_active_tasks` or `max_active_runs`? Is `parallelism` reached (visible as many task instances queued while workers are idle)? Or is real worker capacity exhausted (workers pinned, no free Celery slots, or pods stuck `Pending`)? Only the last one is fixed by adding machines; the others are configuration. ## A worked shape A DAG with `max_active_runs=1` and `max_active_tasks=8`, whose loader tasks sit in a 4-slot `warehouse` pool with `pool_slots=2`, will never run more than two loaders at once no matter how large the Celery fleet is — because the pool, not the fleet, is binding.

  • What is the difference between an Airflow pool and a Celery queue?
    A pool is a concurrency semaphore: a named set of slots that any task in any DAG must acquire before running, so it caps parallelism across the whole deployment. A Celery queue is a routing key: it decides which worker machines are eligible to pick the task up. A pool limits how many run; a queue decides where they run. They are often used together — GPU tasks routed to a `gpu` queue and also capped by a `gpu` pool.
  • You raise `[core] parallelism` from 32 to 128 and nothing gets faster. What are the likely reasons?
    Either another cap is binding — a pool, `max_active_tasks` on the DAG, `max_active_runs`, or a per-task limit — or there is no real capacity to absorb the extra work: Celery workers are already saturated at their `worker_concurrency`, the LocalExecutor host is out of CPU, or Kubernetes cannot schedule more pods. Check whether task instances are queued while workers sit idle; that pattern points at configuration rather than capacity.
  • Why do teams set `max_active_runs=1` on DAGs that write to the same table?
    Because catchup, a backfill, or a slow run overlapping the next schedule can otherwise put several runs in flight at once, and two runs writing the same target concurrently produce duplicated or interleaved rows. Serialising runs is the cheapest guarantee that the table has one writer, at the cost of the DAG falling behind rather than parallelising when it does.

saying these in an interview costs you the question

  • Treats a Celery queue as if it capped concurrency
  • Thinks raising parallelism alone increases throughput
  • Confuses max_active_runs with max_active_tasks
  • Assumes pools are scoped to one DAG
  • Cannot say where a task waits when a pool is full

context

open as a page

In Airflow, how do LocalExecutor, CeleryExecutor and KubernetesExecutor differ in where tasks run?

level: middleimportance: must knowfreq 80%

basics

~10 s

LocalExecutor runs tasks as subprocesses on the scheduler's machine, CeleryExecutor sends them through a broker to separate long-lived worker processes on other machines, and KubernetesExecutor creates one short-lived pod per task instance.

open as a page

In Airflow, which long-running processes make up a deployment and what does each do?

level: juniorimportance: should knowfreq 62%

basics

~20 s

An Airflow deployment runs a scheduler that decides what should run, workers that execute task code, a web UI/API process, a triggerer for deferred tasks, and a metadata database that all of them read and write.

open as a page

In Airflow, how does KubernetesExecutor differ from using KubernetesPodOperator?

level: middleimportance: should knowfreq 48%

basics

~20 s

KubernetesExecutor is a deployment-wide choice that runs every Airflow task in its own worker pod. KubernetesPodOperator is a single operator, usable under any executor, that launches an arbitrary pod of your choosing and waits for it.

open as a page

Airflow task instances sit in the queued state for hours and never start — how do you diagnose it?

level: seniorimportance: should knowfreq 55%

basics

~20 s

Queued means Airflow's scheduler handed the task to the executor but nothing picked it up. Check in order: pool and concurrency limits, whether any worker consumes the task's queue, worker or broker health, and for Kubernetes whether the pod is stuck Pending.

open as a page

How would you choose between CeleryExecutor and KubernetesExecutor for a team's Airflow platform?

level: principalimportance: should knowfreq 40%

basics

~10 s

Characterise the workload first. Many short tasks and stable dependencies favour CeleryExecutor's warm workers; heterogeneous resource needs, per-team images and strong isolation favour KubernetesExecutor's pod-per-task, which costs seconds of startup on every task.

open as a page

In Airflow 2, how can several schedulers run at once without duplicating task instances?

level: seniorimportance: nice to knowfreq 30%

basics

~20 s

Airflow 2 supports multiple active schedulers with no leader election. They coordinate through the metadata database, taking row-level locks with SELECT ... FOR UPDATE so only one scheduler can claim and queue a given task instance.

open as a page