Skip to content

Fix duplicate task execution when running multiple schedulers (HA) - #60330

Merged
kaxil merged 1 commit into
apache:mainfrom
astronomer:fix-ha-scheduler-race
Mar 19, 2026
Merged

Fix duplicate task execution when running multiple schedulers (HA)#60330
kaxil merged 1 commit into
apache:mainfrom
astronomer:fix-ha-scheduler-race

Conversation

@ephraimbuddy

Copy link
Copy Markdown
Contributor

In HA, two scheduler processes can race to schedule the same TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id alone, so a scheduler could increment try_number and transition state even after another scheduler had already advanced the TI (e.g. to SCHEDULED/QUEUED), resulting in duplicate attempts being queued.

This change makes scheduling idempotent under HA races by:

  • Guarding schedule_tis() DB updates to only apply when the TI is still in schedulable states (derived from SCHEDULEABLE_STATES, handling NULL explicitly).

  • Using a single CASE (next_try_number) so reschedules (UP_FOR_RESCHEDULE) do not start a new try, and applying this consistently to both normal scheduling and the EmptyOperator fast-path.

Adds regression tests covering:

  • TI already queued by another scheduler.
  • EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
  • UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
  • Only one “scheduler” update succeeds when competing.

Closes: #57618

Note: The reproduction of this issue was based on unit tests

Comment threadairflow-core/src/airflow/models/dagrun.py
ashb
ashb approved these changes Jan 9, 2026

@ashbashb 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 once we've done some perf testing

@dabla

Copy link
Copy Markdown
Contributor

@ephraimbuddy Good job, we are actually testing this fix right now, keep you guys posted if it helped. The weird thing is it seems (we're not sure yet) that we only have this issue with DAG's that have tasks that use the WinRMOperator.

@dabla

Copy link
Copy Markdown
Contributor

@ephraimbuddy Unfortunately, even after applying this fix, we still run onto the same issue.

@ephraimbuddy

Copy link
Copy Markdown
ContributorAuthor

@ephraimbuddy Unfortunately, even after applying this fix, we still run onto the same issue.

Major problem with the issue is reliability of reproduction path. I have reproduced it once but trying again, I can't. Do you have steps for reliable reproduction?

@dabladabla mentioned this pull request Jan 16, 2026
1 task
@dabla

Copy link
Copy Markdown
Contributor

Now we also experience the same issue with HttpOperator:

[2026-01-21 04:44:19] INFO - Calling HTTP method
[2026-01-21 04:44:19] INFO - The hook_class 'airflow.providers.http.hooks.http.HttpHook' is not fully initialized (UI widgets will be missing), because the 'flask_appbuilder' package is not installed, however it is not required for Airflow components to work
[2026-01-21 04:44:25] ERROR - Server indicated the task shouldn't be running anymore. Terminating process detail={"detail":{"reason":"not_found","message":"Task Instance not found"}}
[2026-01-21 04:44:30] ERROR - Task killed!

@yennysu

Copy link
Copy Markdown

@ephraimbuddy Unfortunately, even after applying this fix, we still run onto the same issue.

Major problem with the issue is reliability of reproduction path. I have reproduced it once but trying again, I can't. Do you have steps for reliable reproduction?

While we haven’t been able to reproduce this reliably, we’ve observed that the frequency of occurrences seems to scale with the number of schedulers running.

@dstandish

Copy link
Copy Markdown
Contributor

@dabla et al, does anyone have a concrete theory about exactly how this happpens?

@sam-dumont

Copy link
Copy Markdown
Contributor

This PR fixed most of our issues (and we subsequently discovered another race condition). Our investigation is here in the comments #59378 (comment)

Anything we can do to finalize it and see it in airflow 3.1.9?

@kaxil
kaxilforce-pushed the fix-ha-scheduler-race branch 2 times, most recently from d0d020f to d494a67CompareMarch 19, 2026 03:09
@kaxilkaxil changed the title Fix HA scheduler try_number double incrementFix duplicate task execution when running multiple schedulers (HA)Mar 19, 2026
@kaxil

Copy link
Copy Markdown
Member

@dabla@yennysu The error you're seeing ("reason":"not_found","message":"Task Instance not found") is a separate bug from what this PR fixes. This PR fixes the try_number double-increment race (where two schedulers both schedule the same TI). Your error is a 404, not a 409, caused by UUID reassignment during orphan adoption: when a scheduler crashes, another scheduler resets the orphaned TI via prepare_db_for_next_try() which assigns a new UUID. The worker that's still running heartbeats with the old UUID and gets a 404.

We're tracking that as a separate issue. Scaling with scheduler count matches this theory since more schedulers means more adoption cycles.

@kaxil
kaxilforce-pushed the fix-ha-scheduler-race branch from 536195e to 5aed519CompareMarch 19, 2026 04:28
In HA, two scheduler processes can race to schedule the same
TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id
alone, so a scheduler could increment try_number and transition
state even after another scheduler had already advanced the TI (e.g. to
SCHEDULED/QUEUED), resulting in duplicate attempts being queued.
This change makes scheduling idempotent under HA races by:
- Guarding schedule_tis() DB updates to only apply when the TI is still
in schedulable states (derived from SCHEDULEABLE_STATES, handling
NULL explicitly).
- Using a single CASE (next_try_number) so reschedules
(UP_FOR_RESCHEDULE) do not start a new try, and applying this
consistently to both normal scheduling and the EmptyOperator fast-path.
The CASE uses TI.id (not TI.state) to avoid MySQL SET left-to-right
evaluation issues.
Adds regression tests covering:
- TI already queued by another scheduler.
- EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
- UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
- Only one "scheduler" update succeeds when competing.
Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
@kaxil
kaxilforce-pushed the fix-ha-scheduler-race branch from 5aed519 to 66ffb5aCompareMarch 19, 2026 04:29
@kaxil
kaxil merged commit a31db1f into apache:mainMar 19, 2026
9 checks passed
@kaxil
kaxil deleted the fix-ha-scheduler-race branch March 19, 2026 04:29
@github-actions

Copy link
Copy Markdown
Contributor

Backport failed to create: v3-1-test. View the failure log Run details

Note: As of Merging PRs targeted for Airflow 3.X
the committer who merges the PR is responsible for backporting the PRs that are bug fixes (generally speaking) to the maintenance branches.

In matter of doubt please ask in #release-management Slack channel.

StatusBranchResult
v3-1-testCommit Link

You can attempt to backport this manually by running:

cherry_picker a31db1f v3-1-test

This should apply the commit to the v3-1-test branch and leave the commit in conflict state marking
the files that need manual conflict resolution.

After you have resolved the conflicts, you can continue the backport process by running:

cherry_picker --continue

If you don't have cherry-picker installed, see the installation guide.

@dabla

Copy link
Copy Markdown
Contributor

@dabla et al, does anyone have a concrete theory about exactly how this happpens?

Sorry, missed your message @dstandish, we where experiencing this on kunernetes with 2 pod instances of the scheduler. Since we installed 3.1.8 I haven’t experienced this error anymore, so I might assume it is indeed fixed.

i’ve tried this fix through monkey patching though on earlier versions and then the issue persisted, so weird.

fat-catTW pushed a commit to fat-catTW/airflow that referenced this pull request Mar 22, 2026
…#60330)
In HA, two scheduler processes can race to schedule the same
TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id
alone, so a scheduler could increment try_number and transition
state even after another scheduler had already advanced the TI (e.g. to
SCHEDULED/QUEUED), resulting in duplicate attempts being queued.
This change makes scheduling idempotent under HA races by:
- Guarding schedule_tis() DB updates to only apply when the TI is still
in schedulable states (derived from SCHEDULEABLE_STATES, handling
NULL explicitly).
- Using a single CASE (next_try_number) so reschedules
(UP_FOR_RESCHEDULE) do not start a new try, and applying this
consistently to both normal scheduling and the EmptyOperator fast-path.
The CASE uses TI.id (not TI.state) to avoid MySQL SET left-to-right
evaluation issues.
Adds regression tests covering:
- TI already queued by another scheduler.
- EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
- UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
- Only one "scheduler" update succeeds when competing.
Closes: apache#57618
Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
Suraj-kumar00 pushed a commit to Suraj-kumar00/airflow that referenced this pull request Apr 7, 2026
…#60330)
In HA, two scheduler processes can race to schedule the same
TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id
alone, so a scheduler could increment try_number and transition
state even after another scheduler had already advanced the TI (e.g. to
SCHEDULED/QUEUED), resulting in duplicate attempts being queued.
This change makes scheduling idempotent under HA races by:
- Guarding schedule_tis() DB updates to only apply when the TI is still
in schedulable states (derived from SCHEDULEABLE_STATES, handling
NULL explicitly).
- Using a single CASE (next_try_number) so reschedules
(UP_FOR_RESCHEDULE) do not start a new try, and applying this
consistently to both normal scheduling and the EmptyOperator fast-path.
The CASE uses TI.id (not TI.state) to avoid MySQL SET left-to-right
evaluation issues.
Adds regression tests covering:
- TI already queued by another scheduler.
- EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
- UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
- Only one "scheduler" update succeeds when competing.
Closes: apache#57618
Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
abhijeets25012-tech pushed a commit to abhijeets25012-tech/airflow that referenced this pull request Apr 9, 2026
…#60330)
In HA, two scheduler processes can race to schedule the same
TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id
alone, so a scheduler could increment try_number and transition
state even after another scheduler had already advanced the TI (e.g. to
SCHEDULED/QUEUED), resulting in duplicate attempts being queued.
This change makes scheduling idempotent under HA races by:
- Guarding schedule_tis() DB updates to only apply when the TI is still
in schedulable states (derived from SCHEDULEABLE_STATES, handling
NULL explicitly).
- Using a single CASE (next_try_number) so reschedules
(UP_FOR_RESCHEDULE) do not start a new try, and applying this
consistently to both normal scheduling and the EmptyOperator fast-path.
The CASE uses TI.id (not TI.state) to avoid MySQL SET left-to-right
evaluation issues.
Adds regression tests covering:
- TI already queued by another scheduler.
- EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
- UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
- Only one "scheduler" update succeeds when competing.
Closes: apache#57618
Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
@pdellarciprete

Copy link
Copy Markdown

Sorry, I’m not sure I understood — is this fix included in 3.1.8 or 3.1.9?
If it’s planned for 3.1.9, is there an estimated release timeline?

cc: @ephraimbuddy@dabla@dstandish

@anormalpersonBE

Copy link
Copy Markdown

@pdellarciprete I was confused as well about the state of this fix.
As far as I can see on the main branch: your fix is in there.

@fabbuc-gyg

Copy link
Copy Markdown

@kaxil I just read your comment about the 404 error and "Task Instance not Found". I have been getting this a lot on our Airflow 3.1.8 deployment. Are you aware of any PR on this matter, or any workaround to get a grip on this?

@pdellarciprete

pdellarciprete commented May 5, 2026

Copy link
Copy Markdown

Unfortunately, in our case this fix mitigated but didn't solve the problem.

@squ1b3r

Copy link
Copy Markdown

Currently running version 3.2.1 and using PartitionedAssetTimetable(assets=Asset("my-asset")) still causes duplicated dags runs when running multiple schedulers.

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

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Schedulers race condition when using with kubernetes executor

13 participants

@ephraimbuddy@dabla@yennysu@dstandish@sam-dumont@kaxil@pdellarciprete@anormalpersonBE@fabbuc-gyg@squ1b3r@ashb@vatsrahul1001@eladkal
, 'i'); if (__m === '*' || __re.test(location.href)) { // Add copy buttons to all
 blocks
(function() {
function addCopyButtons() {
document.querySelectorAll('pre code').forEach(function(codeBlock) {
if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;
codeBlock.parentElement.setAttribute('data-copy-added', 'true');
var btn = document.createElement('button');
btn.textContent = 'Copy';
btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';
btn.onmouseover = function() { this.style.opacity = '1'; };
btn.onmouseout = function() { this.style.opacity = '0.7'; };
btn.onclick = function() {
navigator.clipboard.writeText(codeBlock.textContent).then(function() {
btn.textContent = 'Copied!';
setTimeout(function() { btn.textContent = 'Copy'; }, 1500);
});
};
codeBlock.parentElement.style.position = 'relative';
codeBlock.parentElement.appendChild(btn);
});
}
addCopyButtons();
// Re-run on dynamic content
var observer = new MutationObserver(addCopyButtons);
observer.observe(document.body, { childList: true, subtree: true });
})();
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Fix duplicate task execution when running multiple schedulers (HA) by ephraimbuddy · Pull Request #60330 · apache/airflow · GitHub
Skip to content

Fix duplicate task execution when running multiple schedulers (HA) - #60330

Merged
kaxil merged 1 commit into
apache:mainfrom
astronomer:fix-ha-scheduler-race
Mar 19, 2026
Merged

Fix duplicate task execution when running multiple schedulers (HA)#60330
kaxil merged 1 commit into
apache:mainfrom
astronomer:fix-ha-scheduler-race

Conversation

@ephraimbuddy

Copy link
Copy Markdown
Contributor

In HA, two scheduler processes can race to schedule the same TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id alone, so a scheduler could increment try_number and transition state even after another scheduler had already advanced the TI (e.g. to SCHEDULED/QUEUED), resulting in duplicate attempts being queued.

This change makes scheduling idempotent under HA races by:

  • Guarding schedule_tis() DB updates to only apply when the TI is still in schedulable states (derived from SCHEDULEABLE_STATES, handling NULL explicitly).

  • Using a single CASE (next_try_number) so reschedules (UP_FOR_RESCHEDULE) do not start a new try, and applying this consistently to both normal scheduling and the EmptyOperator fast-path.

Adds regression tests covering:

  • TI already queued by another scheduler.
  • EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
  • UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
  • Only one “scheduler” update succeeds when competing.

Closes: #57618

Note: The reproduction of this issue was based on unit tests

Comment threadairflow-core/src/airflow/models/dagrun.py
ashb
ashb approved these changes Jan 9, 2026

@ashbashb 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 once we've done some perf testing

@dabla

Copy link
Copy Markdown
Contributor

@ephraimbuddy Good job, we are actually testing this fix right now, keep you guys posted if it helped. The weird thing is it seems (we're not sure yet) that we only have this issue with DAG's that have tasks that use the WinRMOperator.

@dabla

Copy link
Copy Markdown
Contributor

@ephraimbuddy Unfortunately, even after applying this fix, we still run onto the same issue.

@ephraimbuddy

Copy link
Copy Markdown
ContributorAuthor

@ephraimbuddy Unfortunately, even after applying this fix, we still run onto the same issue.

Major problem with the issue is reliability of reproduction path. I have reproduced it once but trying again, I can't. Do you have steps for reliable reproduction?

@dabladabla mentioned this pull request Jan 16, 2026
1 task
@dabla

Copy link
Copy Markdown
Contributor

Now we also experience the same issue with HttpOperator:

[2026-01-21 04:44:19] INFO - Calling HTTP method
[2026-01-21 04:44:19] INFO - The hook_class 'airflow.providers.http.hooks.http.HttpHook' is not fully initialized (UI widgets will be missing), because the 'flask_appbuilder' package is not installed, however it is not required for Airflow components to work
[2026-01-21 04:44:25] ERROR - Server indicated the task shouldn't be running anymore. Terminating process detail={"detail":{"reason":"not_found","message":"Task Instance not found"}}
[2026-01-21 04:44:30] ERROR - Task killed!

@yennysu

Copy link
Copy Markdown

@ephraimbuddy Unfortunately, even after applying this fix, we still run onto the same issue.

Major problem with the issue is reliability of reproduction path. I have reproduced it once but trying again, I can't. Do you have steps for reliable reproduction?

While we haven’t been able to reproduce this reliably, we’ve observed that the frequency of occurrences seems to scale with the number of schedulers running.

@dstandish

Copy link
Copy Markdown
Contributor

@dabla et al, does anyone have a concrete theory about exactly how this happpens?

@sam-dumont

Copy link
Copy Markdown
Contributor

This PR fixed most of our issues (and we subsequently discovered another race condition). Our investigation is here in the comments #59378 (comment)

Anything we can do to finalize it and see it in airflow 3.1.9?

@kaxil
kaxilforce-pushed the fix-ha-scheduler-race branch 2 times, most recently from d0d020f to d494a67CompareMarch 19, 2026 03:09
@kaxilkaxil changed the title Fix HA scheduler try_number double incrementFix duplicate task execution when running multiple schedulers (HA)Mar 19, 2026
@kaxil

Copy link
Copy Markdown
Member

@dabla@yennysu The error you're seeing ("reason":"not_found","message":"Task Instance not found") is a separate bug from what this PR fixes. This PR fixes the try_number double-increment race (where two schedulers both schedule the same TI). Your error is a 404, not a 409, caused by UUID reassignment during orphan adoption: when a scheduler crashes, another scheduler resets the orphaned TI via prepare_db_for_next_try() which assigns a new UUID. The worker that's still running heartbeats with the old UUID and gets a 404.

We're tracking that as a separate issue. Scaling with scheduler count matches this theory since more schedulers means more adoption cycles.

@kaxil
kaxilforce-pushed the fix-ha-scheduler-race branch from 536195e to 5aed519CompareMarch 19, 2026 04:28
In HA, two scheduler processes can race to schedule the same
TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id
alone, so a scheduler could increment try_number and transition
state even after another scheduler had already advanced the TI (e.g. to
SCHEDULED/QUEUED), resulting in duplicate attempts being queued.
This change makes scheduling idempotent under HA races by:
- Guarding schedule_tis() DB updates to only apply when the TI is still
in schedulable states (derived from SCHEDULEABLE_STATES, handling
NULL explicitly).
- Using a single CASE (next_try_number) so reschedules
(UP_FOR_RESCHEDULE) do not start a new try, and applying this
consistently to both normal scheduling and the EmptyOperator fast-path.
The CASE uses TI.id (not TI.state) to avoid MySQL SET left-to-right
evaluation issues.
Adds regression tests covering:
- TI already queued by another scheduler.
- EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
- UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
- Only one "scheduler" update succeeds when competing.
Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
@kaxil
kaxilforce-pushed the fix-ha-scheduler-race branch from 5aed519 to 66ffb5aCompareMarch 19, 2026 04:29
@kaxil
kaxil merged commit a31db1f into apache:mainMar 19, 2026
9 checks passed
@kaxil
kaxil deleted the fix-ha-scheduler-race branch March 19, 2026 04:29
@github-actions

Copy link
Copy Markdown
Contributor

Backport failed to create: v3-1-test. View the failure log Run details

Note: As of Merging PRs targeted for Airflow 3.X
the committer who merges the PR is responsible for backporting the PRs that are bug fixes (generally speaking) to the maintenance branches.

In matter of doubt please ask in #release-management Slack channel.

StatusBranchResult
v3-1-testCommit Link

You can attempt to backport this manually by running:

cherry_picker a31db1f v3-1-test

This should apply the commit to the v3-1-test branch and leave the commit in conflict state marking
the files that need manual conflict resolution.

After you have resolved the conflicts, you can continue the backport process by running:

cherry_picker --continue

If you don't have cherry-picker installed, see the installation guide.

@dabla

Copy link
Copy Markdown
Contributor

@dabla et al, does anyone have a concrete theory about exactly how this happpens?

Sorry, missed your message @dstandish, we where experiencing this on kunernetes with 2 pod instances of the scheduler. Since we installed 3.1.8 I haven’t experienced this error anymore, so I might assume it is indeed fixed.

i’ve tried this fix through monkey patching though on earlier versions and then the issue persisted, so weird.

fat-catTW pushed a commit to fat-catTW/airflow that referenced this pull request Mar 22, 2026
…#60330)
In HA, two scheduler processes can race to schedule the same
TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id
alone, so a scheduler could increment try_number and transition
state even after another scheduler had already advanced the TI (e.g. to
SCHEDULED/QUEUED), resulting in duplicate attempts being queued.
This change makes scheduling idempotent under HA races by:
- Guarding schedule_tis() DB updates to only apply when the TI is still
in schedulable states (derived from SCHEDULEABLE_STATES, handling
NULL explicitly).
- Using a single CASE (next_try_number) so reschedules
(UP_FOR_RESCHEDULE) do not start a new try, and applying this
consistently to both normal scheduling and the EmptyOperator fast-path.
The CASE uses TI.id (not TI.state) to avoid MySQL SET left-to-right
evaluation issues.
Adds regression tests covering:
- TI already queued by another scheduler.
- EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
- UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
- Only one "scheduler" update succeeds when competing.
Closes: apache#57618
Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
Suraj-kumar00 pushed a commit to Suraj-kumar00/airflow that referenced this pull request Apr 7, 2026
…#60330)
In HA, two scheduler processes can race to schedule the same
TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id
alone, so a scheduler could increment try_number and transition
state even after another scheduler had already advanced the TI (e.g. to
SCHEDULED/QUEUED), resulting in duplicate attempts being queued.
This change makes scheduling idempotent under HA races by:
- Guarding schedule_tis() DB updates to only apply when the TI is still
in schedulable states (derived from SCHEDULEABLE_STATES, handling
NULL explicitly).
- Using a single CASE (next_try_number) so reschedules
(UP_FOR_RESCHEDULE) do not start a new try, and applying this
consistently to both normal scheduling and the EmptyOperator fast-path.
The CASE uses TI.id (not TI.state) to avoid MySQL SET left-to-right
evaluation issues.
Adds regression tests covering:
- TI already queued by another scheduler.
- EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
- UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
- Only one "scheduler" update succeeds when competing.
Closes: apache#57618
Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
abhijeets25012-tech pushed a commit to abhijeets25012-tech/airflow that referenced this pull request Apr 9, 2026
…#60330)
In HA, two scheduler processes can race to schedule the same
TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id
alone, so a scheduler could increment try_number and transition
state even after another scheduler had already advanced the TI (e.g. to
SCHEDULED/QUEUED), resulting in duplicate attempts being queued.
This change makes scheduling idempotent under HA races by:
- Guarding schedule_tis() DB updates to only apply when the TI is still
in schedulable states (derived from SCHEDULEABLE_STATES, handling
NULL explicitly).
- Using a single CASE (next_try_number) so reschedules
(UP_FOR_RESCHEDULE) do not start a new try, and applying this
consistently to both normal scheduling and the EmptyOperator fast-path.
The CASE uses TI.id (not TI.state) to avoid MySQL SET left-to-right
evaluation issues.
Adds regression tests covering:
- TI already queued by another scheduler.
- EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
- UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
- Only one "scheduler" update succeeds when competing.
Closes: apache#57618
Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
@pdellarciprete

Copy link
Copy Markdown

Sorry, I’m not sure I understood — is this fix included in 3.1.8 or 3.1.9?
If it’s planned for 3.1.9, is there an estimated release timeline?

cc: @ephraimbuddy@dabla@dstandish

@anormalpersonBE

Copy link
Copy Markdown

@pdellarciprete I was confused as well about the state of this fix.
As far as I can see on the main branch: your fix is in there.

@fabbuc-gyg

Copy link
Copy Markdown

@kaxil I just read your comment about the 404 error and "Task Instance not Found". I have been getting this a lot on our Airflow 3.1.8 deployment. Are you aware of any PR on this matter, or any workaround to get a grip on this?

@pdellarciprete

pdellarciprete commented May 5, 2026

Copy link
Copy Markdown

Unfortunately, in our case this fix mitigated but didn't solve the problem.

@squ1b3r

Copy link
Copy Markdown

Currently running version 3.2.1 and using PartitionedAssetTimetable(assets=Asset("my-asset")) still causes duplicated dags runs when running multiple schedulers.

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

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Schedulers race condition when using with kubernetes executor

13 participants

@ephraimbuddy@dabla@yennysu@dstandish@sam-dumont@kaxil@pdellarciprete@anormalpersonBE@fabbuc-gyg@squ1b3r@ashb@vatsrahul1001@eladkal
, 'i'); if (__m === '*' || __re.test(location.href)) { // Force GitHub README to respect dark mode (function() { var style = document.createElement('style'); style.textContent = ' .markdown-body { color-scheme: dark light; } .markdown-body pre { background: #161b22 !important; } .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; } .markdown-body table th, .markdown-body table td { border-color: #30363d !important; } .markdown-body img { background: #0d1117; } .markdown-body blockquote { border-left-color: #8b949e; } .markdown-body hr { border-color: #30363d; } '; document.head.appendChild(style); })(); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' Fix duplicate task execution when running multiple schedulers (HA) by ephraimbuddy · Pull Request #60330 · apache/airflow · GitHub
Skip to content

Fix duplicate task execution when running multiple schedulers (HA) - #60330

Merged
kaxil merged 1 commit into
apache:mainfrom
astronomer:fix-ha-scheduler-race
Mar 19, 2026
Merged

Fix duplicate task execution when running multiple schedulers (HA)#60330
kaxil merged 1 commit into
apache:mainfrom
astronomer:fix-ha-scheduler-race

Conversation

@ephraimbuddy

Copy link
Copy Markdown
Contributor

In HA, two scheduler processes can race to schedule the same TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id alone, so a scheduler could increment try_number and transition state even after another scheduler had already advanced the TI (e.g. to SCHEDULED/QUEUED), resulting in duplicate attempts being queued.

This change makes scheduling idempotent under HA races by:

  • Guarding schedule_tis() DB updates to only apply when the TI is still in schedulable states (derived from SCHEDULEABLE_STATES, handling NULL explicitly).

  • Using a single CASE (next_try_number) so reschedules (UP_FOR_RESCHEDULE) do not start a new try, and applying this consistently to both normal scheduling and the EmptyOperator fast-path.

Adds regression tests covering:

  • TI already queued by another scheduler.
  • EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
  • UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
  • Only one “scheduler” update succeeds when competing.

Closes: #57618

Note: The reproduction of this issue was based on unit tests

Comment threadairflow-core/src/airflow/models/dagrun.py
ashb
ashb approved these changes Jan 9, 2026

@ashbashb 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 once we've done some perf testing

@dabla

Copy link
Copy Markdown
Contributor

@ephraimbuddy Good job, we are actually testing this fix right now, keep you guys posted if it helped. The weird thing is it seems (we're not sure yet) that we only have this issue with DAG's that have tasks that use the WinRMOperator.

@dabla

Copy link
Copy Markdown
Contributor

@ephraimbuddy Unfortunately, even after applying this fix, we still run onto the same issue.

@ephraimbuddy

Copy link
Copy Markdown
ContributorAuthor

@ephraimbuddy Unfortunately, even after applying this fix, we still run onto the same issue.

Major problem with the issue is reliability of reproduction path. I have reproduced it once but trying again, I can't. Do you have steps for reliable reproduction?

@dabladabla mentioned this pull request Jan 16, 2026
1 task
@dabla

Copy link
Copy Markdown
Contributor

Now we also experience the same issue with HttpOperator:

[2026-01-21 04:44:19] INFO - Calling HTTP method
[2026-01-21 04:44:19] INFO - The hook_class 'airflow.providers.http.hooks.http.HttpHook' is not fully initialized (UI widgets will be missing), because the 'flask_appbuilder' package is not installed, however it is not required for Airflow components to work
[2026-01-21 04:44:25] ERROR - Server indicated the task shouldn't be running anymore. Terminating process detail={"detail":{"reason":"not_found","message":"Task Instance not found"}}
[2026-01-21 04:44:30] ERROR - Task killed!

@yennysu

Copy link
Copy Markdown

@ephraimbuddy Unfortunately, even after applying this fix, we still run onto the same issue.

Major problem with the issue is reliability of reproduction path. I have reproduced it once but trying again, I can't. Do you have steps for reliable reproduction?

While we haven’t been able to reproduce this reliably, we’ve observed that the frequency of occurrences seems to scale with the number of schedulers running.

@dstandish

Copy link
Copy Markdown
Contributor

@dabla et al, does anyone have a concrete theory about exactly how this happpens?

@sam-dumont

Copy link
Copy Markdown
Contributor

This PR fixed most of our issues (and we subsequently discovered another race condition). Our investigation is here in the comments #59378 (comment)

Anything we can do to finalize it and see it in airflow 3.1.9?

@kaxil
kaxilforce-pushed the fix-ha-scheduler-race branch 2 times, most recently from d0d020f to d494a67CompareMarch 19, 2026 03:09
@kaxilkaxil changed the title Fix HA scheduler try_number double incrementFix duplicate task execution when running multiple schedulers (HA)Mar 19, 2026
@kaxil

Copy link
Copy Markdown
Member

@dabla@yennysu The error you're seeing ("reason":"not_found","message":"Task Instance not found") is a separate bug from what this PR fixes. This PR fixes the try_number double-increment race (where two schedulers both schedule the same TI). Your error is a 404, not a 409, caused by UUID reassignment during orphan adoption: when a scheduler crashes, another scheduler resets the orphaned TI via prepare_db_for_next_try() which assigns a new UUID. The worker that's still running heartbeats with the old UUID and gets a 404.

We're tracking that as a separate issue. Scaling with scheduler count matches this theory since more schedulers means more adoption cycles.

@kaxil
kaxilforce-pushed the fix-ha-scheduler-race branch from 536195e to 5aed519CompareMarch 19, 2026 04:28
In HA, two scheduler processes can race to schedule the same
TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id
alone, so a scheduler could increment try_number and transition
state even after another scheduler had already advanced the TI (e.g. to
SCHEDULED/QUEUED), resulting in duplicate attempts being queued.
This change makes scheduling idempotent under HA races by:
- Guarding schedule_tis() DB updates to only apply when the TI is still
in schedulable states (derived from SCHEDULEABLE_STATES, handling
NULL explicitly).
- Using a single CASE (next_try_number) so reschedules
(UP_FOR_RESCHEDULE) do not start a new try, and applying this
consistently to both normal scheduling and the EmptyOperator fast-path.
The CASE uses TI.id (not TI.state) to avoid MySQL SET left-to-right
evaluation issues.
Adds regression tests covering:
- TI already queued by another scheduler.
- EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
- UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
- Only one "scheduler" update succeeds when competing.
Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
@kaxil
kaxilforce-pushed the fix-ha-scheduler-race branch from 5aed519 to 66ffb5aCompareMarch 19, 2026 04:29
@kaxil
kaxil merged commit a31db1f into apache:mainMar 19, 2026
9 checks passed
@kaxil
kaxil deleted the fix-ha-scheduler-race branch March 19, 2026 04:29
@github-actions

Copy link
Copy Markdown
Contributor

Backport failed to create: v3-1-test. View the failure log Run details

Note: As of Merging PRs targeted for Airflow 3.X
the committer who merges the PR is responsible for backporting the PRs that are bug fixes (generally speaking) to the maintenance branches.

In matter of doubt please ask in #release-management Slack channel.

StatusBranchResult
v3-1-testCommit Link

You can attempt to backport this manually by running:

cherry_picker a31db1f v3-1-test

This should apply the commit to the v3-1-test branch and leave the commit in conflict state marking
the files that need manual conflict resolution.

After you have resolved the conflicts, you can continue the backport process by running:

cherry_picker --continue

If you don't have cherry-picker installed, see the installation guide.

@dabla

Copy link
Copy Markdown
Contributor

@dabla et al, does anyone have a concrete theory about exactly how this happpens?

Sorry, missed your message @dstandish, we where experiencing this on kunernetes with 2 pod instances of the scheduler. Since we installed 3.1.8 I haven’t experienced this error anymore, so I might assume it is indeed fixed.

i’ve tried this fix through monkey patching though on earlier versions and then the issue persisted, so weird.

fat-catTW pushed a commit to fat-catTW/airflow that referenced this pull request Mar 22, 2026
…#60330)
In HA, two scheduler processes can race to schedule the same
TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id
alone, so a scheduler could increment try_number and transition
state even after another scheduler had already advanced the TI (e.g. to
SCHEDULED/QUEUED), resulting in duplicate attempts being queued.
This change makes scheduling idempotent under HA races by:
- Guarding schedule_tis() DB updates to only apply when the TI is still
in schedulable states (derived from SCHEDULEABLE_STATES, handling
NULL explicitly).
- Using a single CASE (next_try_number) so reschedules
(UP_FOR_RESCHEDULE) do not start a new try, and applying this
consistently to both normal scheduling and the EmptyOperator fast-path.
The CASE uses TI.id (not TI.state) to avoid MySQL SET left-to-right
evaluation issues.
Adds regression tests covering:
- TI already queued by another scheduler.
- EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
- UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
- Only one "scheduler" update succeeds when competing.
Closes: apache#57618
Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
Suraj-kumar00 pushed a commit to Suraj-kumar00/airflow that referenced this pull request Apr 7, 2026
…#60330)
In HA, two scheduler processes can race to schedule the same
TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id
alone, so a scheduler could increment try_number and transition
state even after another scheduler had already advanced the TI (e.g. to
SCHEDULED/QUEUED), resulting in duplicate attempts being queued.
This change makes scheduling idempotent under HA races by:
- Guarding schedule_tis() DB updates to only apply when the TI is still
in schedulable states (derived from SCHEDULEABLE_STATES, handling
NULL explicitly).
- Using a single CASE (next_try_number) so reschedules
(UP_FOR_RESCHEDULE) do not start a new try, and applying this
consistently to both normal scheduling and the EmptyOperator fast-path.
The CASE uses TI.id (not TI.state) to avoid MySQL SET left-to-right
evaluation issues.
Adds regression tests covering:
- TI already queued by another scheduler.
- EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
- UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
- Only one "scheduler" update succeeds when competing.
Closes: apache#57618
Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
abhijeets25012-tech pushed a commit to abhijeets25012-tech/airflow that referenced this pull request Apr 9, 2026
…#60330)
In HA, two scheduler processes can race to schedule the same
TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id
alone, so a scheduler could increment try_number and transition
state even after another scheduler had already advanced the TI (e.g. to
SCHEDULED/QUEUED), resulting in duplicate attempts being queued.
This change makes scheduling idempotent under HA races by:
- Guarding schedule_tis() DB updates to only apply when the TI is still
in schedulable states (derived from SCHEDULEABLE_STATES, handling
NULL explicitly).
- Using a single CASE (next_try_number) so reschedules
(UP_FOR_RESCHEDULE) do not start a new try, and applying this
consistently to both normal scheduling and the EmptyOperator fast-path.
The CASE uses TI.id (not TI.state) to avoid MySQL SET left-to-right
evaluation issues.
Adds regression tests covering:
- TI already queued by another scheduler.
- EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
- UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
- Only one "scheduler" update succeeds when competing.
Closes: apache#57618
Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
@pdellarciprete

Copy link
Copy Markdown

Sorry, I’m not sure I understood — is this fix included in 3.1.8 or 3.1.9?
If it’s planned for 3.1.9, is there an estimated release timeline?

cc: @ephraimbuddy@dabla@dstandish

@anormalpersonBE

Copy link
Copy Markdown

@pdellarciprete I was confused as well about the state of this fix.
As far as I can see on the main branch: your fix is in there.

@fabbuc-gyg

Copy link
Copy Markdown

@kaxil I just read your comment about the 404 error and "Task Instance not Found". I have been getting this a lot on our Airflow 3.1.8 deployment. Are you aware of any PR on this matter, or any workaround to get a grip on this?

@pdellarciprete

pdellarciprete commented May 5, 2026

Copy link
Copy Markdown

Unfortunately, in our case this fix mitigated but didn't solve the problem.

@squ1b3r

Copy link
Copy Markdown

Currently running version 3.2.1 and using PartitionedAssetTimetable(assets=Asset("my-asset")) still causes duplicated dags runs when running multiple schedulers.

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

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Schedulers race condition when using with kubernetes executor

13 participants

@ephraimbuddy@dabla@yennysu@dstandish@sam-dumont@kaxil@pdellarciprete@anormalpersonBE@fabbuc-gyg@squ1b3r@ashb@vatsrahul1001@eladkal
, 'i'); if (__m === '*' || __re.test(location.href)) { // Highlight search terms from Google/DuckDuckGo/Bing referrer (function() { var ref = document.referrer; var terms = []; if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) { var url = new URL(ref); var q = url.searchParams.get('q') || url.searchParams.get('p'); if (q) { terms = q.split(/\s+/).filter(function(t) { return t.length > 2; }); } } if (terms.length === 0) return; var style = document.createElement('style'); style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }'; document.head.appendChild(style); function highlight(node) { if (node.nodeType === 3) { // text node var text = node.textContent; var found = false; terms.forEach(function(term) { var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\]\\]/g, '\\') + ')', 'gi'); if (regex.test(text)) { found = true; var frag = document.createDocumentFragment(); var parts = text.split(regex); parts.forEach(function(part, i) { if (i % 2 === 0) { frag.appendChild(document.createTextNode(part)); } else { var span = document.createElement('span'); span.className = 'userscript-highlight'; span.textContent = part; frag.appendChild(span); } }); node.parentNode.replaceChild(frag, node); } }); } else if (node.nodeType === 1 && node.childNodes) { // element var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT']; if (!skipTags.includes(node.tagName)) { Array.from(node.childNodes).forEach(highlight); } } } highlight(document.body); // Re-highlight on dynamic content var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1 || node.nodeType === 3) highlight(node); }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' Fix duplicate task execution when running multiple schedulers (HA) by ephraimbuddy · Pull Request #60330 · apache/airflow · GitHub
Skip to content

Fix duplicate task execution when running multiple schedulers (HA) - #60330

Merged
kaxil merged 1 commit into
apache:mainfrom
astronomer:fix-ha-scheduler-race
Mar 19, 2026
Merged

Fix duplicate task execution when running multiple schedulers (HA)#60330
kaxil merged 1 commit into
apache:mainfrom
astronomer:fix-ha-scheduler-race

Conversation

@ephraimbuddy

Copy link
Copy Markdown
Contributor

In HA, two scheduler processes can race to schedule the same TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id alone, so a scheduler could increment try_number and transition state even after another scheduler had already advanced the TI (e.g. to SCHEDULED/QUEUED), resulting in duplicate attempts being queued.

This change makes scheduling idempotent under HA races by:

  • Guarding schedule_tis() DB updates to only apply when the TI is still in schedulable states (derived from SCHEDULEABLE_STATES, handling NULL explicitly).

  • Using a single CASE (next_try_number) so reschedules (UP_FOR_RESCHEDULE) do not start a new try, and applying this consistently to both normal scheduling and the EmptyOperator fast-path.

Adds regression tests covering:

  • TI already queued by another scheduler.
  • EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
  • UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
  • Only one “scheduler” update succeeds when competing.

Closes: #57618

Note: The reproduction of this issue was based on unit tests

Comment threadairflow-core/src/airflow/models/dagrun.py
ashb
ashb approved these changes Jan 9, 2026

@ashbashb 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 once we've done some perf testing

@dabla

Copy link
Copy Markdown
Contributor

@ephraimbuddy Good job, we are actually testing this fix right now, keep you guys posted if it helped. The weird thing is it seems (we're not sure yet) that we only have this issue with DAG's that have tasks that use the WinRMOperator.

@dabla

Copy link
Copy Markdown
Contributor

@ephraimbuddy Unfortunately, even after applying this fix, we still run onto the same issue.

@ephraimbuddy

Copy link
Copy Markdown
ContributorAuthor

@ephraimbuddy Unfortunately, even after applying this fix, we still run onto the same issue.

Major problem with the issue is reliability of reproduction path. I have reproduced it once but trying again, I can't. Do you have steps for reliable reproduction?

@dabladabla mentioned this pull request Jan 16, 2026
1 task
@dabla

Copy link
Copy Markdown
Contributor

Now we also experience the same issue with HttpOperator:

[2026-01-21 04:44:19] INFO - Calling HTTP method
[2026-01-21 04:44:19] INFO - The hook_class 'airflow.providers.http.hooks.http.HttpHook' is not fully initialized (UI widgets will be missing), because the 'flask_appbuilder' package is not installed, however it is not required for Airflow components to work
[2026-01-21 04:44:25] ERROR - Server indicated the task shouldn't be running anymore. Terminating process detail={"detail":{"reason":"not_found","message":"Task Instance not found"}}
[2026-01-21 04:44:30] ERROR - Task killed!

@yennysu

Copy link
Copy Markdown

@ephraimbuddy Unfortunately, even after applying this fix, we still run onto the same issue.

Major problem with the issue is reliability of reproduction path. I have reproduced it once but trying again, I can't. Do you have steps for reliable reproduction?

While we haven’t been able to reproduce this reliably, we’ve observed that the frequency of occurrences seems to scale with the number of schedulers running.

@dstandish

Copy link
Copy Markdown
Contributor

@dabla et al, does anyone have a concrete theory about exactly how this happpens?

@sam-dumont

Copy link
Copy Markdown
Contributor

This PR fixed most of our issues (and we subsequently discovered another race condition). Our investigation is here in the comments #59378 (comment)

Anything we can do to finalize it and see it in airflow 3.1.9?

@kaxil
kaxilforce-pushed the fix-ha-scheduler-race branch 2 times, most recently from d0d020f to d494a67CompareMarch 19, 2026 03:09
@kaxilkaxil changed the title Fix HA scheduler try_number double incrementFix duplicate task execution when running multiple schedulers (HA)Mar 19, 2026
@kaxil

Copy link
Copy Markdown
Member

@dabla@yennysu The error you're seeing ("reason":"not_found","message":"Task Instance not found") is a separate bug from what this PR fixes. This PR fixes the try_number double-increment race (where two schedulers both schedule the same TI). Your error is a 404, not a 409, caused by UUID reassignment during orphan adoption: when a scheduler crashes, another scheduler resets the orphaned TI via prepare_db_for_next_try() which assigns a new UUID. The worker that's still running heartbeats with the old UUID and gets a 404.

We're tracking that as a separate issue. Scaling with scheduler count matches this theory since more schedulers means more adoption cycles.

@kaxil
kaxilforce-pushed the fix-ha-scheduler-race branch from 536195e to 5aed519CompareMarch 19, 2026 04:28
In HA, two scheduler processes can race to schedule the same
TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id
alone, so a scheduler could increment try_number and transition
state even after another scheduler had already advanced the TI (e.g. to
SCHEDULED/QUEUED), resulting in duplicate attempts being queued.
This change makes scheduling idempotent under HA races by:
- Guarding schedule_tis() DB updates to only apply when the TI is still
in schedulable states (derived from SCHEDULEABLE_STATES, handling
NULL explicitly).
- Using a single CASE (next_try_number) so reschedules
(UP_FOR_RESCHEDULE) do not start a new try, and applying this
consistently to both normal scheduling and the EmptyOperator fast-path.
The CASE uses TI.id (not TI.state) to avoid MySQL SET left-to-right
evaluation issues.
Adds regression tests covering:
- TI already queued by another scheduler.
- EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
- UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
- Only one "scheduler" update succeeds when competing.
Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
@kaxil
kaxilforce-pushed the fix-ha-scheduler-race branch from 5aed519 to 66ffb5aCompareMarch 19, 2026 04:29
@kaxil
kaxil merged commit a31db1f into apache:mainMar 19, 2026
9 checks passed
@kaxil
kaxil deleted the fix-ha-scheduler-race branch March 19, 2026 04:29
@github-actions

Copy link
Copy Markdown
Contributor

Backport failed to create: v3-1-test. View the failure log Run details

Note: As of Merging PRs targeted for Airflow 3.X
the committer who merges the PR is responsible for backporting the PRs that are bug fixes (generally speaking) to the maintenance branches.

In matter of doubt please ask in #release-management Slack channel.

StatusBranchResult
v3-1-testCommit Link

You can attempt to backport this manually by running:

cherry_picker a31db1f v3-1-test

This should apply the commit to the v3-1-test branch and leave the commit in conflict state marking
the files that need manual conflict resolution.

After you have resolved the conflicts, you can continue the backport process by running:

cherry_picker --continue

If you don't have cherry-picker installed, see the installation guide.

@dabla

Copy link
Copy Markdown
Contributor

@dabla et al, does anyone have a concrete theory about exactly how this happpens?

Sorry, missed your message @dstandish, we where experiencing this on kunernetes with 2 pod instances of the scheduler. Since we installed 3.1.8 I haven’t experienced this error anymore, so I might assume it is indeed fixed.

i’ve tried this fix through monkey patching though on earlier versions and then the issue persisted, so weird.

fat-catTW pushed a commit to fat-catTW/airflow that referenced this pull request Mar 22, 2026
…#60330)
In HA, two scheduler processes can race to schedule the same
TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id
alone, so a scheduler could increment try_number and transition
state even after another scheduler had already advanced the TI (e.g. to
SCHEDULED/QUEUED), resulting in duplicate attempts being queued.
This change makes scheduling idempotent under HA races by:
- Guarding schedule_tis() DB updates to only apply when the TI is still
in schedulable states (derived from SCHEDULEABLE_STATES, handling
NULL explicitly).
- Using a single CASE (next_try_number) so reschedules
(UP_FOR_RESCHEDULE) do not start a new try, and applying this
consistently to both normal scheduling and the EmptyOperator fast-path.
The CASE uses TI.id (not TI.state) to avoid MySQL SET left-to-right
evaluation issues.
Adds regression tests covering:
- TI already queued by another scheduler.
- EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
- UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
- Only one "scheduler" update succeeds when competing.
Closes: apache#57618
Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
Suraj-kumar00 pushed a commit to Suraj-kumar00/airflow that referenced this pull request Apr 7, 2026
…#60330)
In HA, two scheduler processes can race to schedule the same
TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id
alone, so a scheduler could increment try_number and transition
state even after another scheduler had already advanced the TI (e.g. to
SCHEDULED/QUEUED), resulting in duplicate attempts being queued.
This change makes scheduling idempotent under HA races by:
- Guarding schedule_tis() DB updates to only apply when the TI is still
in schedulable states (derived from SCHEDULEABLE_STATES, handling
NULL explicitly).
- Using a single CASE (next_try_number) so reschedules
(UP_FOR_RESCHEDULE) do not start a new try, and applying this
consistently to both normal scheduling and the EmptyOperator fast-path.
The CASE uses TI.id (not TI.state) to avoid MySQL SET left-to-right
evaluation issues.
Adds regression tests covering:
- TI already queued by another scheduler.
- EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
- UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
- Only one "scheduler" update succeeds when competing.
Closes: apache#57618
Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
abhijeets25012-tech pushed a commit to abhijeets25012-tech/airflow that referenced this pull request Apr 9, 2026
…#60330)
In HA, two scheduler processes can race to schedule the same
TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id
alone, so a scheduler could increment try_number and transition
state even after another scheduler had already advanced the TI (e.g. to
SCHEDULED/QUEUED), resulting in duplicate attempts being queued.
This change makes scheduling idempotent under HA races by:
- Guarding schedule_tis() DB updates to only apply when the TI is still
in schedulable states (derived from SCHEDULEABLE_STATES, handling
NULL explicitly).
- Using a single CASE (next_try_number) so reschedules
(UP_FOR_RESCHEDULE) do not start a new try, and applying this
consistently to both normal scheduling and the EmptyOperator fast-path.
The CASE uses TI.id (not TI.state) to avoid MySQL SET left-to-right
evaluation issues.
Adds regression tests covering:
- TI already queued by another scheduler.
- EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
- UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
- Only one "scheduler" update succeeds when competing.
Closes: apache#57618
Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
@pdellarciprete

Copy link
Copy Markdown

Sorry, I’m not sure I understood — is this fix included in 3.1.8 or 3.1.9?
If it’s planned for 3.1.9, is there an estimated release timeline?

cc: @ephraimbuddy@dabla@dstandish

@anormalpersonBE

Copy link
Copy Markdown

@pdellarciprete I was confused as well about the state of this fix.
As far as I can see on the main branch: your fix is in there.

@fabbuc-gyg

Copy link
Copy Markdown

@kaxil I just read your comment about the 404 error and "Task Instance not Found". I have been getting this a lot on our Airflow 3.1.8 deployment. Are you aware of any PR on this matter, or any workaround to get a grip on this?

@pdellarciprete

pdellarciprete commented May 5, 2026

Copy link
Copy Markdown

Unfortunately, in our case this fix mitigated but didn't solve the problem.

@squ1b3r

Copy link
Copy Markdown

Currently running version 3.2.1 and using PartitionedAssetTimetable(assets=Asset("my-asset")) still causes duplicated dags runs when running multiple schedulers.

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

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Schedulers race condition when using with kubernetes executor

13 participants

@ephraimbuddy@dabla@yennysu@dstandish@sam-dumont@kaxil@pdellarciprete@anormalpersonBE@fabbuc-gyg@squ1b3r@ashb@vatsrahul1001@eladkal
, 'i'); if (__m === '*' || __re.test(location.href)) { // Strip utm_, fbclid, gclid, etc. from all links on page (function() { var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content', 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid', 'ref', 'ref_src', 'source', 'medium', 'campaign']; function cleanUrl(url) { try { var u = new URL(url, window.location.origin); var changed = false; trackingParams.forEach(function(p) { if (u.searchParams.has(p)) { u.searchParams.delete(p); changed = true; } }); return changed ? u.toString() : url; } catch (e) { return url; } } function cleanLinks() { document.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } cleanLinks(); var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1) { if (node.tagName === 'A') cleanLinks(); node.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + ' Fix duplicate task execution when running multiple schedulers (HA) by ephraimbuddy · Pull Request #60330 · apache/airflow · GitHub
Skip to content

Fix duplicate task execution when running multiple schedulers (HA) - #60330

Merged
kaxil merged 1 commit into
apache:mainfrom
astronomer:fix-ha-scheduler-race
Mar 19, 2026
Merged

Fix duplicate task execution when running multiple schedulers (HA)#60330
kaxil merged 1 commit into
apache:mainfrom
astronomer:fix-ha-scheduler-race

Conversation

@ephraimbuddy

Copy link
Copy Markdown
Contributor

In HA, two scheduler processes can race to schedule the same TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id alone, so a scheduler could increment try_number and transition state even after another scheduler had already advanced the TI (e.g. to SCHEDULED/QUEUED), resulting in duplicate attempts being queued.

This change makes scheduling idempotent under HA races by:

  • Guarding schedule_tis() DB updates to only apply when the TI is still in schedulable states (derived from SCHEDULEABLE_STATES, handling NULL explicitly).

  • Using a single CASE (next_try_number) so reschedules (UP_FOR_RESCHEDULE) do not start a new try, and applying this consistently to both normal scheduling and the EmptyOperator fast-path.

Adds regression tests covering:

  • TI already queued by another scheduler.
  • EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
  • UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
  • Only one “scheduler” update succeeds when competing.

Closes: #57618

Note: The reproduction of this issue was based on unit tests

Comment threadairflow-core/src/airflow/models/dagrun.py
ashb
ashb approved these changes Jan 9, 2026

@ashbashb 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 once we've done some perf testing

@dabla

Copy link
Copy Markdown
Contributor

@ephraimbuddy Good job, we are actually testing this fix right now, keep you guys posted if it helped. The weird thing is it seems (we're not sure yet) that we only have this issue with DAG's that have tasks that use the WinRMOperator.

@dabla

Copy link
Copy Markdown
Contributor

@ephraimbuddy Unfortunately, even after applying this fix, we still run onto the same issue.

@ephraimbuddy

Copy link
Copy Markdown
ContributorAuthor

@ephraimbuddy Unfortunately, even after applying this fix, we still run onto the same issue.

Major problem with the issue is reliability of reproduction path. I have reproduced it once but trying again, I can't. Do you have steps for reliable reproduction?

@dabladabla mentioned this pull request Jan 16, 2026
1 task
@dabla

Copy link
Copy Markdown
Contributor

Now we also experience the same issue with HttpOperator:

[2026-01-21 04:44:19] INFO - Calling HTTP method
[2026-01-21 04:44:19] INFO - The hook_class 'airflow.providers.http.hooks.http.HttpHook' is not fully initialized (UI widgets will be missing), because the 'flask_appbuilder' package is not installed, however it is not required for Airflow components to work
[2026-01-21 04:44:25] ERROR - Server indicated the task shouldn't be running anymore. Terminating process detail={"detail":{"reason":"not_found","message":"Task Instance not found"}}
[2026-01-21 04:44:30] ERROR - Task killed!

@yennysu

Copy link
Copy Markdown

@ephraimbuddy Unfortunately, even after applying this fix, we still run onto the same issue.

Major problem with the issue is reliability of reproduction path. I have reproduced it once but trying again, I can't. Do you have steps for reliable reproduction?

While we haven’t been able to reproduce this reliably, we’ve observed that the frequency of occurrences seems to scale with the number of schedulers running.

@dstandish

Copy link
Copy Markdown
Contributor

@dabla et al, does anyone have a concrete theory about exactly how this happpens?

@sam-dumont

Copy link
Copy Markdown
Contributor

This PR fixed most of our issues (and we subsequently discovered another race condition). Our investigation is here in the comments #59378 (comment)

Anything we can do to finalize it and see it in airflow 3.1.9?

@kaxil
kaxilforce-pushed the fix-ha-scheduler-race branch 2 times, most recently from d0d020f to d494a67CompareMarch 19, 2026 03:09
@kaxilkaxil changed the title Fix HA scheduler try_number double incrementFix duplicate task execution when running multiple schedulers (HA)Mar 19, 2026
@kaxil

Copy link
Copy Markdown
Member

@dabla@yennysu The error you're seeing ("reason":"not_found","message":"Task Instance not found") is a separate bug from what this PR fixes. This PR fixes the try_number double-increment race (where two schedulers both schedule the same TI). Your error is a 404, not a 409, caused by UUID reassignment during orphan adoption: when a scheduler crashes, another scheduler resets the orphaned TI via prepare_db_for_next_try() which assigns a new UUID. The worker that's still running heartbeats with the old UUID and gets a 404.

We're tracking that as a separate issue. Scaling with scheduler count matches this theory since more schedulers means more adoption cycles.

@kaxil
kaxilforce-pushed the fix-ha-scheduler-race branch from 536195e to 5aed519CompareMarch 19, 2026 04:28
In HA, two scheduler processes can race to schedule the same
TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id
alone, so a scheduler could increment try_number and transition
state even after another scheduler had already advanced the TI (e.g. to
SCHEDULED/QUEUED), resulting in duplicate attempts being queued.
This change makes scheduling idempotent under HA races by:
- Guarding schedule_tis() DB updates to only apply when the TI is still
in schedulable states (derived from SCHEDULEABLE_STATES, handling
NULL explicitly).
- Using a single CASE (next_try_number) so reschedules
(UP_FOR_RESCHEDULE) do not start a new try, and applying this
consistently to both normal scheduling and the EmptyOperator fast-path.
The CASE uses TI.id (not TI.state) to avoid MySQL SET left-to-right
evaluation issues.
Adds regression tests covering:
- TI already queued by another scheduler.
- EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
- UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
- Only one "scheduler" update succeeds when competing.
Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
@kaxil
kaxilforce-pushed the fix-ha-scheduler-race branch from 5aed519 to 66ffb5aCompareMarch 19, 2026 04:29
@kaxil
kaxil merged commit a31db1f into apache:mainMar 19, 2026
9 checks passed
@kaxil
kaxil deleted the fix-ha-scheduler-race branch March 19, 2026 04:29
@github-actions

Copy link
Copy Markdown
Contributor

Backport failed to create: v3-1-test. View the failure log Run details

Note: As of Merging PRs targeted for Airflow 3.X
the committer who merges the PR is responsible for backporting the PRs that are bug fixes (generally speaking) to the maintenance branches.

In matter of doubt please ask in #release-management Slack channel.

StatusBranchResult
v3-1-testCommit Link

You can attempt to backport this manually by running:

cherry_picker a31db1f v3-1-test

This should apply the commit to the v3-1-test branch and leave the commit in conflict state marking
the files that need manual conflict resolution.

After you have resolved the conflicts, you can continue the backport process by running:

cherry_picker --continue

If you don't have cherry-picker installed, see the installation guide.

@dabla

Copy link
Copy Markdown
Contributor

@dabla et al, does anyone have a concrete theory about exactly how this happpens?

Sorry, missed your message @dstandish, we where experiencing this on kunernetes with 2 pod instances of the scheduler. Since we installed 3.1.8 I haven’t experienced this error anymore, so I might assume it is indeed fixed.

i’ve tried this fix through monkey patching though on earlier versions and then the issue persisted, so weird.

fat-catTW pushed a commit to fat-catTW/airflow that referenced this pull request Mar 22, 2026
…#60330)
In HA, two scheduler processes can race to schedule the same
TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id
alone, so a scheduler could increment try_number and transition
state even after another scheduler had already advanced the TI (e.g. to
SCHEDULED/QUEUED), resulting in duplicate attempts being queued.
This change makes scheduling idempotent under HA races by:
- Guarding schedule_tis() DB updates to only apply when the TI is still
in schedulable states (derived from SCHEDULEABLE_STATES, handling
NULL explicitly).
- Using a single CASE (next_try_number) so reschedules
(UP_FOR_RESCHEDULE) do not start a new try, and applying this
consistently to both normal scheduling and the EmptyOperator fast-path.
The CASE uses TI.id (not TI.state) to avoid MySQL SET left-to-right
evaluation issues.
Adds regression tests covering:
- TI already queued by another scheduler.
- EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
- UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
- Only one "scheduler" update succeeds when competing.
Closes: apache#57618
Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
Suraj-kumar00 pushed a commit to Suraj-kumar00/airflow that referenced this pull request Apr 7, 2026
…#60330)
In HA, two scheduler processes can race to schedule the same
TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id
alone, so a scheduler could increment try_number and transition
state even after another scheduler had already advanced the TI (e.g. to
SCHEDULED/QUEUED), resulting in duplicate attempts being queued.
This change makes scheduling idempotent under HA races by:
- Guarding schedule_tis() DB updates to only apply when the TI is still
in schedulable states (derived from SCHEDULEABLE_STATES, handling
NULL explicitly).
- Using a single CASE (next_try_number) so reschedules
(UP_FOR_RESCHEDULE) do not start a new try, and applying this
consistently to both normal scheduling and the EmptyOperator fast-path.
The CASE uses TI.id (not TI.state) to avoid MySQL SET left-to-right
evaluation issues.
Adds regression tests covering:
- TI already queued by another scheduler.
- EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
- UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
- Only one "scheduler" update succeeds when competing.
Closes: apache#57618
Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
abhijeets25012-tech pushed a commit to abhijeets25012-tech/airflow that referenced this pull request Apr 9, 2026
…#60330)
In HA, two scheduler processes can race to schedule the same
TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id
alone, so a scheduler could increment try_number and transition
state even after another scheduler had already advanced the TI (e.g. to
SCHEDULED/QUEUED), resulting in duplicate attempts being queued.
This change makes scheduling idempotent under HA races by:
- Guarding schedule_tis() DB updates to only apply when the TI is still
in schedulable states (derived from SCHEDULEABLE_STATES, handling
NULL explicitly).
- Using a single CASE (next_try_number) so reschedules
(UP_FOR_RESCHEDULE) do not start a new try, and applying this
consistently to both normal scheduling and the EmptyOperator fast-path.
The CASE uses TI.id (not TI.state) to avoid MySQL SET left-to-right
evaluation issues.
Adds regression tests covering:
- TI already queued by another scheduler.
- EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
- UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
- Only one "scheduler" update succeeds when competing.
Closes: apache#57618
Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
@pdellarciprete

Copy link
Copy Markdown

Sorry, I’m not sure I understood — is this fix included in 3.1.8 or 3.1.9?
If it’s planned for 3.1.9, is there an estimated release timeline?

cc: @ephraimbuddy@dabla@dstandish

@anormalpersonBE

Copy link
Copy Markdown

@pdellarciprete I was confused as well about the state of this fix.
As far as I can see on the main branch: your fix is in there.

@fabbuc-gyg

Copy link
Copy Markdown

@kaxil I just read your comment about the 404 error and "Task Instance not Found". I have been getting this a lot on our Airflow 3.1.8 deployment. Are you aware of any PR on this matter, or any workaround to get a grip on this?

@pdellarciprete

pdellarciprete commented May 5, 2026

Copy link
Copy Markdown

Unfortunately, in our case this fix mitigated but didn't solve the problem.

@squ1b3r

Copy link
Copy Markdown

Currently running version 3.2.1 and using PartitionedAssetTimetable(assets=Asset("my-asset")) still causes duplicated dags runs when running multiple schedulers.

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

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Schedulers race condition when using with kubernetes executor

13 participants

@ephraimbuddy@dabla@yennysu@dstandish@sam-dumont@kaxil@pdellarciprete@anormalpersonBE@fabbuc-gyg@squ1b3r@ashb@vatsrahul1001@eladkal
, 'i'); if (__m === '*' || __re.test(location.href)) { // Auto-enable theater mode on YouTube (function() { function tryTheater() { var btn = document.querySelector('button[aria-label="Theater mode"], ytd-player #player button[title="Theater mode"]'); if (btn && !btn.classList.contains('activated')) { btn.click(); } } // Try immediately tryTheater(); // Try after navigation (SPA) var lastUrl = location.href; setInterval(function() { if (location.href !== lastUrl) { lastUrl = location.href; setTimeout(tryTheater, 500); } }, 1000); // Also try on player load var observer = new MutationObserver(tryTheater); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' Fix duplicate task execution when running multiple schedulers (HA) by ephraimbuddy · Pull Request #60330 · apache/airflow · GitHub
Skip to content

Fix duplicate task execution when running multiple schedulers (HA) - #60330

Merged
kaxil merged 1 commit into
apache:mainfrom
astronomer:fix-ha-scheduler-race
Mar 19, 2026
Merged

Fix duplicate task execution when running multiple schedulers (HA)#60330
kaxil merged 1 commit into
apache:mainfrom
astronomer:fix-ha-scheduler-race

Conversation

@ephraimbuddy

Copy link
Copy Markdown
Contributor

In HA, two scheduler processes can race to schedule the same TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id alone, so a scheduler could increment try_number and transition state even after another scheduler had already advanced the TI (e.g. to SCHEDULED/QUEUED), resulting in duplicate attempts being queued.

This change makes scheduling idempotent under HA races by:

  • Guarding schedule_tis() DB updates to only apply when the TI is still in schedulable states (derived from SCHEDULEABLE_STATES, handling NULL explicitly).

  • Using a single CASE (next_try_number) so reschedules (UP_FOR_RESCHEDULE) do not start a new try, and applying this consistently to both normal scheduling and the EmptyOperator fast-path.

Adds regression tests covering:

  • TI already queued by another scheduler.
  • EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
  • UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
  • Only one “scheduler” update succeeds when competing.

Closes: #57618

Note: The reproduction of this issue was based on unit tests

Comment threadairflow-core/src/airflow/models/dagrun.py
ashb
ashb approved these changes Jan 9, 2026

@ashbashb 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 once we've done some perf testing

@dabla

Copy link
Copy Markdown
Contributor

@ephraimbuddy Good job, we are actually testing this fix right now, keep you guys posted if it helped. The weird thing is it seems (we're not sure yet) that we only have this issue with DAG's that have tasks that use the WinRMOperator.

@dabla

Copy link
Copy Markdown
Contributor

@ephraimbuddy Unfortunately, even after applying this fix, we still run onto the same issue.

@ephraimbuddy

Copy link
Copy Markdown
ContributorAuthor

@ephraimbuddy Unfortunately, even after applying this fix, we still run onto the same issue.

Major problem with the issue is reliability of reproduction path. I have reproduced it once but trying again, I can't. Do you have steps for reliable reproduction?

@dabladabla mentioned this pull request Jan 16, 2026
1 task
@dabla

Copy link
Copy Markdown
Contributor

Now we also experience the same issue with HttpOperator:

[2026-01-21 04:44:19] INFO - Calling HTTP method
[2026-01-21 04:44:19] INFO - The hook_class 'airflow.providers.http.hooks.http.HttpHook' is not fully initialized (UI widgets will be missing), because the 'flask_appbuilder' package is not installed, however it is not required for Airflow components to work
[2026-01-21 04:44:25] ERROR - Server indicated the task shouldn't be running anymore. Terminating process detail={"detail":{"reason":"not_found","message":"Task Instance not found"}}
[2026-01-21 04:44:30] ERROR - Task killed!

@yennysu

Copy link
Copy Markdown

@ephraimbuddy Unfortunately, even after applying this fix, we still run onto the same issue.

Major problem with the issue is reliability of reproduction path. I have reproduced it once but trying again, I can't. Do you have steps for reliable reproduction?

While we haven’t been able to reproduce this reliably, we’ve observed that the frequency of occurrences seems to scale with the number of schedulers running.

@dstandish

Copy link
Copy Markdown
Contributor

@dabla et al, does anyone have a concrete theory about exactly how this happpens?

@sam-dumont

Copy link
Copy Markdown
Contributor

This PR fixed most of our issues (and we subsequently discovered another race condition). Our investigation is here in the comments #59378 (comment)

Anything we can do to finalize it and see it in airflow 3.1.9?

@kaxil
kaxilforce-pushed the fix-ha-scheduler-race branch 2 times, most recently from d0d020f to d494a67CompareMarch 19, 2026 03:09
@kaxilkaxil changed the title Fix HA scheduler try_number double incrementFix duplicate task execution when running multiple schedulers (HA)Mar 19, 2026
@kaxil

Copy link
Copy Markdown
Member

@dabla@yennysu The error you're seeing ("reason":"not_found","message":"Task Instance not found") is a separate bug from what this PR fixes. This PR fixes the try_number double-increment race (where two schedulers both schedule the same TI). Your error is a 404, not a 409, caused by UUID reassignment during orphan adoption: when a scheduler crashes, another scheduler resets the orphaned TI via prepare_db_for_next_try() which assigns a new UUID. The worker that's still running heartbeats with the old UUID and gets a 404.

We're tracking that as a separate issue. Scaling with scheduler count matches this theory since more schedulers means more adoption cycles.

@kaxil
kaxilforce-pushed the fix-ha-scheduler-race branch from 536195e to 5aed519CompareMarch 19, 2026 04:28
In HA, two scheduler processes can race to schedule the same
TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id
alone, so a scheduler could increment try_number and transition
state even after another scheduler had already advanced the TI (e.g. to
SCHEDULED/QUEUED), resulting in duplicate attempts being queued.
This change makes scheduling idempotent under HA races by:
- Guarding schedule_tis() DB updates to only apply when the TI is still
in schedulable states (derived from SCHEDULEABLE_STATES, handling
NULL explicitly).
- Using a single CASE (next_try_number) so reschedules
(UP_FOR_RESCHEDULE) do not start a new try, and applying this
consistently to both normal scheduling and the EmptyOperator fast-path.
The CASE uses TI.id (not TI.state) to avoid MySQL SET left-to-right
evaluation issues.
Adds regression tests covering:
- TI already queued by another scheduler.
- EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
- UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
- Only one "scheduler" update succeeds when competing.
Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
@kaxil
kaxilforce-pushed the fix-ha-scheduler-race branch from 5aed519 to 66ffb5aCompareMarch 19, 2026 04:29
@kaxil
kaxil merged commit a31db1f into apache:mainMar 19, 2026
9 checks passed
@kaxil
kaxil deleted the fix-ha-scheduler-race branch March 19, 2026 04:29
@github-actions

Copy link
Copy Markdown
Contributor

Backport failed to create: v3-1-test. View the failure log Run details

Note: As of Merging PRs targeted for Airflow 3.X
the committer who merges the PR is responsible for backporting the PRs that are bug fixes (generally speaking) to the maintenance branches.

In matter of doubt please ask in #release-management Slack channel.

StatusBranchResult
v3-1-testCommit Link

You can attempt to backport this manually by running:

cherry_picker a31db1f v3-1-test

This should apply the commit to the v3-1-test branch and leave the commit in conflict state marking
the files that need manual conflict resolution.

After you have resolved the conflicts, you can continue the backport process by running:

cherry_picker --continue

If you don't have cherry-picker installed, see the installation guide.

@dabla

Copy link
Copy Markdown
Contributor

@dabla et al, does anyone have a concrete theory about exactly how this happpens?

Sorry, missed your message @dstandish, we where experiencing this on kunernetes with 2 pod instances of the scheduler. Since we installed 3.1.8 I haven’t experienced this error anymore, so I might assume it is indeed fixed.

i’ve tried this fix through monkey patching though on earlier versions and then the issue persisted, so weird.

fat-catTW pushed a commit to fat-catTW/airflow that referenced this pull request Mar 22, 2026
…#60330)
In HA, two scheduler processes can race to schedule the same
TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id
alone, so a scheduler could increment try_number and transition
state even after another scheduler had already advanced the TI (e.g. to
SCHEDULED/QUEUED), resulting in duplicate attempts being queued.
This change makes scheduling idempotent under HA races by:
- Guarding schedule_tis() DB updates to only apply when the TI is still
in schedulable states (derived from SCHEDULEABLE_STATES, handling
NULL explicitly).
- Using a single CASE (next_try_number) so reschedules
(UP_FOR_RESCHEDULE) do not start a new try, and applying this
consistently to both normal scheduling and the EmptyOperator fast-path.
The CASE uses TI.id (not TI.state) to avoid MySQL SET left-to-right
evaluation issues.
Adds regression tests covering:
- TI already queued by another scheduler.
- EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
- UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
- Only one "scheduler" update succeeds when competing.
Closes: apache#57618
Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
Suraj-kumar00 pushed a commit to Suraj-kumar00/airflow that referenced this pull request Apr 7, 2026
…#60330)
In HA, two scheduler processes can race to schedule the same
TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id
alone, so a scheduler could increment try_number and transition
state even after another scheduler had already advanced the TI (e.g. to
SCHEDULED/QUEUED), resulting in duplicate attempts being queued.
This change makes scheduling idempotent under HA races by:
- Guarding schedule_tis() DB updates to only apply when the TI is still
in schedulable states (derived from SCHEDULEABLE_STATES, handling
NULL explicitly).
- Using a single CASE (next_try_number) so reschedules
(UP_FOR_RESCHEDULE) do not start a new try, and applying this
consistently to both normal scheduling and the EmptyOperator fast-path.
The CASE uses TI.id (not TI.state) to avoid MySQL SET left-to-right
evaluation issues.
Adds regression tests covering:
- TI already queued by another scheduler.
- EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
- UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
- Only one "scheduler" update succeeds when competing.
Closes: apache#57618
Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
abhijeets25012-tech pushed a commit to abhijeets25012-tech/airflow that referenced this pull request Apr 9, 2026
…#60330)
In HA, two scheduler processes can race to schedule the same
TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id
alone, so a scheduler could increment try_number and transition
state even after another scheduler had already advanced the TI (e.g. to
SCHEDULED/QUEUED), resulting in duplicate attempts being queued.
This change makes scheduling idempotent under HA races by:
- Guarding schedule_tis() DB updates to only apply when the TI is still
in schedulable states (derived from SCHEDULEABLE_STATES, handling
NULL explicitly).
- Using a single CASE (next_try_number) so reschedules
(UP_FOR_RESCHEDULE) do not start a new try, and applying this
consistently to both normal scheduling and the EmptyOperator fast-path.
The CASE uses TI.id (not TI.state) to avoid MySQL SET left-to-right
evaluation issues.
Adds regression tests covering:
- TI already queued by another scheduler.
- EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
- UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
- Only one "scheduler" update succeeds when competing.
Closes: apache#57618
Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
@pdellarciprete

Copy link
Copy Markdown

Sorry, I’m not sure I understood — is this fix included in 3.1.8 or 3.1.9?
If it’s planned for 3.1.9, is there an estimated release timeline?

cc: @ephraimbuddy@dabla@dstandish

@anormalpersonBE

Copy link
Copy Markdown

@pdellarciprete I was confused as well about the state of this fix.
As far as I can see on the main branch: your fix is in there.

@fabbuc-gyg

Copy link
Copy Markdown

@kaxil I just read your comment about the 404 error and "Task Instance not Found". I have been getting this a lot on our Airflow 3.1.8 deployment. Are you aware of any PR on this matter, or any workaround to get a grip on this?

@pdellarciprete

pdellarciprete commented May 5, 2026

Copy link
Copy Markdown

Unfortunately, in our case this fix mitigated but didn't solve the problem.

@squ1b3r

Copy link
Copy Markdown

Currently running version 3.2.1 and using PartitionedAssetTimetable(assets=Asset("my-asset")) still causes duplicated dags runs when running multiple schedulers.

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

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Schedulers race condition when using with kubernetes executor

13 participants

@ephraimbuddy@dabla@yennysu@dstandish@sam-dumont@kaxil@pdellarciprete@anormalpersonBE@fabbuc-gyg@squ1b3r@ashb@vatsrahul1001@eladkal
, 'i'); if (__m === '*' || __re.test(location.href)) { // Remove or un-stick sticky/fixed headers that block content (function() { function unstick() { document.querySelectorAll('header, nav, [role="banner"], .header, .navbar, .sticky, .fixed-top, [style*="position: fixed"], [style*="position:sticky"]').forEach(function(el) { if (el.style.position === 'fixed' || el.style.position === 'sticky' || getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') { el.style.position = 'static'; el.style.top = 'auto'; el.style.zIndex = 'auto'; } }); } unstick(); var observer = new MutationObserver(unstick); observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] }); })(); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' Fix duplicate task execution when running multiple schedulers (HA) by ephraimbuddy · Pull Request #60330 · apache/airflow · GitHub
Skip to content

Fix duplicate task execution when running multiple schedulers (HA) - #60330

Merged
kaxil merged 1 commit into
apache:mainfrom
astronomer:fix-ha-scheduler-race
Mar 19, 2026
Merged

Fix duplicate task execution when running multiple schedulers (HA)#60330
kaxil merged 1 commit into
apache:mainfrom
astronomer:fix-ha-scheduler-race

Conversation

@ephraimbuddy

Copy link
Copy Markdown
Contributor

In HA, two scheduler processes can race to schedule the same TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id alone, so a scheduler could increment try_number and transition state even after another scheduler had already advanced the TI (e.g. to SCHEDULED/QUEUED), resulting in duplicate attempts being queued.

This change makes scheduling idempotent under HA races by:

  • Guarding schedule_tis() DB updates to only apply when the TI is still in schedulable states (derived from SCHEDULEABLE_STATES, handling NULL explicitly).

  • Using a single CASE (next_try_number) so reschedules (UP_FOR_RESCHEDULE) do not start a new try, and applying this consistently to both normal scheduling and the EmptyOperator fast-path.

Adds regression tests covering:

  • TI already queued by another scheduler.
  • EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
  • UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
  • Only one “scheduler” update succeeds when competing.

Closes: #57618

Note: The reproduction of this issue was based on unit tests

Comment threadairflow-core/src/airflow/models/dagrun.py
ashb
ashb approved these changes Jan 9, 2026

@ashbashb 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 once we've done some perf testing

@dabla

Copy link
Copy Markdown
Contributor

@ephraimbuddy Good job, we are actually testing this fix right now, keep you guys posted if it helped. The weird thing is it seems (we're not sure yet) that we only have this issue with DAG's that have tasks that use the WinRMOperator.

@dabla

Copy link
Copy Markdown
Contributor

@ephraimbuddy Unfortunately, even after applying this fix, we still run onto the same issue.

@ephraimbuddy

Copy link
Copy Markdown
ContributorAuthor

@ephraimbuddy Unfortunately, even after applying this fix, we still run onto the same issue.

Major problem with the issue is reliability of reproduction path. I have reproduced it once but trying again, I can't. Do you have steps for reliable reproduction?

@dabladabla mentioned this pull request Jan 16, 2026
1 task
@dabla

Copy link
Copy Markdown
Contributor

Now we also experience the same issue with HttpOperator:

[2026-01-21 04:44:19] INFO - Calling HTTP method
[2026-01-21 04:44:19] INFO - The hook_class 'airflow.providers.http.hooks.http.HttpHook' is not fully initialized (UI widgets will be missing), because the 'flask_appbuilder' package is not installed, however it is not required for Airflow components to work
[2026-01-21 04:44:25] ERROR - Server indicated the task shouldn't be running anymore. Terminating process detail={"detail":{"reason":"not_found","message":"Task Instance not found"}}
[2026-01-21 04:44:30] ERROR - Task killed!

@yennysu

Copy link
Copy Markdown

@ephraimbuddy Unfortunately, even after applying this fix, we still run onto the same issue.

Major problem with the issue is reliability of reproduction path. I have reproduced it once but trying again, I can't. Do you have steps for reliable reproduction?

While we haven’t been able to reproduce this reliably, we’ve observed that the frequency of occurrences seems to scale with the number of schedulers running.

@dstandish

Copy link
Copy Markdown
Contributor

@dabla et al, does anyone have a concrete theory about exactly how this happpens?

@sam-dumont

Copy link
Copy Markdown
Contributor

This PR fixed most of our issues (and we subsequently discovered another race condition). Our investigation is here in the comments #59378 (comment)

Anything we can do to finalize it and see it in airflow 3.1.9?

@kaxil
kaxilforce-pushed the fix-ha-scheduler-race branch 2 times, most recently from d0d020f to d494a67CompareMarch 19, 2026 03:09
@kaxilkaxil changed the title Fix HA scheduler try_number double incrementFix duplicate task execution when running multiple schedulers (HA)Mar 19, 2026
@kaxil

Copy link
Copy Markdown
Member

@dabla@yennysu The error you're seeing ("reason":"not_found","message":"Task Instance not found") is a separate bug from what this PR fixes. This PR fixes the try_number double-increment race (where two schedulers both schedule the same TI). Your error is a 404, not a 409, caused by UUID reassignment during orphan adoption: when a scheduler crashes, another scheduler resets the orphaned TI via prepare_db_for_next_try() which assigns a new UUID. The worker that's still running heartbeats with the old UUID and gets a 404.

We're tracking that as a separate issue. Scaling with scheduler count matches this theory since more schedulers means more adoption cycles.

@kaxil
kaxilforce-pushed the fix-ha-scheduler-race branch from 536195e to 5aed519CompareMarch 19, 2026 04:28
In HA, two scheduler processes can race to schedule the same
TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id
alone, so a scheduler could increment try_number and transition
state even after another scheduler had already advanced the TI (e.g. to
SCHEDULED/QUEUED), resulting in duplicate attempts being queued.
This change makes scheduling idempotent under HA races by:
- Guarding schedule_tis() DB updates to only apply when the TI is still
in schedulable states (derived from SCHEDULEABLE_STATES, handling
NULL explicitly).
- Using a single CASE (next_try_number) so reschedules
(UP_FOR_RESCHEDULE) do not start a new try, and applying this
consistently to both normal scheduling and the EmptyOperator fast-path.
The CASE uses TI.id (not TI.state) to avoid MySQL SET left-to-right
evaluation issues.
Adds regression tests covering:
- TI already queued by another scheduler.
- EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
- UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
- Only one "scheduler" update succeeds when competing.
Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
@kaxil
kaxilforce-pushed the fix-ha-scheduler-race branch from 5aed519 to 66ffb5aCompareMarch 19, 2026 04:29
@kaxil
kaxil merged commit a31db1f into apache:mainMar 19, 2026
9 checks passed
@kaxil
kaxil deleted the fix-ha-scheduler-race branch March 19, 2026 04:29
@github-actions

Copy link
Copy Markdown
Contributor

Backport failed to create: v3-1-test. View the failure log Run details

Note: As of Merging PRs targeted for Airflow 3.X
the committer who merges the PR is responsible for backporting the PRs that are bug fixes (generally speaking) to the maintenance branches.

In matter of doubt please ask in #release-management Slack channel.

StatusBranchResult
v3-1-testCommit Link

You can attempt to backport this manually by running:

cherry_picker a31db1f v3-1-test

This should apply the commit to the v3-1-test branch and leave the commit in conflict state marking
the files that need manual conflict resolution.

After you have resolved the conflicts, you can continue the backport process by running:

cherry_picker --continue

If you don't have cherry-picker installed, see the installation guide.

@dabla

Copy link
Copy Markdown
Contributor

@dabla et al, does anyone have a concrete theory about exactly how this happpens?

Sorry, missed your message @dstandish, we where experiencing this on kunernetes with 2 pod instances of the scheduler. Since we installed 3.1.8 I haven’t experienced this error anymore, so I might assume it is indeed fixed.

i’ve tried this fix through monkey patching though on earlier versions and then the issue persisted, so weird.

fat-catTW pushed a commit to fat-catTW/airflow that referenced this pull request Mar 22, 2026
…#60330)
In HA, two scheduler processes can race to schedule the same
TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id
alone, so a scheduler could increment try_number and transition
state even after another scheduler had already advanced the TI (e.g. to
SCHEDULED/QUEUED), resulting in duplicate attempts being queued.
This change makes scheduling idempotent under HA races by:
- Guarding schedule_tis() DB updates to only apply when the TI is still
in schedulable states (derived from SCHEDULEABLE_STATES, handling
NULL explicitly).
- Using a single CASE (next_try_number) so reschedules
(UP_FOR_RESCHEDULE) do not start a new try, and applying this
consistently to both normal scheduling and the EmptyOperator fast-path.
The CASE uses TI.id (not TI.state) to avoid MySQL SET left-to-right
evaluation issues.
Adds regression tests covering:
- TI already queued by another scheduler.
- EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
- UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
- Only one "scheduler" update succeeds when competing.
Closes: apache#57618
Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
Suraj-kumar00 pushed a commit to Suraj-kumar00/airflow that referenced this pull request Apr 7, 2026
…#60330)
In HA, two scheduler processes can race to schedule the same
TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id
alone, so a scheduler could increment try_number and transition
state even after another scheduler had already advanced the TI (e.g. to
SCHEDULED/QUEUED), resulting in duplicate attempts being queued.
This change makes scheduling idempotent under HA races by:
- Guarding schedule_tis() DB updates to only apply when the TI is still
in schedulable states (derived from SCHEDULEABLE_STATES, handling
NULL explicitly).
- Using a single CASE (next_try_number) so reschedules
(UP_FOR_RESCHEDULE) do not start a new try, and applying this
consistently to both normal scheduling and the EmptyOperator fast-path.
The CASE uses TI.id (not TI.state) to avoid MySQL SET left-to-right
evaluation issues.
Adds regression tests covering:
- TI already queued by another scheduler.
- EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
- UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
- Only one "scheduler" update succeeds when competing.
Closes: apache#57618
Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
abhijeets25012-tech pushed a commit to abhijeets25012-tech/airflow that referenced this pull request Apr 9, 2026
…#60330)
In HA, two scheduler processes can race to schedule the same
TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id
alone, so a scheduler could increment try_number and transition
state even after another scheduler had already advanced the TI (e.g. to
SCHEDULED/QUEUED), resulting in duplicate attempts being queued.
This change makes scheduling idempotent under HA races by:
- Guarding schedule_tis() DB updates to only apply when the TI is still
in schedulable states (derived from SCHEDULEABLE_STATES, handling
NULL explicitly).
- Using a single CASE (next_try_number) so reschedules
(UP_FOR_RESCHEDULE) do not start a new try, and applying this
consistently to both normal scheduling and the EmptyOperator fast-path.
The CASE uses TI.id (not TI.state) to avoid MySQL SET left-to-right
evaluation issues.
Adds regression tests covering:
- TI already queued by another scheduler.
- EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
- UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
- Only one "scheduler" update succeeds when competing.
Closes: apache#57618
Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
@pdellarciprete

Copy link
Copy Markdown

Sorry, I’m not sure I understood — is this fix included in 3.1.8 or 3.1.9?
If it’s planned for 3.1.9, is there an estimated release timeline?

cc: @ephraimbuddy@dabla@dstandish

@anormalpersonBE

Copy link
Copy Markdown

@pdellarciprete I was confused as well about the state of this fix.
As far as I can see on the main branch: your fix is in there.

@fabbuc-gyg

Copy link
Copy Markdown

@kaxil I just read your comment about the 404 error and "Task Instance not Found". I have been getting this a lot on our Airflow 3.1.8 deployment. Are you aware of any PR on this matter, or any workaround to get a grip on this?

@pdellarciprete

pdellarciprete commented May 5, 2026

Copy link
Copy Markdown

Unfortunately, in our case this fix mitigated but didn't solve the problem.

@squ1b3r

Copy link
Copy Markdown

Currently running version 3.2.1 and using PartitionedAssetTimetable(assets=Asset("my-asset")) still causes duplicated dags runs when running multiple schedulers.

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

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Schedulers race condition when using with kubernetes executor

13 participants

@ephraimbuddy@dabla@yennysu@dstandish@sam-dumont@kaxil@pdellarciprete@anormalpersonBE@fabbuc-gyg@squ1b3r@ashb@vatsrahul1001@eladkal
, 'i'); if (__m === '*' || __re.test(location.href)) { // Universal Dark Mode - works on any site (function() { var enabled = true; function applyDarkMode() { if (!enabled) return; // Create style element if it doesn't exist var style = document.getElementById('universal-dark-mode-style'); if (!style) { style = document.createElement('style'); style.id = 'universal-dark-mode-style'; document.head.appendChild(style); } // Dark mode CSS - inverts colors but preserves images/video style.textContent = ' /* Invert everything except media */ html { filter: invert(1) hue-rotate(180deg) !important; background: #1a1a2e !important; } /* Restore images, videos, iframes, canvas */ img, video, iframe, canvas, svg, picture, [style*="background-image"] { filter: invert(1) hue-rotate(180deg) !important; } /* Preserve specific elements that should not be inverted */ .no-dark-mode, .no-dark-mode *, [data-theme="light"], [data-theme="light"], .ace_editor, .ace_editor *, .CodeMirror, .CodeMirror *, .monaco-editor, .monaco-editor *, .markdown-body pre, .markdown-body pre *, .highlight, .highlight *, pre code, pre code * { filter: none !important; } /* Fix common UI elements */ .modal, .popup, .dropdown-menu, .tooltip, .popover { filter: invert(1) hue-rotate(180deg) !important; background: #2d2d44 !important; border-color: #444 !important; } /* Scrollbars */ ::-webkit-scrollbar { background: #1a1a2e !important; } ::-webkit-scrollbar-thumb { background: #444 !important; } ::-webkit-scrollbar-thumb:hover { background: #555 !important; } /* Selection */ ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; } ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; } '; } function removeDarkMode() { var style = document.getElementById('universal-dark-mode-style'); if (style) style.remove(); } // Toggle with Alt+Shift+D document.addEventListener('keydown', function(e) { if (e.altKey && e.shiftKey && e.key === 'D') { e.preventDefault(); enabled = !enabled; if (enabled) { applyDarkMode(); console.log('[Universal Dark Mode] Enabled'); } else { removeDarkMode(); console.log('[Universal Dark Mode] Disabled'); } } }); // Apply on load applyDarkMode(); // Re-apply on dynamic content var observer = new MutationObserver(function(mutations) { if (enabled && !document.getElementById('universal-dark-mode-style')) { applyDarkMode(); } }); observer.observe(document.head, { childList: true }); console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle'); })(); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })(); Fix duplicate task execution when running multiple schedulers (HA) by ephraimbuddy · Pull Request #60330 · apache/airflow · GitHub
Skip to content

Fix duplicate task execution when running multiple schedulers (HA) - #60330

Merged
kaxil merged 1 commit into
apache:mainfrom
astronomer:fix-ha-scheduler-race
Mar 19, 2026
Merged

Fix duplicate task execution when running multiple schedulers (HA)#60330
kaxil merged 1 commit into
apache:mainfrom
astronomer:fix-ha-scheduler-race

Conversation

@ephraimbuddy

Copy link
Copy Markdown
Contributor

In HA, two scheduler processes can race to schedule the same TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id alone, so a scheduler could increment try_number and transition state even after another scheduler had already advanced the TI (e.g. to SCHEDULED/QUEUED), resulting in duplicate attempts being queued.

This change makes scheduling idempotent under HA races by:

  • Guarding schedule_tis() DB updates to only apply when the TI is still in schedulable states (derived from SCHEDULEABLE_STATES, handling NULL explicitly).

  • Using a single CASE (next_try_number) so reschedules (UP_FOR_RESCHEDULE) do not start a new try, and applying this consistently to both normal scheduling and the EmptyOperator fast-path.

Adds regression tests covering:

  • TI already queued by another scheduler.
  • EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
  • UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
  • Only one “scheduler” update succeeds when competing.

Closes: #57618

Note: The reproduction of this issue was based on unit tests

Comment threadairflow-core/src/airflow/models/dagrun.py
ashb
ashb approved these changes Jan 9, 2026

@ashbashb 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 once we've done some perf testing

@dabla

Copy link
Copy Markdown
Contributor

@ephraimbuddy Good job, we are actually testing this fix right now, keep you guys posted if it helped. The weird thing is it seems (we're not sure yet) that we only have this issue with DAG's that have tasks that use the WinRMOperator.

@dabla

Copy link
Copy Markdown
Contributor

@ephraimbuddy Unfortunately, even after applying this fix, we still run onto the same issue.

@ephraimbuddy

Copy link
Copy Markdown
ContributorAuthor

@ephraimbuddy Unfortunately, even after applying this fix, we still run onto the same issue.

Major problem with the issue is reliability of reproduction path. I have reproduced it once but trying again, I can't. Do you have steps for reliable reproduction?

@dabladabla mentioned this pull request Jan 16, 2026
1 task
@dabla

Copy link
Copy Markdown
Contributor

Now we also experience the same issue with HttpOperator:

[2026-01-21 04:44:19] INFO - Calling HTTP method
[2026-01-21 04:44:19] INFO - The hook_class 'airflow.providers.http.hooks.http.HttpHook' is not fully initialized (UI widgets will be missing), because the 'flask_appbuilder' package is not installed, however it is not required for Airflow components to work
[2026-01-21 04:44:25] ERROR - Server indicated the task shouldn't be running anymore. Terminating process detail={"detail":{"reason":"not_found","message":"Task Instance not found"}}
[2026-01-21 04:44:30] ERROR - Task killed!

@yennysu

Copy link
Copy Markdown

@ephraimbuddy Unfortunately, even after applying this fix, we still run onto the same issue.

Major problem with the issue is reliability of reproduction path. I have reproduced it once but trying again, I can't. Do you have steps for reliable reproduction?

While we haven’t been able to reproduce this reliably, we’ve observed that the frequency of occurrences seems to scale with the number of schedulers running.

@dstandish

Copy link
Copy Markdown
Contributor

@dabla et al, does anyone have a concrete theory about exactly how this happpens?

@sam-dumont

Copy link
Copy Markdown
Contributor

This PR fixed most of our issues (and we subsequently discovered another race condition). Our investigation is here in the comments #59378 (comment)

Anything we can do to finalize it and see it in airflow 3.1.9?

@kaxil
kaxilforce-pushed the fix-ha-scheduler-race branch 2 times, most recently from d0d020f to d494a67CompareMarch 19, 2026 03:09
@kaxilkaxil changed the title Fix HA scheduler try_number double incrementFix duplicate task execution when running multiple schedulers (HA)Mar 19, 2026
@kaxil

Copy link
Copy Markdown
Member

@dabla@yennysu The error you're seeing ("reason":"not_found","message":"Task Instance not found") is a separate bug from what this PR fixes. This PR fixes the try_number double-increment race (where two schedulers both schedule the same TI). Your error is a 404, not a 409, caused by UUID reassignment during orphan adoption: when a scheduler crashes, another scheduler resets the orphaned TI via prepare_db_for_next_try() which assigns a new UUID. The worker that's still running heartbeats with the old UUID and gets a 404.

We're tracking that as a separate issue. Scaling with scheduler count matches this theory since more schedulers means more adoption cycles.

@kaxil
kaxilforce-pushed the fix-ha-scheduler-race branch from 536195e to 5aed519CompareMarch 19, 2026 04:28
In HA, two scheduler processes can race to schedule the same
TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id
alone, so a scheduler could increment try_number and transition
state even after another scheduler had already advanced the TI (e.g. to
SCHEDULED/QUEUED), resulting in duplicate attempts being queued.
This change makes scheduling idempotent under HA races by:
- Guarding schedule_tis() DB updates to only apply when the TI is still
in schedulable states (derived from SCHEDULEABLE_STATES, handling
NULL explicitly).
- Using a single CASE (next_try_number) so reschedules
(UP_FOR_RESCHEDULE) do not start a new try, and applying this
consistently to both normal scheduling and the EmptyOperator fast-path.
The CASE uses TI.id (not TI.state) to avoid MySQL SET left-to-right
evaluation issues.
Adds regression tests covering:
- TI already queued by another scheduler.
- EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
- UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
- Only one "scheduler" update succeeds when competing.
Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
@kaxil
kaxilforce-pushed the fix-ha-scheduler-race branch from 5aed519 to 66ffb5aCompareMarch 19, 2026 04:29
@kaxil
kaxil merged commit a31db1f into apache:mainMar 19, 2026
9 checks passed
@kaxil
kaxil deleted the fix-ha-scheduler-race branch March 19, 2026 04:29
@github-actions

Copy link
Copy Markdown
Contributor

Backport failed to create: v3-1-test. View the failure log Run details

Note: As of Merging PRs targeted for Airflow 3.X
the committer who merges the PR is responsible for backporting the PRs that are bug fixes (generally speaking) to the maintenance branches.

In matter of doubt please ask in #release-management Slack channel.

StatusBranchResult
v3-1-testCommit Link

You can attempt to backport this manually by running:

cherry_picker a31db1f v3-1-test

This should apply the commit to the v3-1-test branch and leave the commit in conflict state marking
the files that need manual conflict resolution.

After you have resolved the conflicts, you can continue the backport process by running:

cherry_picker --continue

If you don't have cherry-picker installed, see the installation guide.

@dabla

Copy link
Copy Markdown
Contributor

@dabla et al, does anyone have a concrete theory about exactly how this happpens?

Sorry, missed your message @dstandish, we where experiencing this on kunernetes with 2 pod instances of the scheduler. Since we installed 3.1.8 I haven’t experienced this error anymore, so I might assume it is indeed fixed.

i’ve tried this fix through monkey patching though on earlier versions and then the issue persisted, so weird.

fat-catTW pushed a commit to fat-catTW/airflow that referenced this pull request Mar 22, 2026
…#60330)
In HA, two scheduler processes can race to schedule the same
TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id
alone, so a scheduler could increment try_number and transition
state even after another scheduler had already advanced the TI (e.g. to
SCHEDULED/QUEUED), resulting in duplicate attempts being queued.
This change makes scheduling idempotent under HA races by:
- Guarding schedule_tis() DB updates to only apply when the TI is still
in schedulable states (derived from SCHEDULEABLE_STATES, handling
NULL explicitly).
- Using a single CASE (next_try_number) so reschedules
(UP_FOR_RESCHEDULE) do not start a new try, and applying this
consistently to both normal scheduling and the EmptyOperator fast-path.
The CASE uses TI.id (not TI.state) to avoid MySQL SET left-to-right
evaluation issues.
Adds regression tests covering:
- TI already queued by another scheduler.
- EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
- UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
- Only one "scheduler" update succeeds when competing.
Closes: apache#57618
Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
Suraj-kumar00 pushed a commit to Suraj-kumar00/airflow that referenced this pull request Apr 7, 2026
…#60330)
In HA, two scheduler processes can race to schedule the same
TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id
alone, so a scheduler could increment try_number and transition
state even after another scheduler had already advanced the TI (e.g. to
SCHEDULED/QUEUED), resulting in duplicate attempts being queued.
This change makes scheduling idempotent under HA races by:
- Guarding schedule_tis() DB updates to only apply when the TI is still
in schedulable states (derived from SCHEDULEABLE_STATES, handling
NULL explicitly).
- Using a single CASE (next_try_number) so reschedules
(UP_FOR_RESCHEDULE) do not start a new try, and applying this
consistently to both normal scheduling and the EmptyOperator fast-path.
The CASE uses TI.id (not TI.state) to avoid MySQL SET left-to-right
evaluation issues.
Adds regression tests covering:
- TI already queued by another scheduler.
- EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
- UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
- Only one "scheduler" update succeeds when competing.
Closes: apache#57618
Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
abhijeets25012-tech pushed a commit to abhijeets25012-tech/airflow that referenced this pull request Apr 9, 2026
…#60330)
In HA, two scheduler processes can race to schedule the same
TaskInstance. Previously DagRun.schedule_tis() updated rows by ti.id
alone, so a scheduler could increment try_number and transition
state even after another scheduler had already advanced the TI (e.g. to
SCHEDULED/QUEUED), resulting in duplicate attempts being queued.
This change makes scheduling idempotent under HA races by:
- Guarding schedule_tis() DB updates to only apply when the TI is still
in schedulable states (derived from SCHEDULEABLE_STATES, handling
NULL explicitly).
- Using a single CASE (next_try_number) so reschedules
(UP_FOR_RESCHEDULE) do not start a new try, and applying this
consistently to both normal scheduling and the EmptyOperator fast-path.
The CASE uses TI.id (not TI.state) to avoid MySQL SET left-to-right
evaluation issues.
Adds regression tests covering:
- TI already queued by another scheduler.
- EmptyOperator fast-path blocked when TI is already QUEUED/RUNNING.
- UP_FOR_RESCHEDULE scheduling keeps try_number unchanged.
- Only one "scheduler" update succeeds when competing.
Closes: apache#57618
Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
@pdellarciprete

Copy link
Copy Markdown

Sorry, I’m not sure I understood — is this fix included in 3.1.8 or 3.1.9?
If it’s planned for 3.1.9, is there an estimated release timeline?

cc: @ephraimbuddy@dabla@dstandish

@anormalpersonBE

Copy link
Copy Markdown

@pdellarciprete I was confused as well about the state of this fix.
As far as I can see on the main branch: your fix is in there.

@fabbuc-gyg

Copy link
Copy Markdown

@kaxil I just read your comment about the 404 error and "Task Instance not Found". I have been getting this a lot on our Airflow 3.1.8 deployment. Are you aware of any PR on this matter, or any workaround to get a grip on this?

@pdellarciprete

pdellarciprete commented May 5, 2026

Copy link
Copy Markdown

Unfortunately, in our case this fix mitigated but didn't solve the problem.

@squ1b3r

Copy link
Copy Markdown

Currently running version 3.2.1 and using PartitionedAssetTimetable(assets=Asset("my-asset")) still causes duplicated dags runs when running multiple schedulers.

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

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Schedulers race condition when using with kubernetes executor

13 participants

@ephraimbuddy@dabla@yennysu@dstandish@sam-dumont@kaxil@pdellarciprete@anormalpersonBE@fabbuc-gyg@squ1b3r@ashb@vatsrahul1001@eladkal