Menu

Airflow course · Lesson 8 of 8

Airflow Executors, Pools and Scaling

Compare Airflow 3 executors (Local, Celery, Kubernetes, multiple executors), learn what replaced SequentialExecutor and hybrids, and tune pools and scaling.

  • Advanced
  • 17 min read
  • Updated Oct 2026
On this page
  1. Reading the defaults
  2. SequentialExecutor (removed)
  3. What it was
  4. Why it went
  5. In interviews
  6. LocalExecutor
  7. What it is
  8. How it works
  9. Pitfalls
  10. In interviews
  11. CeleryExecutor
  12. What it is
  13. How it works
  14. Pitfalls
  15. In interviews
  16. KubernetesExecutor
  17. What it is
  18. How it works
  19. Pitfalls
  20. In interviews
  21. CeleryKubernetes hybrid and multiple executors
  22. What it was
  23. Why it went
  24. In interviews
  25. Parallelism and pools
  26. What it is
  27. How it works
  28. Pitfalls
  29. In interviews
  30. Resource management
  31. What it is
  32. How it works
  33. Pitfalls
  34. In interviews
  35. Scaling Airflow
  36. What it is
  37. How it works
  38. Pitfalls
  39. In interviews
  40. Practice questions
  41. Key takeaways

The executor decides where and how task instances run: in local processes, on a fleet of Celery workers, or in Kubernetes pods. Concurrency settings and pools decide how many run at once. Together they determine whether Airflow keeps up with your workload or quietly builds a backlog. This lesson covers every executor an interviewer may ask about, including the ones Airflow 3 removed, and how to scale each component.

Reading the defaults

The concurrency numbers in this lesson come from a default Airflow 3.3 configuration. You can read any setting the same way, or with airflow config get-value core parallelism on the command line:

from airflow.configuration import conf

for section, key in [
    ("core", "executor"),
    ("core", "parallelism"),
    ("core", "max_active_tasks_per_dag"),
    ("core", "max_active_runs_per_dag"),
    ("core", "default_pool_task_slot_count"),
    ("celery", "worker_concurrency"),
    ("scheduler", "task_queued_timeout"),
    ("triggerer", "capacity"),
]:
    print(f"[{section}] {key} = {conf.get(section, key)}")
[core] executor = LocalExecutor
[core] parallelism = 32
[core] max_active_tasks_per_dag = 16
[core] max_active_runs_per_dag = 16
[core] default_pool_task_slot_count = 128
[celery] worker_concurrency = 16
[scheduler] task_queued_timeout = 600.0
[triggerer] capacity = 1000

Every setting can also be set as an environment variable AIRFLOW__<SECTION>__<KEY>, for example AIRFLOW__CORE__PARALLELISM=64, which is how containerised deployments usually configure Airflow.

SequentialExecutor (removed)

What it was

The SequentialExecutor ran one task instance at a time, in a subprocess of the scheduler. It was the out-of-the-box executor in Airflow 1 and 2 because it was the only one that worked with SQLite, which cannot handle concurrent writers. Tutorials and airflow standalone used it, so many people’s first Airflow ran everything one task at a time.

Why it went

It was removed in Airflow 3.0, together with the DebugExecutor. The upgrade guide’s reasoning: it was only useful for local testing, and the LocalExecutor now works with SQLite (using its write-ahead-log mode) while running tasks in parallel. The replacements:

  • LocalExecutor is the new default ([core] executor = LocalExecutor, as the output above shows), including for local development.
  • dag.test() (or airflow dags test) replaces the DebugExecutor for running a DAG in one process under a debugger.

In interviews

“What is the SequentialExecutor and when would you use it?” Answer: it ran one task at a time and was the default only because of SQLite; it was removed in Airflow 3, LocalExecutor is the default now, and you use dag.test() for debugging. Mentioning that SQLite is still only for development, not production, avoids a follow-up trap.

LocalExecutor

What it is

The LocalExecutor runs task instances as parallel processes on the same machine as the scheduler. It needs no broker or cluster, which makes it the simplest real executor.

How it works

[core]
executor = LocalExecutor
parallelism = 32            # max task instances running at once for this scheduler
[database]
sql_alchemy_conn = postgresql+psycopg2://airflow:***@db.internal:5432/airflow
  • The scheduler starts worker processes and hands them queued tasks, up to parallelism.
  • Use PostgreSQL or MySQL in production; SQLite is for development.
  • Scaling is vertical: a bigger machine. With several HA schedulers, each runs its own local workers, so capacity grows with schedulers, but tasks run where the scheduler runs.

Pitfalls

  • Heavy tasks on the scheduler machine starve the scheduler of CPU and memory, which slows scheduling for everyone. Push heavy work out to Spark, the warehouse or pods.
  • No isolation between tasks: they share the machine’s Python environment and resources.
  • A scheduler restart kills running tasks on that machine (they are retried if they have retries).

In interviews

LocalExecutor is the right answer for small teams and single-machine deployments, and for development. Say when you would outgrow it: when tasks need more CPU or memory than one box, need different dependencies, or when you need to scale workers independently of the scheduler.

CeleryExecutor

What it is

The CeleryExecutor sends tasks through a message broker (Redis or RabbitMQ) to a pool of long-running Celery workers on other machines. It ships in the apache-airflow-providers-celery package.

How it works

[core]
executor = CeleryExecutor

[celery]
broker_url = redis://redis.internal:6379/0
result_backend = db+postgresql://airflow:***@db.internal:5432/airflow
worker_concurrency = 16       # task slots per worker process
airflow celery worker                      # serves the default queue
airflow celery worker --queues heavy,ml    # a worker class for big tasks
airflow celery flower                      # optional monitoring UI
  • The scheduler puts a message on the broker for each queued task; any worker listening on that task’s queue picks it up and runs it.
  • Each task’s queue argument (default default) routes it to a class of workers, which is how you send memory-hungry tasks to large machines.
  • Capacity is roughly workers × worker_concurrency, capped by parallelism per scheduler and by pools.
  • Workers stay up between tasks, so start-up latency is low, and you scale by adding workers (the official Helm chart can autoscale Celery workers with KEDA based on queued tasks).

Pitfalls

  • The broker becomes critical infrastructure. If Redis loses messages or fills up, tasks sit in queued.
  • Workers need every dependency every DAG uses, which leads to large, conflicting images; isolate with queues, virtualenv operators or pods.
  • A task queue that no worker listens to: tasks stay queued until task_queued_timeout (600 seconds by default) fails them.
  • Long-running tasks block a worker slot; use deferrable operators for waiting.
  • Visibility timeouts in the broker shorter than your longest task can cause duplicate execution; check your broker settings.

In interviews

Expect “How does CeleryExecutor work?” Explain scheduler, broker, workers and the result backend; queues for routing; capacity as workers × concurrency; and failure modes (broker down, unserved queues). Compare it with Kubernetes: Celery has warm workers and low latency but shared environments; Kubernetes gives isolation per task at the cost of pod start-up time.

KubernetesExecutor

What it is

The KubernetesExecutor launches one pod per task instance in a Kubernetes cluster. The pod runs Airflow’s task runner for that single task and exits. It ships in the apache-airflow-providers-cncf-kubernetes package.

How it works

[core]
executor = KubernetesExecutor

[kubernetes_executor]
namespace = airflow
pod_template_file = /opt/airflow/pod_templates/default.yaml
worker_pods_creation_batch_size = 16      # pods the scheduler creates per loop (default 1)
delete_worker_pods = True

Per-task resources and images through executor_config:

from kubernetes.client import models as k8s

@task(executor_config={
    "pod_override": k8s.V1Pod(spec=k8s.V1PodSpec(containers=[
        k8s.V1Container(
            name="base",                      # must be "base" to override the task container
            image="registry.example.com/airflow-ml:3.3.2-1",
            resources=k8s.V1ResourceRequirements(requests={"cpu": "2", "memory": "8Gi"},
                                                 limits={"memory": "8Gi"}),
        )
    ]))
})
def train_model():
    ...
  • No idle workers: capacity scales with the cluster, and each task is isolated, with its own resources and optionally its own image.
  • Start-up cost per task (scheduling the pod, pulling the image, starting Python and loading the DAG) is typically several seconds or more, so many tiny tasks are slow and expensive.
  • Logs must be shipped to remote storage because the pod disappears when the task ends.

Pitfalls

  • Image pull time dominating short tasks; keep images small and cached on nodes.
  • Pods pending forever because of quotas or no capacity; tasks then sit in queued.
  • Confusing it with KubernetesPodOperator. The executor runs every Airflow task in a pod; KPO is one task that runs your container, with any executor.
  • Forgetting remote logging, so logs vanish with the pod.

In interviews

Explain one pod per task, pod_template_file and pod_override for resources, the start-up latency trade-off, and when you would choose it (spiky workloads, strong isolation, teams with different dependencies). Be ready to contrast it with KPO.

CeleryKubernetes hybrid and multiple executors

What it was

The CeleryKubernetesExecutor (and the LocalKubernetesExecutor) combined two executors in one class: tasks on a special queue (by default kubernetes) went to Kubernetes pods, all others to Celery workers. Teams used it to keep fast, frequent tasks on warm Celery workers while sending heavy or unusual tasks to isolated pods.

Why it went

Both hybrids are not supported in Airflow 3.0 and later. The provider documentation gives two reasons: they misused each task’s queue field to pick the sub-executor, which made queues unusable for their real purpose, and each combination needed its own hand-written class. Since Airflow 2.10 you can instead configure multiple executors concurrently:

[core]
# Comma-separated; the first one is the default for tasks that do not choose
executor = CeleryExecutor,KubernetesExecutor
# Aliases are allowed: executor = CeleryExecutor,my_company.executors.GpuExecutor:gpu
@task                                   # runs on the default executor (CeleryExecutor)
def small_step(): ...

@task(executor="KubernetesExecutor")    # this task only runs in its own pod
def heavy_step(): ...

with DAG("mixed", default_args={"executor": "KubernetesExecutor"}, ...):   # or per DAG
    ...

A task’s queue keeps its normal meaning (Celery routing) and the executor argument picks where it runs. Other executors can join the list too, such as the AWS ECS, Batch and Lambda executors in the Amazon provider or the EdgeExecutor for remote workers.

In interviews

If asked about CeleryKubernetesExecutor, describe what it did and why teams liked it, then say it was dropped in Airflow 3 in favour of the multiple-executors configuration, where each task or DAG chooses an executor by name. That shows you know both the classic answer and the current one.

Parallelism and pools

What it is

Several limits work together to decide how many task instances may run. Pools are named sets of slots you create to protect a shared resource, such as “no more than 4 concurrent queries against the warehouse”.

How it works

Limit Scope Default (3.3)
[core] parallelism running task instances per scheduler 32
max_active_tasks on the DAG ([core] max_active_tasks_per_dag) running tasks across all runs of one DAG 16
max_active_runs on the DAG ([core] max_active_runs_per_dag) concurrent DAG runs of one DAG 16
max_active_tis_per_dag on a task running instances of one task across runs unlimited
max_active_tis_per_dagrun on a task running instances of one task in a run (useful for mapped tasks) unlimited
Pools (pool, pool_slots on a task) tasks sharing a named slot budget default_pool with 128 slots
Executor capacity Celery worker_concurrency × workers, cluster size depends

The scheduler picks scheduled tasks in priority order in its critical section. Each task’s effective priority comes from priority_weight and the weight_rule: downstream (the default) adds up the weights of all downstream tasks, so tasks early in long chains go first; upstream does the opposite; absolute uses the task’s own weight.

from datetime import datetime
from airflow.sdk import DAG
from airflow.models.pool import Pool
from airflow.providers.standard.operators.empty import EmptyOperator
from airflow.utils.session import create_session

# Same as: airflow pools set warehouse 4 "Max concurrent warehouse queries"
Pool.create_or_update_pool("warehouse", slots=4, description="Max concurrent warehouse queries",
                           include_deferred=False)

with DAG("concurrency_demo", schedule="@hourly", start_date=datetime(2026, 1, 1), catchup=False,
         max_active_runs=1, max_active_tasks=4) as concurrency_dag:
    extract = EmptyOperator(task_id="extract", pool="warehouse", pool_slots=2)   # a heavy query: 2 slots
    transform = EmptyOperator(task_id="transform", priority_weight=5)
    load = EmptyOperator(task_id="load", priority_weight=10, weight_rule="absolute",
                         max_active_tis_per_dag=1)
    extract >> transform >> load

print("DAG limits:", concurrency_dag.max_active_runs, "run(s),", concurrency_dag.max_active_tasks, "tasks")
for t in concurrency_dag.tasks:
    print(f"{t.task_id:10} pool={t.pool:12} slots={t.pool_slots} weight={t.priority_weight} rule={t.weight_rule}")
with create_session() as session:
    for p in session.query(Pool).order_by(Pool.pool):
        print(f"pool {p.pool:12} slots={p.slots} include_deferred={p.include_deferred}")
DAG limits: 1 run(s), 4 tasks
extract    pool=warehouse    slots=2 weight=1 rule=WeightRule.DOWNSTREAM
transform  pool=default_pool slots=1 weight=5 rule=WeightRule.DOWNSTREAM
load       pool=default_pool slots=1 weight=10 rule=absolute
pool default_pool slots=128 include_deferred=False
pool warehouse    slots=4 include_deferred=False

With 4 slots in warehouse and 2 per extract, at most two extract tasks (from different DAGs using the pool) query the warehouse at once, however many workers are free. include_deferred=True makes deferred tasks keep counting against the pool.

Pitfalls

  • Raising parallelism without raising executor capacity (or the reverse): the smaller number wins.
  • Pools that are smaller than a task’s pool_slots: the task can never run.
  • Forgetting that every task is in default_pool, which caps the whole installation at 128 running tasks unless resized.
  • Sensors in poke mode occupying pool slots; put them in their own pool or use deferrable mode.
  • max_active_runs=1 plus catchup: history is processed one run at a time, which may be what you want (ordered loads) or a very long backfill.

In interviews

Be ready to list the limits from global to per-task, explain pools as “protect a shared external resource”, and explain priority weights. A typical scenario: “Twenty DAGs hammer the same database at 02:00.” Answer: a pool sized to what the database can take, priority weights for the critical DAGs, and staggered schedules.

Resource management

What it is

Resource management is making sure each task gets the CPU, memory and software it needs without starving others, and that Airflow’s own components keep the resources they need.

How it works

  • Do not process data on workers. Submit work to Spark, the warehouse or a container, and let the Airflow task wait (deferrably, ideally). Then workers need little CPU or memory.
  • Celery: route by queue to worker classes sized for the work (default, heavy, gpu), each started with --queues and its own concurrency.
  • Kubernetes: set requests and limits in the pod_template_file and per task with executor_config={"pod_override": ...}; for tasks that run your own image, use KubernetesPodOperator with container_resources.
  • Python dependencies: isolate conflicting libraries with @task.virtualenv, @task.external_python, Docker or KPO rather than installing everything on every worker.
  • Time limits: execution_timeout per task and dagrun_timeout per DAG stop runaway work.
  • External systems: pools cap concurrency against databases and APIs; deferrable operators avoid holding slots while waiting.
  • The resources argument on operators exists, but most executors ignore it, so use executor-specific mechanisms instead.

Pitfalls

  • One giant worker image with every library: slow to build, slow to pull, and full of version conflicts.
  • No memory limits on pods, so one task’s leak gets the node’s other pods evicted.
  • Scheduler, DAG processor and workers on the same small machine; one busy component slows the others.

In interviews

“A task keeps getting OOM-killed. What do you do?” Check whether the task processes data that belongs in an external engine; if it must run in Airflow, give it a larger queue or a pod override with higher memory limits, and add execution_timeout. Mentioning dependency isolation options shows practical experience.

Scaling Airflow

What it is

Scaling Airflow means scaling each component for its own bottleneck: parsing (DAG processor), scheduling decisions (scheduler and database), task execution (executor and workers), waiting (triggerer) and serving users and workers (API server).

How it works

Component Scale by Watch
DAG processor more [dag_processor] parsing_processes (2 by default), more DAG processor instances, faster DAG files parse time per file, import errors
Scheduler several schedulers (HA, active-active), PostgreSQL or MySQL 8 scheduler loop duration, tasks starving, scheduled backlog
Metadata database a bigger instance, connection pooling (for example PgBouncer), cleaning old rows with airflow db clean connections, slow queries
Workers more Celery workers (autoscale on queue length) or more cluster capacity queued backlog, open slots
Triggerer more triggerer replicas; each handles up to [triggerer] capacity triggers (1000 by default) triggerer capacity left, blocked event loop
API server more replicas behind a load balancer (it also serves the Task Execution API in Airflow 3) latency, errors from workers

Practical order of work when Airflow is slow:

  1. Find which state tasks pile up in. Many scheduled tasks point at concurrency limits or the scheduler; many queued tasks point at executor capacity.
  2. Fix DAG files first: remove top-level code, split huge files, avoid thousands of tiny tasks. This is usually the cheapest win.
  3. Raise the right limit: pool sizes, parallelism, max_active_tasks, worker counts.
  4. Scale components horizontally: more schedulers, DAG processors, triggerers, workers.
  5. Keep the metadata database healthy: pooling, enough connections, regular airflow db clean for old task instances, logs and XComs.

The official Helm chart deploys each component as its own Kubernetes workload with replica settings, which is the usual way to scale a self-managed Airflow on Kubernetes. Managed services (Amazon MWAA, Google Cloud Composer, Astronomer) expose similar knobs.

Pitfalls

  • Adding schedulers when the bottleneck is the database or the DAG files.
  • Scaling workers to hundreds while parallelism stays at 32 per scheduler.
  • Running SQLite or a tiny database instance in production.
  • Measuring nothing: use the metrics in the operations lesson (scheduler loop duration, pool open slots, executor queued tasks).

In interviews

“Airflow is getting slow as we add DAGs. How would you scale it?” Walk through the components, say how you would find the bottleneck (task state backlog and metrics), fix DAG parsing first, then scale horizontally, and keep the database healthy. Mentioning that HA schedulers rely on database row locks, so they need PostgreSQL or MySQL 8, shows depth.

Practice questions

Which executor would you choose for a 5-person team running 30 DAGs on one VM? For a platform team serving 40 teams?

The small team: LocalExecutor with PostgreSQL, keeping heavy work in external systems. The platform team: CeleryExecutor for low-latency common tasks and KubernetesExecutor (through the multiple-executors setting) for isolated or heavy tasks, or KubernetesExecutor alone if isolation matters more than start-up latency.

Why was the SequentialExecutor removed, and what do you use instead?

It ran one task at a time and existed mainly because SQLite could not handle concurrent writers. In Airflow 3 the LocalExecutor works with SQLite and runs tasks in parallel, so it became the default. For debugging a DAG in a single process, use dag.test() or airflow dags test.

How do you send only your ML training tasks to Kubernetes while everything else runs on Celery in Airflow 3?

Configure [core] executor = CeleryExecutor,KubernetesExecutor (the first is the default) and set executor="KubernetesExecutor" on the training tasks, with executor_config pod overrides for resources. In Airflow 2 this used the CeleryKubernetesExecutor and the kubernetes queue, which is not supported in Airflow 3.

Tasks pile up in scheduled while workers are idle. What limits do you check?

Concurrency limits applied by the scheduler: max_active_tasks on the DAG, max_active_tis_per_dag/per_dagrun on tasks, pool slots (including default_pool at 128), parallelism per scheduler, and depends_on_past. Idle workers with a scheduled backlog means the executor is not the bottleneck.

How do you stop many DAGs from overloading a shared database?

Create a pool sized to what the database can handle and assign all tasks that query it to the pool (with pool_slots larger than 1 for heavier queries). Use priority weights so critical DAGs go first, and stagger schedules.

What are the trade-offs between CeleryExecutor and KubernetesExecutor?

Celery: always-on workers, low start-up latency and simple per-task cost, but shared environments, a broker to operate and idle capacity to pay for. Kubernetes: one pod per task with its own resources and image and elastic capacity, but pod start-up latency per task, image management and the need for remote logging.

Key takeaways

  • The executor decides where tasks run; Airflow 3’s default is LocalExecutor, and Celery and Kubernetes executors come from providers.
  • SequentialExecutor and DebugExecutor were removed in 3.0; use LocalExecutor and dag.test().
  • CeleryKubernetesExecutor and LocalKubernetesExecutor are gone in Airflow 3; list several executors in [core] executor and choose per task with executor=.
  • Concurrency is the minimum of parallelism, DAG and task limits, pool slots and executor capacity; pools protect shared systems and priority weights order the queue.
  • Scale by component: fix DAG parsing first, then add schedulers, DAG processors, triggerers and workers, and keep the metadata database healthy.

By DataDank Editorial · Last reviewed Oct 2026 · Configuration defaults were read from an Apache Airflow 3.3.2 installation, and the pool and concurrency examples ran on it with a SQLite metadata database. Executor, Celery, Kubernetes and Helm configuration is written from the Airflow and provider documentation and was not executed (it needs a broker, a cluster or several machines).

Progress is saved in this browser only. No account needed.

Search
Filter by type