Skip to content

fix: skip transiently missing serialized DAGs instead of bulk-failing tasks - #62878

Closed
YoannAbriel wants to merge 1 commit into
apache:mainfrom
YoannAbriel:fix/issue-62050
Closed

YoannAbriel wants to merge 1 commit into
apache:mainfrom
YoannAbriel:fix/issue-62050

Conversation

@YoannAbriel

Copy link
Copy Markdown
Contributor

Problem

On multi-scheduler Airflow deployments (e.g. AWS MWAA with 4+ schedulers), tasks intermittently fail in bulk with no apparent cause, but succeed on manual retry. The root issue is in _executable_task_instances_to_queued: when the scheduler checks task concurrency limits and cannot find the serialized DAG, it immediately bulk-fails all SCHEDULED task instances for that DAG via a raw SQL UPDATE. This is an overly aggressive response to what is often a transient race condition during DAG file parsing or serialization refresh cycles.

Root Cause

In scheduler_job_runner.py, the _executable_task_instances_to_queued method loads the serialized DAG only when checking per-task or per-dagrun concurrency limits. If scheduler_dag_bag.get_dag_for_run() returns None (serialized DAG transiently absent), the code executed a bulk UPDATE task_instance SET state='failed' for all SCHEDULED tasks of that DAG — instead of treating it as a transient miss.

Fix

Replace the bulk-fail with a graceful skip: when the serialized DAG is not found, log a WARNING (down from ERROR), add the dag_id to starved_dags so the rest of its tasks are also skipped in this iteration, and let the scheduler retry naturally on the next heartbeat.

Changes:

  • airflow-core/src/airflow/jobs/scheduler_job_runner.py: removed the bulk UPDATE … SET state=FAILED block; replaced with starved_dags.add(dag_id) + warning log.
  • airflow-core/tests/unit/jobs/test_scheduler_job.py:
    • Renamed test_queued_task_instances_fails_with_missing_dag → test_queued_task_instances_skips_with_missing_dag and updated assertions to expect SCHEDULED state (not FAILED).
    • Added new regression test test_missing_serialized_dag_does_not_bulk_fail_tasks (based on the reproducer from the issue) that explicitly asserts tasks remain SCHEDULED when the serialized DAG is transiently missing.

Both new/updated tests pass. The full test_executable_task_instances suite (24 tests) passes with no regressions.

Closes: #62050



Was generative AI tooling used to co-author this PR?
  • Yes (Claude Code)

Generated-by: Claude Code following the guidelines


  • Read the Pull Request Guidelines for more information. Note: commit author/co-author name and email in commits become permanently public when merged.
  • For fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
  • When adding dependency, check compliance with the ASF 3rd Party License Policy.
  • For significant user-facing changes create newsfragment: {pr_number}.significant.rst or {issue_number}.significant.rst, in airflow-core/newsfragments.

@YoannAbriel
YoannAbriel requested review from XD-DENG and ashb as code owners March 4, 2026 15:25
@YoannAbriel
YoannAbriel force-pushed the fix/issue-62050 branch 4 times, most recently from c0bc313 to c38181d Compare March 8, 2026 19:04

@potiuk potiuk left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM. But I would like others who are better versed in scheduler to take a look @ashb @kaxil @ephraimbuddy

.values(state=TaskInstanceState.FAILED)
.execution_options(synchronize_session="fetch")
)
starved_dags.add(dag_id)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

If a DAG is permanently deleted (not just transiently missing), tasks will stay SCHEDULED forever and this warning will fire every scheduler iteration. Would it be worth tracking consecutive misses per dag_id and escalating to failure after N iterations? The issue thread mentioned this approach too. At minimum, this should probably be documented as a known limitation in the PR description so follow-up work can address it.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Added a code comment documenting this as a known limitation. Tracking consecutive misses can be a follow-up.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

By reading the code I had exactly the same comment like kaxil.

I think we can not leave it like this, because it would maybe improve some cases but leave endless loop of scheduled tasks. Therefore I'd suggest to add a check such that tasks are not endless looped over.

if not serialized_dag:
self.log.error(
"DAG '%s' for task instance %s not found in serialized_dag table",
self.log.warning(

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

When a batch has multiple TIs for the same DAG, each one will hit get_dag_for_run and log this warning individually before the starved_dags filter kicks in on the next query iteration. Consider checking if dag_id in starved_dags: continue before entering the has_task_concurrency_limits block, both to avoid redundant DB lookups and to reduce log noise.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done — added early starved_dags check before the has_task_concurrency_limits block.

assert all(ti.state == State.FAILED for ti in tis)
assert all(ti.state == State.SCHEDULED for ti in tis)

def test_missing_serialized_dag_does_not_bulk_fail_tasks(self, dag_maker, session):

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This test covers the same scenario as test_queued_task_instances_skips_with_missing_dag above (mock get_dag_for_run to return None, assert tasks stay SCHEDULED). The only differences are the dag_id string, max_active_tis_per_dag value (1 vs 2), and the assertion style. I'd pick one or the other rather than keeping both.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Removed the duplicate test — the existing one covers this.

self.job_runner = SchedulerJobRunner(job=scheduler_job)

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

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

mock.MagicMock() without spec won't catch typos if someone later renames get_dag_for_run. Consider mock.MagicMock(spec=DBDagBag) (same applies to the existing test on line 1813, but that's pre-existing).

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Switched to spec=DBDagBag.

@eladkal

eladkal commented Mar 15, 2026

Copy link
Copy Markdown
Contributor

static checks are failing

@eladkal eladkal added this to the Airflow 3.1.9 milestone Mar 15, 2026
@YoannAbriel
YoannAbriel force-pushed the fix/issue-62050 branch 6 times, most recently from 4b2d28a to 4f3e5ae Compare March 16, 2026 16:09
@YoannAbriel
YoannAbriel force-pushed the fix/issue-62050 branch 3 times, most recently from b1580dc to d7c6064 Compare March 24, 2026 22:03
@YoannAbriel
YoannAbriel force-pushed the fix/issue-62050 branch 2 times, most recently from 5508b92 to f604fd5 Compare April 5, 2026 12:03
@eladkal eladkal removed this from the Airflow 3.1.9 milestone Apr 6, 2026
if not serialized_dag:
self.log.error(
"DAG '%s' for task instance %s not found in serialized_dag table",
self.log.warning(

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Replicating the comment here such that it is not over-looked:

By reading the code I had exactly the same comment like kaxil earlier.

I think we can not leave it like this, because it would maybe improve some cases but leave endless loop of scheduled tasks. Therefore I'd suggest to add a check such that tasks are not endless looped over.

@potiuk
potiuk marked this pull request as draft April 6, 2026 22:40
@potiuk

potiuk commented Apr 6, 2026 •

Copy link
Copy Markdown
Member

@YoannAbriel Converting to draft — this PR doesn't yet meet our Pull Request quality criteria.

  • ❌ Other failing CI checks: Failing: Postgres tests: core / DB-core:Postgres:14:3.10:Core...Serialization, MySQL tests: core / DB-core:MySQL:8.0:3.10:Core...Serialization, Sqlite tests: core / DB-core:Sqlite:3.10:Core...Serialization, Low dep tests:core / All-core:LowestDeps:14:3.10:Core...Serialization. Run prek run --from-ref main locally to reproduce. See static checks docs.
  • ⚠️ Unresolved review comments: This PR has 5 unresolved review threads from maintainers: @kaxil (MEMBER): 4 unresolved threads; @jscheffl (MEMBER): 1 unresolved threads. Please review and resolve all inline review comments before requesting another review. You can resolve a conversation by clicking 'Resolve conversation' on each thread after addressing the feedback. See pull request guidelines.

See the linked criteria for how to fix each item, then mark the PR "Ready for review". This is not a rejection — just an invitation to bring the PR up to standard. No rush.


Note: This comment was drafted by an AI-assisted triage tool and may contain mistakes. Once you have addressed the points above, an Apache Airflow maintainer — a real person — will take the next look at your PR. We use this two-stage triage process so that our maintainers' limited time is spent where it matters most: the conversation with you.

@YoannAbriel
YoannAbriel force-pushed the fix/issue-62050 branch 11 times, most recently from 76b522e to 4e7aadd Compare April 10, 2026 09:04
@kaxil
kaxil requested a review from Copilot April 10, 2026 19:55

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

This PR changes the scheduler’s handling of transiently missing serialized DAGs during concurrency checks to avoid bulk-failing all SCHEDULED task instances for the DAG, which can happen intermittently in multi-scheduler deployments.

Changes:

  • Scheduler: replace the bulk UPDATE ... SET state=FAILED behavior with a warning + “skip this DAG for this iteration” behavior.
  • Tests: update the existing missing-serialized-DAG test to expect tasks remain SCHEDULED, and add a regression test to ensure tasks are not bulk-failed.

Reviewed changes

Copilot reviewed 2 out of 2 changed files in this pull request and generated 3 comments.

File Description
airflow-core/src/airflow/jobs/scheduler_job_runner.py Stop bulk-failing scheduled tasks when get_dag_for_run() returns None; instead warn and mark the DAG as starved for the current iteration.
airflow-core/tests/unit/jobs/test_scheduler_job.py Adjust expectations for missing serialized DAG behavior and add regression coverage for issue #62050.
Comments suppressed due to low confidence (1)

airflow-core/tests/unit/jobs/test_scheduler_job.py:1822

  • The dag_id string looks like a copy/paste from another test (...not_in_dagbag) and no longer matches this test’s intent/name (missing serialized DAG). Consider updating the dag_id to something aligned with test_queued_task_instances_skips_with_missing_dag to make failures easier to interpret.
        dag_id = "SchedulerJobTest.test_find_executable_task_instances_not_in_dagbag"
        task_id_1 = "dummy"
        task_id_2 = "dummydummy"

        with dag_maker(dag_id=dag_id, session=session, default_args={"max_active_tis_per_dag": 1}):

self.job_runner = SchedulerJobRunner(job=scheduler_job)

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

Copilot AI Apr 10, 2026

Copy link

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This uses an unspec’d mock.MagicMock() for scheduler_dag_bag. Consider switching to a spec’d/autospec’d mock so the test fails if the production code calls an unexpected attribute/method on the dag bag.

Suggested change
self.job_runner.scheduler_dag_bag = mock.MagicMock()
self.job_runner.scheduler_dag_bag = mock.create_autospec(DagBag, instance=True)

Copilot uses AI. Check for mistakes.
Previously, the scheduler would bulk-fail all SCHEDULED tasks when it couldn't find the
serialized DAG in the DagBag.
"""
dag_id = "SchedulerJobTest.test_missing_serialized_dag_bulk_fails"

Copilot AI Apr 10, 2026

Copy link

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The dag_id value still includes ...bulk_fails, but this test now asserts the opposite behavior (does not bulk-fail). Renaming the dag_id to match the new behavior would improve clarity when debugging failures.

Suggested change
dag_id = "SchedulerJobTest.test_missing_serialized_dag_bulk_fails"
dag_id = "SchedulerJobTest.test_missing_serialized_dag_does_not_bulk_fail_tasks"

Copilot uses AI. Check for mistakes.
Comment on lines 1838 to 1840
session.flush()
res = self.job_runner._executable_task_instances_to_queued(max_tis=32, session=session)
session.flush()

Copilot AI Apr 10, 2026

Copy link

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

In this test, scheduler_dag_bag is set to mock.MagicMock() (above) without a spec/autospec, which can hide attribute typos and makes the test less strict. Consider using a spec’d mock (e.g., mock.create_autospec(..., instance=True) or mock.Mock(spec=["get_dag_for_run"])) so the test only allows the expected API.

Copilot uses AI. Check for mistakes.
@YoannAbriel
YoannAbriel force-pushed the fix/issue-62050 branch 3 times, most recently from e69d496 to bb88ce7 Compare April 13, 2026 06:06
…ransiently missing

When the scheduler cannot find a DAG in the serialized_dag table while checking
task concurrency limits, it previously set all SCHEDULED task instances for that
DAG to FAILED via a bulk UPDATE. This caused intermittent mass task failures on
multi-scheduler setups (e.g. MWAA) where the serialized DAG may be transiently
absent during a DAG file parse cycle.

Instead, treat the missing serialized DAG as a transient condition: log a warning,
add the dag_id to starved_dags so subsequent tasks for the same DAG are also
skipped, and let the scheduler retry on the next iteration.

The existing test 'test_queued_task_instances_fails_with_missing_dag' has been
updated to reflect the new expected behavior (tasks remain SCHEDULED). A new
regression test 'test_missing_serialized_dag_does_not_bulk_fail_tasks' explicitly
verifies no bulk-failure occurs.

Closes apache#62050
@potiuk

potiuk commented May 6, 2026

Copy link
Copy Markdown
Member

@YoannAbriel This draft PR has been inactive for 29 days since the last triage comment and no response from the author. Closing to keep the queue clean.

You are welcome to reopen this PR when you resume work, or to open a new one addressing the issues previously raised. There is no rush — take your time.


Note: This comment was drafted by an AI-assisted triage tool and may contain mistakes. Once you have addressed the points above, an Apache Airflow maintainer — a real person — will take the next look at your PR. We use this two-stage triage process so that our maintainers' limited time is spent where it matters most: the conversation with you.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:scheduler type:bug-fix Changelog: Bug Fixes

Projects

None yet

Development

Successfully merging this pull request may close these issues.

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

8 participants