skip to content

An Airflow DAG paused two weeks with catchup=True is unpaused and floods the scheduler — how do you control it?

level: seniorimportance: should knowfreq 58%

answer

  1. contain first, then decide what to replay
  2. pausing does not delete queued runs
  3. one cap is per-DAG, one protects the shared resource
  4. max_active_runs, max_active_tasks_per_dag, pools
  5. replay a chosen range, do not flip the flag

basics

~10 s

Pause the DAG again immediately, then decide which intervals must genuinely be replayed. Cap concurrency with max_active_runs, max_active_tasks_per_dag and a pool, drop the rest, and replay the needed window as an explicit bounded backfill.

solid answer

~50 s

First stop the bleeding: pause the DAG, then mark the queued runs you do not want as failed or delete them so the scheduler stops issuing work. Then decide whether the missed intervals actually need replaying — for a partition-overwrite DAG they usually do, for anything that emails customers or calls a paid API they usually do not. Before unpausing, bound the blast radius: `max_active_runs` on the DAG caps concurrent DAG runs, `max_active_tasks_per_dag` caps its concurrent task instances, and a **pool** with limited `pool_slots` protects the shared resource everyone else is contending for, such as the source database. `depends_on_past=True` serialises the replay in interval order when runs are cumulative. Finally replay a bounded window with an explicit backfill rather than a flag flip, and fix the root cause: set `catchup=False` and treat replays as deliberate operations.

code

python · 19 lines
python
import pendulum
from airflow import DAG
from airflow.operators.python import PythonOperator

with DAG(
    dag_id="daily_orders",
    start_date=pendulum.datetime(2024, 1, 1, tz="UTC"),
    schedule="@daily",
    catchup=False,
    max_active_runs=2,
    max_active_tasks_per_dag=8,
    dagrun_timeout=pendulum.duration(hours=3),
) as dag:
    PythonOperator(
        task_id="extract",
        python_callable=extract,
        pool="oltp_replica",     # shared budget across all DAGs
        pool_slots=1,
    )

go deeper

for a junior

Know that unpausing a DAG can create a run for every elapsed interval, and that pausing it again is the first thing to do when that starts happening.

for a middle

Explain the specific caps and what each bounds: max_active_runs for concurrent DAG runs, max_active_tasks_per_dag for its task instances, pools for a shared external resource.

for a senior

Sequence a real incident — contain, assess irreversible side effects, replay a bounded window explicitly, then change the DAG so it cannot recur — and name the knobs rather than saying 'limit concurrency'.

for a principal

Own the guardrails: mandatory explicit catchup and max_active_runs in every DAG, pools in front of every shared dependency, queue-depth alerting, and a backfill policy that says who may replay what and over which windows.

## What actually happened Pausing a DAG does not stop time; it stops **run creation**. The intervals keep elapsing. When you unpause a DAG that has `catchup=True`, the scheduler discovers every closed interval since the last run and starts creating DAG runs for all of them, oldest first. Fourteen days of an hourly DAG is 336 runs; if each has twenty tasks and there is no concurrency cap, that is thousands of queued task instances competing with every other DAG on the cluster. The damage is rarely limited to the offending DAG. Symptoms are cluster-wide: the global `parallelism` ceiling is saturated, unrelated DAGs miss their windows, the source database sees fourteen days of extraction traffic at once, and the metadata database groans under the task-instance churn. If the DAG has external side effects, the damage is worse and irreversible — two weeks of back-dated notifications actually delivered. ## Immediate response 1. **Pause the DAG.** This stops new runs from being created; it does not kill running tasks. 2. **Stop the queued work.** Mark the unwanted queued/running DAG runs failed, or delete them, so their task instances stop being scheduled. Do this before debating policy — every minute of deliberation is more executed runs. 3. **Check for side effects already committed.** Which of the replayed runs sent messages, wrote to a live table, or called a billable API? That determines whether this is an inconvenience or an incident. 4. **Only then** decide what should legitimately be replayed. ## Bounding the replay The controls, from broad to narrow: - `max_active_runs` (DAG level) — the number of DAG runs of *this DAG* allowed to run concurrently. On any DAG that can accumulate a backlog, this should be set deliberately, not left to the cluster default. - `max_active_tasks_per_dag` (DAG level) — the number of concurrently running task instances for the DAG, across all its runs. - **Pools** — a named pool with a fixed number of slots, assigned to the tasks that touch a fragile shared resource. Tasks take `pool_slots` from it and queue when it is exhausted. This is the right instrument when the constraint is external ("the OLTP replica tolerates four concurrent extracts"), because it bounds *everyone's* use of that resource, not just this DAG's. - `parallelism` / worker capacity — the cluster-wide ceiling. Reaching it is the failure mode you are trying to avoid, not a control you should rely on. - `depends_on_past=True` — each task instance waits for the same task in the *previous* DAG run to succeed. This forces strictly ordered replay, which is required when a run's output depends on the prior run's state, and is the single biggest throughput cost you can impose. `wait_for_downstream=True` goes further, also requiring the previous run's downstream tasks to have succeeded. - `dagrun_timeout` — keeps a stuck replay run from occupying a `max_active_runs` slot forever. ```python with DAG( dag_id="daily_orders", start_date=pendulum.datetime(2024, 1, 1, tz="UTC"), schedule="@daily", catchup=False, max_active_runs=2, max_active_tasks_per_dag=8, dagrun_timeout=pendulum.duration(hours=3), ) as dag: ... ``` ## Replay deliberately, not by flag flip Once the DAG is safe, replay the window you actually need as an explicit, bounded operation over a date range you chose — in Airflow 2 with `airflow dags backfill --start-date ... --end-date ... <dag_id>`; Airflow 3 makes backfills scheduler-managed and startable from the UI and API. A backfill of a known range is reviewable, restartable and stoppable in a way that "unpause and hope" is not. Split a very large range into chunks so a failure costs you one chunk, and run the replay against the same concurrency caps. Note that flipping `catchup` back and forth is not a replay mechanism: the scheduler derives the next interval from the *newest existing run*, so turning `catchup=True` on again fills forward, never backward into a gap that a later run has already passed. ## Prevention - Write `catchup=` explicitly in every DAG, and default to `False` for anything with external side effects. - Set `max_active_runs` on every DAG that could ever accumulate a backlog — which is every scheduled DAG. - Keep `start_date` sensible. A `start_date` two years in the past on a DAG with `catchup=True` is a loaded weapon in the repository. - Put a pool in front of every fragile shared dependency, so no single DAG's replay can starve it. - Alert on scheduler queue depth and on "DAG has more than N queued runs", so this is caught in minutes rather than discovered by the downstream team. - Make tasks idempotent — partition overwrite or merge on a key rather than blind append — so that if a replay does happen, the worst outcome is wasted compute rather than duplicated data. ## What the interviewer is listening for The weak answer is "set catchup=False", which prevents the next occurrence but does nothing about the fourteen days currently executing. The strong answer sequences containment, blast-radius assessment (especially side effects), a deliberate bounded replay, and a prevention change — and names the specific knobs rather than gesturing at "limit concurrency".

  • Why is a pool sometimes the right control rather than max_active_runs?
    Because the constraint is often external. max_active_runs bounds one DAG's runs, but if five DAGs all extract from the same OLTP replica, capping each one individually still lets their sum overload it. A pool with a fixed slot count sits in front of the resource, so every task that touches it queues against one shared budget regardless of which DAG it belongs to.
  • What does depends_on_past=True cost you during a backlog drain?
    Throughput. Each task instance waits for the same task in the previous DAG run to succeed, so the replay proceeds strictly one interval at a time and cannot be parallelised. Worse, a single failed historical run blocks everything after it until someone intervenes. Use it only when a run genuinely depends on the previous run's state; otherwise make tasks independent and cap concurrency instead.
  • Turning catchup back to True after the incident — does that recover the intervals you deleted?
    No. The scheduler computes the next interval from the newest existing DAG run, so it fills forward from there and never reaches back into a gap that a later run has already passed. Recovering deleted intervals requires an explicit backfill over that date range, or clearing runs that still exist.
  • How would you detect this class of incident automatically?
    Alert on queue depth and on per-DAG queued-run count: "DAG X has more than N runs queued" catches a catchup flood within minutes. Complement it with alerts on scheduler task-queue latency and on saturation of the pools guarding shared resources, so you see the contention before the downstream team reports missed data.

saying these in an interview costs you the question

  • Answers only 'set catchup=False' without containing the runs already executing
  • Thinks pausing a DAG cancels its already-queued runs
  • Proposes flipping catchup back to True to recover deleted intervals
  • Ignores whether the replayed runs sent emails or called billable APIs
  • Relies on cluster-wide parallelism as the concurrency control

context