Skip to content

Scheduler bulk-fails all scheduled tasks when serialized DAG is transiently missing #62050

Description

@vincent-heng

Apache Airflow version

Other Airflow 3 version (please specify below)

If "Other Airflow 3 version" selected, which one?

3.0.6, still on main

What happened?

My instance has been hitting intermittent task failures on MWAA (Airflow 3.0.6, ~150 DAGs, 4 schedulers). Tasks failed in bulk with no obvious cause but succeeded on manual retry. I noticed this on scheduler_job_runner.py:

When the scheduler can't find a DAG in the serialized_dag table, it does this:

session.execute(
    update(TI)
    .where(TI.dag_id == dag_id, TI.state == TaskInstanceState.SCHEDULED)
    .values(state=TaskInstanceState.FAILED)
    .execution_options(synchronize_session="fetch")
)

It sets every SCHEDULED task instance for that DAG to FAILED.

With PR #58259 and #56422, it probably happens less often but the bulk-failure issue has never been addressed

What you think should happen instead?

The scheduler could skip scheduling that DAG for the current iteration and try again next time, instead of immediately failing everything.

I thought of logging a warning instead of error, tracking a counter per DAG, and only failing tasks after several consecutive misses to distinguish transient gaps from genuinely missing DAGs. What do you think about this solution?

PR #55126 tried something similar for stale DAGs (skipping and continuing)

How to reproduce

def test_missing_serialized_dag_bulk_fails(self, dag_maker, session):
    dag_id = "SchedulerJobTest.test_missing_serialized_dag_bulk_fails"

    with dag_maker(dag_id=dag_id, session=session):
        EmptyOperator(task_id="task_a")
        EmptyOperator(task_id="task_b")

    scheduler_job = Job()
    self.job_runner = SchedulerJobRunner(job=scheduler_job)

    # Simulate serialized DAG being transiently missing
    self.job_runner.scheduler_dag_bag = mock.MagicMock()
    self.job_runner.scheduler_dag_bag.get_dag_for_run.return_value = None

    dr = dag_maker.create_dagrun(state=DagRunState.RUNNING, run_type=DagRunType.SCHEDULED)
    for ti in dr.task_instances:
        ti.state = State.SCHEDULED
        session.merge(ti)
    session.flush()

    res = self.job_runner._executable_task_instances_to_queued(max_tis=32, session=session)
    session.flush()

    assert len(res) == 0
    tis = dr.get_task_instances(session=session)  # Both tasks are FAILED instead of SCHEDULED
    for ti in tis:
        print(f"{ti.task_id}: {ti.state}")

Operating System

AWS MWAA

Versions of Apache Airflow Providers

n/a

Deployment

Amazon (AWS) MWAA

Deployment details

  • MWAA environment running 3.0.6 (latest available on MWAA as of Feb 2026)
  • mw1.2xlarge instance class
  • 4 schedulers
  • 5-20 workers (auto-scaling)
  • ~150 DAGs
  • min_file_process_interval=300

Anything else?

would appreciate guidance on the preferred approach (retry counter or just ignoring)

Are you willing to submit PR?

  • Yes I am willing to submit a PR!

Code of Conduct

Activity

  1. boring-cyborg commented on Feb 16, 2026

    @boring-cyborg

    Thanks for opening your first issue here! Be sure to follow the issue template! If you are willing to raise PR to address this issue please do so, no need to wait for approval.

  2. added 9 commits that reference this issue on Mar 4, 2026
    1193158
    b692da6
    a971ae0
    c0bc313
    c38181d
    6202b33
    5ddfc95
    6aeb017
    94bfda0
  3. 24 remaining items

  4. added 2 commits that reference this issue on Apr 13, 2026
    bb88ce7
    c6c98b2
  5. added and removed
    needs-triagelabel for new issues that we didn't triage yet
    on May 11, 2026
  6. Vamsi-klu commented on Aug 29, 2026

    @Vamsi-klu
  7. rjgoyln commented on Aug 29, 2026

    @rjgoyln
    Contributor

    Hi @kaxil, @jscheffl, and @potiuk — regarding the recent closure of #72243 and the concerns raised in #62878, I did a deep dive and reproduced this locally on Postgres 14 against main without mocking. The results actually prove your instincts right, but change the shape of the fix. I'd like to align on the direction before opening a PR.

    1. The current UPDATE is destructive, not a safety net

    When a serialized DAG is missing, the current bulk UPDATE causes:

    • Cross-run blast radius: It wrongly fails SCHEDULED tasks across all runs of the DAG, not just the affected one.
    • Silent UI failures: Tasks go straight to FAILED. It bypasses handle_failure entirely—no retries, no callbacks, no logs.
    • $O(N)$ DB spam: Every TI of the missing DAG triggers its own full-table UPDATE in the batch.
    • Stuck runs: The tasks are destroyed, but the DagRun remains stuck in RUNNING anyway.

    2. The Starvation Trap: Why "pure skip" (#62878 / #72243) is worse

    The current UPDATE accidentally acts as a starvation filter. If we just remove it ("pure skip"), the broken tasks loop forever, locking out healthy DAGs.

    Here is a 2-pass scheduler test with a broken DAG (max_tis=2) and a healthy DAG:

    Scenario PASS 1 queued PASS 2 queued Broken DAG's tasks Healthy DAG's task
    Current main [] [healthy_task] FAILED (Wrongly killed) Queued after 1 loop delay
    Pure Skip [] [] SCHEDULED Never queued (Cluster Starvation!)
    Skip + starved_dags [healthy_task] [] SCHEDULED Queued immediately

    3. The Proposed Fix

    When get_dag_for_run() returns None in _task_concurrency_allows_execution:

    1. Leave TIs as SCHEDULED (stop the destructive UPDATE).
    2. Add dag_id to starved_dags to prevent the cluster starvation shown above.
    3. Track the missing DAG locally to log the error only once per DAG per batch, not once per TI.

    4. Scope & Next Steps

    The remaining concern (from #62878) is that a permanently deleted DAG leaves TIs SCHEDULED forever. Since this is a DagRun-level gap already being addressed by #70056, I plan to scope my PR strictly to the concurrency and starvation issues above.

    Does this plan look good to you? I have the local reproduction ready as a test and can open the PR if you agree with this direction.

  8. added 2 commits that reference this issue on Sep 9, 2026
    495f5b3
    dd58070
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions