Fix ti_skip_downstream overwriting RUNNING tasks to SKIPPED - #63266

Merged
kaxil merged 1 commit into
apache:mainfrom
sam-dumont:fix/ti-skip-downstream-state-guard
Mar 19, 2026
Merged

Fix ti_skip_downstream overwriting RUNNING tasks to SKIPPED#63266
kaxil merged 1 commit into
apache:mainfrom
sam-dumont:fix/ti-skip-downstream-state-guard

Conversation

@sam-dumont

Copy link
Copy Markdown
Contributor

ti_skip_downstream() issues an UPDATE filtered by (dag_id, run_id, task_id, map_index) without a state guard. When a BranchOperator on one scheduler decides to skip downstream tasks, the UPDATE can overwrite a task already RUNNING on a worker. The worker's next heartbeat returns 409 with current_state: skipped, killing the task mid-execution.

This is a companion fix to #60330, which guards schedule_tis() against the same class of race condition. Different code path (Execution API routes vs dagrun.py), same root cause : unguarded bulk UPDATEs on TI state.

Production data (12 days, 5 schedulers, ~500 concurrent workers)

We deployed both fixes as monkey patches on our prod cluster and monitored 409 heartbeat errors via CloudWatch :

Before any fix 14-169 errors/day
After schedule_tis 3-4/day (all current_state: skipped)
After both fixes 0 errors for 18+ hours
MetricBefore fixesAfter schedule_tis onlyAfter both fixes
Total 409s/day14-1693-40
current_state: scheduledpresent00
current_state: failed47/day00
current_state: skipped8/day2-5/day0

Fix

Add skippable_state_clause to the UPDATE's WHERE clause :

skippable_state_clause=or_(
TI.state.is_(None),
TI.state.not_in([RUNNING, SUCCESS, FAILED]),
)

The or_(IS NULL, NOT IN) pattern handles SQL NULL semantics : NULL NOT IN (...) evaluates to NULL (falsy), so tasks with state=None need an explicit IS NULL check to remain skippable.

QUEUED is intentionally NOT guarded : a QUEUED task hasn't started executing yet, so the BranchOperator's decision should take priority. The worker pod will get a benign 409 on PATCH /run and exit cleanly. Blocking QUEUED would cause a semantic error where the wrong branch executes.

Tests

5 regression tests in TestTISkipDownstreamRaceCondition :

  • RUNNING / SUCCESS / FAILED tasks protected from overwrite (parametrized)
  • QUEUED task correctly skipped (BranchOperator decision wins over queue)
  • None-state task still correctly skipped (happy path)

related: #59378

related: #60330

related: #57618


Was generative AI tooling used to co-author this PR?
  • Yes — Claude Code (claude-opus-4-6)

Generated-by: Claude Code (claude-opus-4-6) following the guidelines

@boring-cyborgboring-cyborgBot added area:API Airflow's REST/HTTP API area:task-sdk labels Mar 10, 2026
@sam-dumont
sam-dumontforce-pushed the fix/ti-skip-downstream-state-guard branch from f4a3ded to 589633aCompareMarch 10, 2026 14:44
@potiuk

potiuk commented Mar 11, 2026

Copy link
Copy Markdown
Member

@sam-dumont This PR has been converted to draft because it does not yet meet our Pull Request quality criteria.

Issues found:

  • Pre-commit / static checks: Failing: CI image checks / Static checks. Run prek run --from-ref main locally to find and fix issues. See Pre-commit / static checks docs.
  • mypy (type checking): Failing: CI image checks / MyPy checks (mypy-airflow-core), CI image checks / MyPy checks (mypy-task-sdk). Run prek --stage manual mypy-airflow-core --all-files && prek --stage manual mypy-task-sdk --all-files locally to reproduce. You need breeze ci-image build --python 3.10 for Docker-based mypy. See mypy (type checking) docs.
  • Provider tests: Failing: provider distributions tests / Compat 2.11.1:P3.10:, Postgres tests: providers / DB-prov:Postgres:14:3.10:-amazon,celer...standard, MySQL tests: providers / DB-prov:MySQL:8.0:3.10:-amazon,celer...standard, Sqlite tests: providers / DB-prov:Sqlite:3.10:-amazon,celer...standard, Non-DB tests: providers / Non-DB-prov::3.10:-amazon,celer...standard (+7 more). Run provider tests with breeze run pytest <provider-test-path> -xvs. See Provider tests docs.
  • Other failing CI checks: Failing: CI image checks / Test Python API client, Postgres tests: core / DB-core:Postgres:14:3.10:API...Serialization, MySQL tests: core / DB-core:MySQL:8.0:3.10:API...Serialization, Sqlite tests: core / DB-core:Sqlite:3.10:API...Serialization, Non-DB tests: core / Non-DB-core::3.10:API...Serialization (+8 more). Run prek run --from-ref main locally to reproduce. See static checks docs.

Note: Your branch is 45 commits behind main. Some check failures may be caused by changes in the base branch rather than by your PR. Please rebase your branch and push again to get up-to-date CI results.

What to do next:

  • The comment informs you what you need to do.
  • Fix each issue, then mark the PR as "Ready for review" in the GitHub UI - but only after making sure that all the issues are fixed.
  • Maintainers will then proceed with a normal review.

Converting a PR to draft is not a rejection — it is an invitation to bring the PR up to the project's standards so that maintainer review time is spent productively. If you have questions, feel free to ask on the Airflow Slack.


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

@potiuk
potiuk marked this pull request as draft March 11, 2026 00:01
@sam-dumont
sam-dumontforce-pushed the fix/ti-skip-downstream-state-guard branch 4 times, most recently from 08ece10 to 9dd13fbCompareMarch 13, 2026 14:29
@sam-dumont
sam-dumont marked this pull request as ready for review March 13, 2026 14:43
@sam-dumont

Copy link
Copy Markdown
ContributorAuthor

@potiuk Hi ! fixed the issues following your advice, thanks a lot ! the 3 failures here seems to be transient, not sure if I can relaunch the job myself, does not seems like it. Hope it's going to be okay now :)

In HA deployments, ti_skip_downstream() issues a bulk UPDATE without
a state guard. When a BranchOperator decides to skip downstream tasks,
it can overwrite a task already RUNNING on a worker to SKIPPED, causing
a 409 heartbeat conflict that kills the task mid-execution.
Add a skippable_state_clause to the UPDATE WHERE clause so RUNNING,
SUCCESS, and FAILED tasks are never overwritten to SKIPPED.
QUEUED tasks are intentionally allowed to be skipped: no work has been
done yet and the BranchOperator's decision should take priority. The
worker pod will get a benign 409 on PATCH /run and exit cleanly.
related: apache#59378
@kaxil
kaxilforce-pushed the fix/ti-skip-downstream-state-guard branch from 9dd13fb to 4d57ba4CompareMarch 19, 2026 01:48
@kaxilkaxil added this to the Airflow 3.2.0 milestone Mar 19, 2026
@kaxil
kaxil merged commit 6d1794a into apache:mainMar 19, 2026
130 of 132 checks passed
fat-catTW pushed a commit to fat-catTW/airflow that referenced this pull request Mar 22, 2026
…pache#63266)
In HA deployments, ti_skip_downstream() issues a bulk UPDATE without
a state guard. When a BranchOperator decides to skip downstream tasks,
it can overwrite a task already RUNNING on a worker to SKIPPED, causing
a 409 heartbeat conflict that kills the task mid-execution.
Add a skippable_state_clause to the UPDATE WHERE clause so RUNNING,
SUCCESS, and FAILED tasks are never overwritten to SKIPPED.
QUEUED tasks are intentionally allowed to be skipped: no work has been
done yet and the BranchOperator's decision should take priority. The
worker pod will get a benign 409 on PATCH /run and exit cleanly.
closes: apache#59378closes: apache#57618
Suraj-kumar00 pushed a commit to Suraj-kumar00/airflow that referenced this pull request Apr 7, 2026
…pache#63266)
In HA deployments, ti_skip_downstream() issues a bulk UPDATE without
a state guard. When a BranchOperator decides to skip downstream tasks,
it can overwrite a task already RUNNING on a worker to SKIPPED, causing
a 409 heartbeat conflict that kills the task mid-execution.
Add a skippable_state_clause to the UPDATE WHERE clause so RUNNING,
SUCCESS, and FAILED tasks are never overwritten to SKIPPED.
QUEUED tasks are intentionally allowed to be skipped: no work has been
done yet and the BranchOperator's decision should take priority. The
worker pod will get a benign 409 on PATCH /run and exit cleanly.
closes: apache#59378closes: apache#57618
abhijeets25012-tech pushed a commit to abhijeets25012-tech/airflow that referenced this pull request Apr 9, 2026
…pache#63266)
In HA deployments, ti_skip_downstream() issues a bulk UPDATE without
a state guard. When a BranchOperator decides to skip downstream tasks,
it can overwrite a task already RUNNING on a worker to SKIPPED, causing
a 409 heartbeat conflict that kills the task mid-execution.
Add a skippable_state_clause to the UPDATE WHERE clause so RUNNING,
SUCCESS, and FAILED tasks are never overwritten to SKIPPED.
QUEUED tasks are intentionally allowed to be skipped: no work has been
done yet and the BranchOperator's decision should take priority. The
worker pod will get a benign 409 on PATCH /run and exit cleanly.
closes: apache#59378closes: apache#57618
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:APIAirflow's REST/HTTP APIarea:task-sdk

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@sam-dumont@potiuk@kaxil
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n 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;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content

Fix ti_skip_downstream overwriting RUNNING tasks to SKIPPED - #63266

Merged
kaxil merged 1 commit into
apache:mainfrom
sam-dumont:fix/ti-skip-downstream-state-guard
Mar 19, 2026
Merged

Fix ti_skip_downstream overwriting RUNNING tasks to SKIPPED#63266
kaxil merged 1 commit into
apache:mainfrom
sam-dumont:fix/ti-skip-downstream-state-guard

Conversation

@sam-dumont

Copy link
Copy Markdown
Contributor

ti_skip_downstream() issues an UPDATE filtered by (dag_id, run_id, task_id, map_index) without a state guard. When a BranchOperator on one scheduler decides to skip downstream tasks, the UPDATE can overwrite a task already RUNNING on a worker. The worker's next heartbeat returns 409 with current_state: skipped, killing the task mid-execution.

This is a companion fix to #60330, which guards schedule_tis() against the same class of race condition. Different code path (Execution API routes vs dagrun.py), same root cause : unguarded bulk UPDATEs on TI state.

Production data (12 days, 5 schedulers, ~500 concurrent workers)

We deployed both fixes as monkey patches on our prod cluster and monitored 409 heartbeat errors via CloudWatch :

Before any fix 14-169 errors/day
After schedule_tis 3-4/day (all current_state: skipped)
After both fixes 0 errors for 18+ hours
MetricBefore fixesAfter schedule_tis onlyAfter both fixes
Total 409s/day14-1693-40
current_state: scheduledpresent00
current_state: failed47/day00
current_state: skipped8/day2-5/day0

Fix

Add skippable_state_clause to the UPDATE's WHERE clause :

skippable_state_clause=or_(
TI.state.is_(None),
TI.state.not_in([RUNNING, SUCCESS, FAILED]),
)

The or_(IS NULL, NOT IN) pattern handles SQL NULL semantics : NULL NOT IN (...) evaluates to NULL (falsy), so tasks with state=None need an explicit IS NULL check to remain skippable.

QUEUED is intentionally NOT guarded : a QUEUED task hasn't started executing yet, so the BranchOperator's decision should take priority. The worker pod will get a benign 409 on PATCH /run and exit cleanly. Blocking QUEUED would cause a semantic error where the wrong branch executes.

Tests

5 regression tests in TestTISkipDownstreamRaceCondition :

  • RUNNING / SUCCESS / FAILED tasks protected from overwrite (parametrized)
  • QUEUED task correctly skipped (BranchOperator decision wins over queue)
  • None-state task still correctly skipped (happy path)

related: #59378

related: #60330

related: #57618


Was generative AI tooling used to co-author this PR?
  • Yes — Claude Code (claude-opus-4-6)

Generated-by: Claude Code (claude-opus-4-6) following the guidelines

@boring-cyborgboring-cyborgBot added area:API Airflow's REST/HTTP API area:task-sdk labels Mar 10, 2026
@sam-dumont
sam-dumontforce-pushed the fix/ti-skip-downstream-state-guard branch from f4a3ded to 589633aCompareMarch 10, 2026 14:44
@potiuk

potiuk commented Mar 11, 2026

Copy link
Copy Markdown
Member

@sam-dumont This PR has been converted to draft because it does not yet meet our Pull Request quality criteria.

Issues found:

  • Pre-commit / static checks: Failing: CI image checks / Static checks. Run prek run --from-ref main locally to find and fix issues. See Pre-commit / static checks docs.
  • mypy (type checking): Failing: CI image checks / MyPy checks (mypy-airflow-core), CI image checks / MyPy checks (mypy-task-sdk). Run prek --stage manual mypy-airflow-core --all-files && prek --stage manual mypy-task-sdk --all-files locally to reproduce. You need breeze ci-image build --python 3.10 for Docker-based mypy. See mypy (type checking) docs.
  • Provider tests: Failing: provider distributions tests / Compat 2.11.1:P3.10:, Postgres tests: providers / DB-prov:Postgres:14:3.10:-amazon,celer...standard, MySQL tests: providers / DB-prov:MySQL:8.0:3.10:-amazon,celer...standard, Sqlite tests: providers / DB-prov:Sqlite:3.10:-amazon,celer...standard, Non-DB tests: providers / Non-DB-prov::3.10:-amazon,celer...standard (+7 more). Run provider tests with breeze run pytest <provider-test-path> -xvs. See Provider tests docs.
  • Other failing CI checks: Failing: CI image checks / Test Python API client, Postgres tests: core / DB-core:Postgres:14:3.10:API...Serialization, MySQL tests: core / DB-core:MySQL:8.0:3.10:API...Serialization, Sqlite tests: core / DB-core:Sqlite:3.10:API...Serialization, Non-DB tests: core / Non-DB-core::3.10:API...Serialization (+8 more). Run prek run --from-ref main locally to reproduce. See static checks docs.

Note: Your branch is 45 commits behind main. Some check failures may be caused by changes in the base branch rather than by your PR. Please rebase your branch and push again to get up-to-date CI results.

What to do next:

  • The comment informs you what you need to do.
  • Fix each issue, then mark the PR as "Ready for review" in the GitHub UI - but only after making sure that all the issues are fixed.
  • Maintainers will then proceed with a normal review.

Converting a PR to draft is not a rejection — it is an invitation to bring the PR up to the project's standards so that maintainer review time is spent productively. If you have questions, feel free to ask on the Airflow Slack.


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

@potiuk
potiuk marked this pull request as draft March 11, 2026 00:01
@sam-dumont
sam-dumontforce-pushed the fix/ti-skip-downstream-state-guard branch 4 times, most recently from 08ece10 to 9dd13fbCompareMarch 13, 2026 14:29
@sam-dumont
sam-dumont marked this pull request as ready for review March 13, 2026 14:43
@sam-dumont

Copy link
Copy Markdown
ContributorAuthor

@potiuk Hi ! fixed the issues following your advice, thanks a lot ! the 3 failures here seems to be transient, not sure if I can relaunch the job myself, does not seems like it. Hope it's going to be okay now :)

In HA deployments, ti_skip_downstream() issues a bulk UPDATE without
a state guard. When a BranchOperator decides to skip downstream tasks,
it can overwrite a task already RUNNING on a worker to SKIPPED, causing
a 409 heartbeat conflict that kills the task mid-execution.
Add a skippable_state_clause to the UPDATE WHERE clause so RUNNING,
SUCCESS, and FAILED tasks are never overwritten to SKIPPED.
QUEUED tasks are intentionally allowed to be skipped: no work has been
done yet and the BranchOperator's decision should take priority. The
worker pod will get a benign 409 on PATCH /run and exit cleanly.
related: apache#59378
@kaxil
kaxilforce-pushed the fix/ti-skip-downstream-state-guard branch from 9dd13fb to 4d57ba4CompareMarch 19, 2026 01:48
@kaxilkaxil added this to the Airflow 3.2.0 milestone Mar 19, 2026
@kaxil
kaxil merged commit 6d1794a into apache:mainMar 19, 2026
130 of 132 checks passed
fat-catTW pushed a commit to fat-catTW/airflow that referenced this pull request Mar 22, 2026
…pache#63266)
In HA deployments, ti_skip_downstream() issues a bulk UPDATE without
a state guard. When a BranchOperator decides to skip downstream tasks,
it can overwrite a task already RUNNING on a worker to SKIPPED, causing
a 409 heartbeat conflict that kills the task mid-execution.
Add a skippable_state_clause to the UPDATE WHERE clause so RUNNING,
SUCCESS, and FAILED tasks are never overwritten to SKIPPED.
QUEUED tasks are intentionally allowed to be skipped: no work has been
done yet and the BranchOperator's decision should take priority. The
worker pod will get a benign 409 on PATCH /run and exit cleanly.
closes: apache#59378closes: apache#57618
Suraj-kumar00 pushed a commit to Suraj-kumar00/airflow that referenced this pull request Apr 7, 2026
…pache#63266)
In HA deployments, ti_skip_downstream() issues a bulk UPDATE without
a state guard. When a BranchOperator decides to skip downstream tasks,
it can overwrite a task already RUNNING on a worker to SKIPPED, causing
a 409 heartbeat conflict that kills the task mid-execution.
Add a skippable_state_clause to the UPDATE WHERE clause so RUNNING,
SUCCESS, and FAILED tasks are never overwritten to SKIPPED.
QUEUED tasks are intentionally allowed to be skipped: no work has been
done yet and the BranchOperator's decision should take priority. The
worker pod will get a benign 409 on PATCH /run and exit cleanly.
closes: apache#59378closes: apache#57618
abhijeets25012-tech pushed a commit to abhijeets25012-tech/airflow that referenced this pull request Apr 9, 2026
…pache#63266)
In HA deployments, ti_skip_downstream() issues a bulk UPDATE without
a state guard. When a BranchOperator decides to skip downstream tasks,
it can overwrite a task already RUNNING on a worker to SKIPPED, causing
a 409 heartbeat conflict that kills the task mid-execution.
Add a skippable_state_clause to the UPDATE WHERE clause so RUNNING,
SUCCESS, and FAILED tasks are never overwritten to SKIPPED.
QUEUED tasks are intentionally allowed to be skipped: no work has been
done yet and the BranchOperator's decision should take priority. The
worker pod will get a benign 409 on PATCH /run and exit cleanly.
closes: apache#59378closes: apache#57618
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:APIAirflow's REST/HTTP APIarea:task-sdk

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@sam-dumont@potiuk@kaxil
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Fix ti_skip_downstream overwriting RUNNING tasks to SKIPPED - #63266

Merged
kaxil merged 1 commit into
apache:mainfrom
sam-dumont:fix/ti-skip-downstream-state-guard
Mar 19, 2026
Merged

Fix ti_skip_downstream overwriting RUNNING tasks to SKIPPED#63266
kaxil merged 1 commit into
apache:mainfrom
sam-dumont:fix/ti-skip-downstream-state-guard

Conversation

@sam-dumont

Copy link
Copy Markdown
Contributor

ti_skip_downstream() issues an UPDATE filtered by (dag_id, run_id, task_id, map_index) without a state guard. When a BranchOperator on one scheduler decides to skip downstream tasks, the UPDATE can overwrite a task already RUNNING on a worker. The worker's next heartbeat returns 409 with current_state: skipped, killing the task mid-execution.

This is a companion fix to #60330, which guards schedule_tis() against the same class of race condition. Different code path (Execution API routes vs dagrun.py), same root cause : unguarded bulk UPDATEs on TI state.

Production data (12 days, 5 schedulers, ~500 concurrent workers)

We deployed both fixes as monkey patches on our prod cluster and monitored 409 heartbeat errors via CloudWatch :

Before any fix 14-169 errors/day
After schedule_tis 3-4/day (all current_state: skipped)
After both fixes 0 errors for 18+ hours
MetricBefore fixesAfter schedule_tis onlyAfter both fixes
Total 409s/day14-1693-40
current_state: scheduledpresent00
current_state: failed47/day00
current_state: skipped8/day2-5/day0

Fix

Add skippable_state_clause to the UPDATE's WHERE clause :

skippable_state_clause=or_(
TI.state.is_(None),
TI.state.not_in([RUNNING, SUCCESS, FAILED]),
)

The or_(IS NULL, NOT IN) pattern handles SQL NULL semantics : NULL NOT IN (...) evaluates to NULL (falsy), so tasks with state=None need an explicit IS NULL check to remain skippable.

QUEUED is intentionally NOT guarded : a QUEUED task hasn't started executing yet, so the BranchOperator's decision should take priority. The worker pod will get a benign 409 on PATCH /run and exit cleanly. Blocking QUEUED would cause a semantic error where the wrong branch executes.

Tests

5 regression tests in TestTISkipDownstreamRaceCondition :

  • RUNNING / SUCCESS / FAILED tasks protected from overwrite (parametrized)
  • QUEUED task correctly skipped (BranchOperator decision wins over queue)
  • None-state task still correctly skipped (happy path)

related: #59378

related: #60330

related: #57618


Was generative AI tooling used to co-author this PR?
  • Yes — Claude Code (claude-opus-4-6)

Generated-by: Claude Code (claude-opus-4-6) following the guidelines

@boring-cyborgboring-cyborgBot added area:API Airflow's REST/HTTP API area:task-sdk labels Mar 10, 2026
@sam-dumont
sam-dumontforce-pushed the fix/ti-skip-downstream-state-guard branch from f4a3ded to 589633aCompareMarch 10, 2026 14:44
@potiuk

potiuk commented Mar 11, 2026

Copy link
Copy Markdown
Member

@sam-dumont This PR has been converted to draft because it does not yet meet our Pull Request quality criteria.

Issues found:

  • Pre-commit / static checks: Failing: CI image checks / Static checks. Run prek run --from-ref main locally to find and fix issues. See Pre-commit / static checks docs.
  • mypy (type checking): Failing: CI image checks / MyPy checks (mypy-airflow-core), CI image checks / MyPy checks (mypy-task-sdk). Run prek --stage manual mypy-airflow-core --all-files && prek --stage manual mypy-task-sdk --all-files locally to reproduce. You need breeze ci-image build --python 3.10 for Docker-based mypy. See mypy (type checking) docs.
  • Provider tests: Failing: provider distributions tests / Compat 2.11.1:P3.10:, Postgres tests: providers / DB-prov:Postgres:14:3.10:-amazon,celer...standard, MySQL tests: providers / DB-prov:MySQL:8.0:3.10:-amazon,celer...standard, Sqlite tests: providers / DB-prov:Sqlite:3.10:-amazon,celer...standard, Non-DB tests: providers / Non-DB-prov::3.10:-amazon,celer...standard (+7 more). Run provider tests with breeze run pytest <provider-test-path> -xvs. See Provider tests docs.
  • Other failing CI checks: Failing: CI image checks / Test Python API client, Postgres tests: core / DB-core:Postgres:14:3.10:API...Serialization, MySQL tests: core / DB-core:MySQL:8.0:3.10:API...Serialization, Sqlite tests: core / DB-core:Sqlite:3.10:API...Serialization, Non-DB tests: core / Non-DB-core::3.10:API...Serialization (+8 more). Run prek run --from-ref main locally to reproduce. See static checks docs.

Note: Your branch is 45 commits behind main. Some check failures may be caused by changes in the base branch rather than by your PR. Please rebase your branch and push again to get up-to-date CI results.

What to do next:

  • The comment informs you what you need to do.
  • Fix each issue, then mark the PR as "Ready for review" in the GitHub UI - but only after making sure that all the issues are fixed.
  • Maintainers will then proceed with a normal review.

Converting a PR to draft is not a rejection — it is an invitation to bring the PR up to the project's standards so that maintainer review time is spent productively. If you have questions, feel free to ask on the Airflow Slack.


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

@potiuk
potiuk marked this pull request as draft March 11, 2026 00:01
@sam-dumont
sam-dumontforce-pushed the fix/ti-skip-downstream-state-guard branch 4 times, most recently from 08ece10 to 9dd13fbCompareMarch 13, 2026 14:29
@sam-dumont
sam-dumont marked this pull request as ready for review March 13, 2026 14:43
@sam-dumont

Copy link
Copy Markdown
ContributorAuthor

@potiuk Hi ! fixed the issues following your advice, thanks a lot ! the 3 failures here seems to be transient, not sure if I can relaunch the job myself, does not seems like it. Hope it's going to be okay now :)

In HA deployments, ti_skip_downstream() issues a bulk UPDATE without
a state guard. When a BranchOperator decides to skip downstream tasks,
it can overwrite a task already RUNNING on a worker to SKIPPED, causing
a 409 heartbeat conflict that kills the task mid-execution.
Add a skippable_state_clause to the UPDATE WHERE clause so RUNNING,
SUCCESS, and FAILED tasks are never overwritten to SKIPPED.
QUEUED tasks are intentionally allowed to be skipped: no work has been
done yet and the BranchOperator's decision should take priority. The
worker pod will get a benign 409 on PATCH /run and exit cleanly.
related: apache#59378
@kaxil
kaxilforce-pushed the fix/ti-skip-downstream-state-guard branch from 9dd13fb to 4d57ba4CompareMarch 19, 2026 01:48
@kaxilkaxil added this to the Airflow 3.2.0 milestone Mar 19, 2026
@kaxil
kaxil merged commit 6d1794a into apache:mainMar 19, 2026
130 of 132 checks passed
fat-catTW pushed a commit to fat-catTW/airflow that referenced this pull request Mar 22, 2026
…pache#63266)
In HA deployments, ti_skip_downstream() issues a bulk UPDATE without
a state guard. When a BranchOperator decides to skip downstream tasks,
it can overwrite a task already RUNNING on a worker to SKIPPED, causing
a 409 heartbeat conflict that kills the task mid-execution.
Add a skippable_state_clause to the UPDATE WHERE clause so RUNNING,
SUCCESS, and FAILED tasks are never overwritten to SKIPPED.
QUEUED tasks are intentionally allowed to be skipped: no work has been
done yet and the BranchOperator's decision should take priority. The
worker pod will get a benign 409 on PATCH /run and exit cleanly.
closes: apache#59378closes: apache#57618
Suraj-kumar00 pushed a commit to Suraj-kumar00/airflow that referenced this pull request Apr 7, 2026
…pache#63266)
In HA deployments, ti_skip_downstream() issues a bulk UPDATE without
a state guard. When a BranchOperator decides to skip downstream tasks,
it can overwrite a task already RUNNING on a worker to SKIPPED, causing
a 409 heartbeat conflict that kills the task mid-execution.
Add a skippable_state_clause to the UPDATE WHERE clause so RUNNING,
SUCCESS, and FAILED tasks are never overwritten to SKIPPED.
QUEUED tasks are intentionally allowed to be skipped: no work has been
done yet and the BranchOperator's decision should take priority. The
worker pod will get a benign 409 on PATCH /run and exit cleanly.
closes: apache#59378closes: apache#57618
abhijeets25012-tech pushed a commit to abhijeets25012-tech/airflow that referenced this pull request Apr 9, 2026
…pache#63266)
In HA deployments, ti_skip_downstream() issues a bulk UPDATE without
a state guard. When a BranchOperator decides to skip downstream tasks,
it can overwrite a task already RUNNING on a worker to SKIPPED, causing
a 409 heartbeat conflict that kills the task mid-execution.
Add a skippable_state_clause to the UPDATE WHERE clause so RUNNING,
SUCCESS, and FAILED tasks are never overwritten to SKIPPED.
QUEUED tasks are intentionally allowed to be skipped: no work has been
done yet and the BranchOperator's decision should take priority. The
worker pod will get a benign 409 on PATCH /run and exit cleanly.
closes: apache#59378closes: apache#57618
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:APIAirflow's REST/HTTP APIarea:task-sdk

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

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

Fix ti_skip_downstream overwriting RUNNING tasks to SKIPPED - #63266

Merged
kaxil merged 1 commit into
apache:mainfrom
sam-dumont:fix/ti-skip-downstream-state-guard
Mar 19, 2026
Merged

Fix ti_skip_downstream overwriting RUNNING tasks to SKIPPED#63266
kaxil merged 1 commit into
apache:mainfrom
sam-dumont:fix/ti-skip-downstream-state-guard

Conversation

@sam-dumont

Copy link
Copy Markdown
Contributor

ti_skip_downstream() issues an UPDATE filtered by (dag_id, run_id, task_id, map_index) without a state guard. When a BranchOperator on one scheduler decides to skip downstream tasks, the UPDATE can overwrite a task already RUNNING on a worker. The worker's next heartbeat returns 409 with current_state: skipped, killing the task mid-execution.

This is a companion fix to #60330, which guards schedule_tis() against the same class of race condition. Different code path (Execution API routes vs dagrun.py), same root cause : unguarded bulk UPDATEs on TI state.

Production data (12 days, 5 schedulers, ~500 concurrent workers)

We deployed both fixes as monkey patches on our prod cluster and monitored 409 heartbeat errors via CloudWatch :

Before any fix 14-169 errors/day
After schedule_tis 3-4/day (all current_state: skipped)
After both fixes 0 errors for 18+ hours
MetricBefore fixesAfter schedule_tis onlyAfter both fixes
Total 409s/day14-1693-40
current_state: scheduledpresent00
current_state: failed47/day00
current_state: skipped8/day2-5/day0

Fix

Add skippable_state_clause to the UPDATE's WHERE clause :

skippable_state_clause=or_(
TI.state.is_(None),
TI.state.not_in([RUNNING, SUCCESS, FAILED]),
)

The or_(IS NULL, NOT IN) pattern handles SQL NULL semantics : NULL NOT IN (...) evaluates to NULL (falsy), so tasks with state=None need an explicit IS NULL check to remain skippable.

QUEUED is intentionally NOT guarded : a QUEUED task hasn't started executing yet, so the BranchOperator's decision should take priority. The worker pod will get a benign 409 on PATCH /run and exit cleanly. Blocking QUEUED would cause a semantic error where the wrong branch executes.

Tests

5 regression tests in TestTISkipDownstreamRaceCondition :

  • RUNNING / SUCCESS / FAILED tasks protected from overwrite (parametrized)
  • QUEUED task correctly skipped (BranchOperator decision wins over queue)
  • None-state task still correctly skipped (happy path)

related: #59378

related: #60330

related: #57618


Was generative AI tooling used to co-author this PR?
  • Yes — Claude Code (claude-opus-4-6)

Generated-by: Claude Code (claude-opus-4-6) following the guidelines

@boring-cyborgboring-cyborgBot added area:API Airflow's REST/HTTP API area:task-sdk labels Mar 10, 2026
@sam-dumont
sam-dumontforce-pushed the fix/ti-skip-downstream-state-guard branch from f4a3ded to 589633aCompareMarch 10, 2026 14:44
@potiuk

potiuk commented Mar 11, 2026

Copy link
Copy Markdown
Member

@sam-dumont This PR has been converted to draft because it does not yet meet our Pull Request quality criteria.

Issues found:

  • Pre-commit / static checks: Failing: CI image checks / Static checks. Run prek run --from-ref main locally to find and fix issues. See Pre-commit / static checks docs.
  • mypy (type checking): Failing: CI image checks / MyPy checks (mypy-airflow-core), CI image checks / MyPy checks (mypy-task-sdk). Run prek --stage manual mypy-airflow-core --all-files && prek --stage manual mypy-task-sdk --all-files locally to reproduce. You need breeze ci-image build --python 3.10 for Docker-based mypy. See mypy (type checking) docs.
  • Provider tests: Failing: provider distributions tests / Compat 2.11.1:P3.10:, Postgres tests: providers / DB-prov:Postgres:14:3.10:-amazon,celer...standard, MySQL tests: providers / DB-prov:MySQL:8.0:3.10:-amazon,celer...standard, Sqlite tests: providers / DB-prov:Sqlite:3.10:-amazon,celer...standard, Non-DB tests: providers / Non-DB-prov::3.10:-amazon,celer...standard (+7 more). Run provider tests with breeze run pytest <provider-test-path> -xvs. See Provider tests docs.
  • Other failing CI checks: Failing: CI image checks / Test Python API client, Postgres tests: core / DB-core:Postgres:14:3.10:API...Serialization, MySQL tests: core / DB-core:MySQL:8.0:3.10:API...Serialization, Sqlite tests: core / DB-core:Sqlite:3.10:API...Serialization, Non-DB tests: core / Non-DB-core::3.10:API...Serialization (+8 more). Run prek run --from-ref main locally to reproduce. See static checks docs.

Note: Your branch is 45 commits behind main. Some check failures may be caused by changes in the base branch rather than by your PR. Please rebase your branch and push again to get up-to-date CI results.

What to do next:

  • The comment informs you what you need to do.
  • Fix each issue, then mark the PR as "Ready for review" in the GitHub UI - but only after making sure that all the issues are fixed.
  • Maintainers will then proceed with a normal review.

Converting a PR to draft is not a rejection — it is an invitation to bring the PR up to the project's standards so that maintainer review time is spent productively. If you have questions, feel free to ask on the Airflow Slack.


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

@potiuk
potiuk marked this pull request as draft March 11, 2026 00:01
@sam-dumont
sam-dumontforce-pushed the fix/ti-skip-downstream-state-guard branch 4 times, most recently from 08ece10 to 9dd13fbCompareMarch 13, 2026 14:29
@sam-dumont
sam-dumont marked this pull request as ready for review March 13, 2026 14:43
@sam-dumont

Copy link
Copy Markdown
ContributorAuthor

@potiuk Hi ! fixed the issues following your advice, thanks a lot ! the 3 failures here seems to be transient, not sure if I can relaunch the job myself, does not seems like it. Hope it's going to be okay now :)

In HA deployments, ti_skip_downstream() issues a bulk UPDATE without
a state guard. When a BranchOperator decides to skip downstream tasks,
it can overwrite a task already RUNNING on a worker to SKIPPED, causing
a 409 heartbeat conflict that kills the task mid-execution.
Add a skippable_state_clause to the UPDATE WHERE clause so RUNNING,
SUCCESS, and FAILED tasks are never overwritten to SKIPPED.
QUEUED tasks are intentionally allowed to be skipped: no work has been
done yet and the BranchOperator's decision should take priority. The
worker pod will get a benign 409 on PATCH /run and exit cleanly.
related: apache#59378
@kaxil
kaxilforce-pushed the fix/ti-skip-downstream-state-guard branch from 9dd13fb to 4d57ba4CompareMarch 19, 2026 01:48
@kaxilkaxil added this to the Airflow 3.2.0 milestone Mar 19, 2026
@kaxil
kaxil merged commit 6d1794a into apache:mainMar 19, 2026
130 of 132 checks passed
fat-catTW pushed a commit to fat-catTW/airflow that referenced this pull request Mar 22, 2026
…pache#63266)
In HA deployments, ti_skip_downstream() issues a bulk UPDATE without
a state guard. When a BranchOperator decides to skip downstream tasks,
it can overwrite a task already RUNNING on a worker to SKIPPED, causing
a 409 heartbeat conflict that kills the task mid-execution.
Add a skippable_state_clause to the UPDATE WHERE clause so RUNNING,
SUCCESS, and FAILED tasks are never overwritten to SKIPPED.
QUEUED tasks are intentionally allowed to be skipped: no work has been
done yet and the BranchOperator's decision should take priority. The
worker pod will get a benign 409 on PATCH /run and exit cleanly.
closes: apache#59378closes: apache#57618
Suraj-kumar00 pushed a commit to Suraj-kumar00/airflow that referenced this pull request Apr 7, 2026
…pache#63266)
In HA deployments, ti_skip_downstream() issues a bulk UPDATE without
a state guard. When a BranchOperator decides to skip downstream tasks,
it can overwrite a task already RUNNING on a worker to SKIPPED, causing
a 409 heartbeat conflict that kills the task mid-execution.
Add a skippable_state_clause to the UPDATE WHERE clause so RUNNING,
SUCCESS, and FAILED tasks are never overwritten to SKIPPED.
QUEUED tasks are intentionally allowed to be skipped: no work has been
done yet and the BranchOperator's decision should take priority. The
worker pod will get a benign 409 on PATCH /run and exit cleanly.
closes: apache#59378closes: apache#57618
abhijeets25012-tech pushed a commit to abhijeets25012-tech/airflow that referenced this pull request Apr 9, 2026
…pache#63266)
In HA deployments, ti_skip_downstream() issues a bulk UPDATE without
a state guard. When a BranchOperator decides to skip downstream tasks,
it can overwrite a task already RUNNING on a worker to SKIPPED, causing
a 409 heartbeat conflict that kills the task mid-execution.
Add a skippable_state_clause to the UPDATE WHERE clause so RUNNING,
SUCCESS, and FAILED tasks are never overwritten to SKIPPED.
QUEUED tasks are intentionally allowed to be skipped: no work has been
done yet and the BranchOperator's decision should take priority. The
worker pod will get a benign 409 on PATCH /run and exit cleanly.
closes: apache#59378closes: apache#57618
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:APIAirflow's REST/HTTP APIarea:task-sdk

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@sam-dumont@potiuk@kaxil
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content

Fix ti_skip_downstream overwriting RUNNING tasks to SKIPPED - #63266

Merged
kaxil merged 1 commit into
apache:mainfrom
sam-dumont:fix/ti-skip-downstream-state-guard
Mar 19, 2026
Merged

Fix ti_skip_downstream overwriting RUNNING tasks to SKIPPED#63266
kaxil merged 1 commit into
apache:mainfrom
sam-dumont:fix/ti-skip-downstream-state-guard

Conversation

@sam-dumont

Copy link
Copy Markdown
Contributor

ti_skip_downstream() issues an UPDATE filtered by (dag_id, run_id, task_id, map_index) without a state guard. When a BranchOperator on one scheduler decides to skip downstream tasks, the UPDATE can overwrite a task already RUNNING on a worker. The worker's next heartbeat returns 409 with current_state: skipped, killing the task mid-execution.

This is a companion fix to #60330, which guards schedule_tis() against the same class of race condition. Different code path (Execution API routes vs dagrun.py), same root cause : unguarded bulk UPDATEs on TI state.

Production data (12 days, 5 schedulers, ~500 concurrent workers)

We deployed both fixes as monkey patches on our prod cluster and monitored 409 heartbeat errors via CloudWatch :

Before any fix 14-169 errors/day
After schedule_tis 3-4/day (all current_state: skipped)
After both fixes 0 errors for 18+ hours
MetricBefore fixesAfter schedule_tis onlyAfter both fixes
Total 409s/day14-1693-40
current_state: scheduledpresent00
current_state: failed47/day00
current_state: skipped8/day2-5/day0

Fix

Add skippable_state_clause to the UPDATE's WHERE clause :

skippable_state_clause=or_(
TI.state.is_(None),
TI.state.not_in([RUNNING, SUCCESS, FAILED]),
)

The or_(IS NULL, NOT IN) pattern handles SQL NULL semantics : NULL NOT IN (...) evaluates to NULL (falsy), so tasks with state=None need an explicit IS NULL check to remain skippable.

QUEUED is intentionally NOT guarded : a QUEUED task hasn't started executing yet, so the BranchOperator's decision should take priority. The worker pod will get a benign 409 on PATCH /run and exit cleanly. Blocking QUEUED would cause a semantic error where the wrong branch executes.

Tests

5 regression tests in TestTISkipDownstreamRaceCondition :

  • RUNNING / SUCCESS / FAILED tasks protected from overwrite (parametrized)
  • QUEUED task correctly skipped (BranchOperator decision wins over queue)
  • None-state task still correctly skipped (happy path)

related: #59378

related: #60330

related: #57618


Was generative AI tooling used to co-author this PR?
  • Yes — Claude Code (claude-opus-4-6)

Generated-by: Claude Code (claude-opus-4-6) following the guidelines

@boring-cyborgboring-cyborgBot added area:API Airflow's REST/HTTP API area:task-sdk labels Mar 10, 2026
@sam-dumont
sam-dumontforce-pushed the fix/ti-skip-downstream-state-guard branch from f4a3ded to 589633aCompareMarch 10, 2026 14:44
@potiuk

potiuk commented Mar 11, 2026

Copy link
Copy Markdown
Member

@sam-dumont This PR has been converted to draft because it does not yet meet our Pull Request quality criteria.

Issues found:

  • Pre-commit / static checks: Failing: CI image checks / Static checks. Run prek run --from-ref main locally to find and fix issues. See Pre-commit / static checks docs.
  • mypy (type checking): Failing: CI image checks / MyPy checks (mypy-airflow-core), CI image checks / MyPy checks (mypy-task-sdk). Run prek --stage manual mypy-airflow-core --all-files && prek --stage manual mypy-task-sdk --all-files locally to reproduce. You need breeze ci-image build --python 3.10 for Docker-based mypy. See mypy (type checking) docs.
  • Provider tests: Failing: provider distributions tests / Compat 2.11.1:P3.10:, Postgres tests: providers / DB-prov:Postgres:14:3.10:-amazon,celer...standard, MySQL tests: providers / DB-prov:MySQL:8.0:3.10:-amazon,celer...standard, Sqlite tests: providers / DB-prov:Sqlite:3.10:-amazon,celer...standard, Non-DB tests: providers / Non-DB-prov::3.10:-amazon,celer...standard (+7 more). Run provider tests with breeze run pytest <provider-test-path> -xvs. See Provider tests docs.
  • Other failing CI checks: Failing: CI image checks / Test Python API client, Postgres tests: core / DB-core:Postgres:14:3.10:API...Serialization, MySQL tests: core / DB-core:MySQL:8.0:3.10:API...Serialization, Sqlite tests: core / DB-core:Sqlite:3.10:API...Serialization, Non-DB tests: core / Non-DB-core::3.10:API...Serialization (+8 more). Run prek run --from-ref main locally to reproduce. See static checks docs.

Note: Your branch is 45 commits behind main. Some check failures may be caused by changes in the base branch rather than by your PR. Please rebase your branch and push again to get up-to-date CI results.

What to do next:

  • The comment informs you what you need to do.
  • Fix each issue, then mark the PR as "Ready for review" in the GitHub UI - but only after making sure that all the issues are fixed.
  • Maintainers will then proceed with a normal review.

Converting a PR to draft is not a rejection — it is an invitation to bring the PR up to the project's standards so that maintainer review time is spent productively. If you have questions, feel free to ask on the Airflow Slack.


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

@potiuk
potiuk marked this pull request as draft March 11, 2026 00:01
@sam-dumont
sam-dumontforce-pushed the fix/ti-skip-downstream-state-guard branch 4 times, most recently from 08ece10 to 9dd13fbCompareMarch 13, 2026 14:29
@sam-dumont
sam-dumont marked this pull request as ready for review March 13, 2026 14:43
@sam-dumont

Copy link
Copy Markdown
ContributorAuthor

@potiuk Hi ! fixed the issues following your advice, thanks a lot ! the 3 failures here seems to be transient, not sure if I can relaunch the job myself, does not seems like it. Hope it's going to be okay now :)

In HA deployments, ti_skip_downstream() issues a bulk UPDATE without
a state guard. When a BranchOperator decides to skip downstream tasks,
it can overwrite a task already RUNNING on a worker to SKIPPED, causing
a 409 heartbeat conflict that kills the task mid-execution.
Add a skippable_state_clause to the UPDATE WHERE clause so RUNNING,
SUCCESS, and FAILED tasks are never overwritten to SKIPPED.
QUEUED tasks are intentionally allowed to be skipped: no work has been
done yet and the BranchOperator's decision should take priority. The
worker pod will get a benign 409 on PATCH /run and exit cleanly.
related: apache#59378
@kaxil
kaxilforce-pushed the fix/ti-skip-downstream-state-guard branch from 9dd13fb to 4d57ba4CompareMarch 19, 2026 01:48
@kaxilkaxil added this to the Airflow 3.2.0 milestone Mar 19, 2026
@kaxil
kaxil merged commit 6d1794a into apache:mainMar 19, 2026
130 of 132 checks passed
fat-catTW pushed a commit to fat-catTW/airflow that referenced this pull request Mar 22, 2026
…pache#63266)
In HA deployments, ti_skip_downstream() issues a bulk UPDATE without
a state guard. When a BranchOperator decides to skip downstream tasks,
it can overwrite a task already RUNNING on a worker to SKIPPED, causing
a 409 heartbeat conflict that kills the task mid-execution.
Add a skippable_state_clause to the UPDATE WHERE clause so RUNNING,
SUCCESS, and FAILED tasks are never overwritten to SKIPPED.
QUEUED tasks are intentionally allowed to be skipped: no work has been
done yet and the BranchOperator's decision should take priority. The
worker pod will get a benign 409 on PATCH /run and exit cleanly.
closes: apache#59378closes: apache#57618
Suraj-kumar00 pushed a commit to Suraj-kumar00/airflow that referenced this pull request Apr 7, 2026
…pache#63266)
In HA deployments, ti_skip_downstream() issues a bulk UPDATE without
a state guard. When a BranchOperator decides to skip downstream tasks,
it can overwrite a task already RUNNING on a worker to SKIPPED, causing
a 409 heartbeat conflict that kills the task mid-execution.
Add a skippable_state_clause to the UPDATE WHERE clause so RUNNING,
SUCCESS, and FAILED tasks are never overwritten to SKIPPED.
QUEUED tasks are intentionally allowed to be skipped: no work has been
done yet and the BranchOperator's decision should take priority. The
worker pod will get a benign 409 on PATCH /run and exit cleanly.
closes: apache#59378closes: apache#57618
abhijeets25012-tech pushed a commit to abhijeets25012-tech/airflow that referenced this pull request Apr 9, 2026
…pache#63266)
In HA deployments, ti_skip_downstream() issues a bulk UPDATE without
a state guard. When a BranchOperator decides to skip downstream tasks,
it can overwrite a task already RUNNING on a worker to SKIPPED, causing
a 409 heartbeat conflict that kills the task mid-execution.
Add a skippable_state_clause to the UPDATE WHERE clause so RUNNING,
SUCCESS, and FAILED tasks are never overwritten to SKIPPED.
QUEUED tasks are intentionally allowed to be skipped: no work has been
done yet and the BranchOperator's decision should take priority. The
worker pod will get a benign 409 on PATCH /run and exit cleanly.
closes: apache#59378closes: apache#57618
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:APIAirflow's REST/HTTP APIarea:task-sdk

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@sam-dumont@potiuk@kaxil
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Fix ti_skip_downstream overwriting RUNNING tasks to SKIPPED - #63266

Merged
kaxil merged 1 commit into
apache:mainfrom
sam-dumont:fix/ti-skip-downstream-state-guard
Mar 19, 2026
Merged

Fix ti_skip_downstream overwriting RUNNING tasks to SKIPPED#63266
kaxil merged 1 commit into
apache:mainfrom
sam-dumont:fix/ti-skip-downstream-state-guard

Conversation

@sam-dumont

Copy link
Copy Markdown
Contributor

ti_skip_downstream() issues an UPDATE filtered by (dag_id, run_id, task_id, map_index) without a state guard. When a BranchOperator on one scheduler decides to skip downstream tasks, the UPDATE can overwrite a task already RUNNING on a worker. The worker's next heartbeat returns 409 with current_state: skipped, killing the task mid-execution.

This is a companion fix to #60330, which guards schedule_tis() against the same class of race condition. Different code path (Execution API routes vs dagrun.py), same root cause : unguarded bulk UPDATEs on TI state.

Production data (12 days, 5 schedulers, ~500 concurrent workers)

We deployed both fixes as monkey patches on our prod cluster and monitored 409 heartbeat errors via CloudWatch :

Before any fix 14-169 errors/day
After schedule_tis 3-4/day (all current_state: skipped)
After both fixes 0 errors for 18+ hours
MetricBefore fixesAfter schedule_tis onlyAfter both fixes
Total 409s/day14-1693-40
current_state: scheduledpresent00
current_state: failed47/day00
current_state: skipped8/day2-5/day0

Fix

Add skippable_state_clause to the UPDATE's WHERE clause :

skippable_state_clause=or_(
TI.state.is_(None),
TI.state.not_in([RUNNING, SUCCESS, FAILED]),
)

The or_(IS NULL, NOT IN) pattern handles SQL NULL semantics : NULL NOT IN (...) evaluates to NULL (falsy), so tasks with state=None need an explicit IS NULL check to remain skippable.

QUEUED is intentionally NOT guarded : a QUEUED task hasn't started executing yet, so the BranchOperator's decision should take priority. The worker pod will get a benign 409 on PATCH /run and exit cleanly. Blocking QUEUED would cause a semantic error where the wrong branch executes.

Tests

5 regression tests in TestTISkipDownstreamRaceCondition :

  • RUNNING / SUCCESS / FAILED tasks protected from overwrite (parametrized)
  • QUEUED task correctly skipped (BranchOperator decision wins over queue)
  • None-state task still correctly skipped (happy path)

related: #59378

related: #60330

related: #57618


Was generative AI tooling used to co-author this PR?
  • Yes — Claude Code (claude-opus-4-6)

Generated-by: Claude Code (claude-opus-4-6) following the guidelines

@boring-cyborgboring-cyborgBot added area:API Airflow's REST/HTTP API area:task-sdk labels Mar 10, 2026
@sam-dumont
sam-dumontforce-pushed the fix/ti-skip-downstream-state-guard branch from f4a3ded to 589633aCompareMarch 10, 2026 14:44
@potiuk

potiuk commented Mar 11, 2026

Copy link
Copy Markdown
Member

@sam-dumont This PR has been converted to draft because it does not yet meet our Pull Request quality criteria.

Issues found:

  • Pre-commit / static checks: Failing: CI image checks / Static checks. Run prek run --from-ref main locally to find and fix issues. See Pre-commit / static checks docs.
  • mypy (type checking): Failing: CI image checks / MyPy checks (mypy-airflow-core), CI image checks / MyPy checks (mypy-task-sdk). Run prek --stage manual mypy-airflow-core --all-files && prek --stage manual mypy-task-sdk --all-files locally to reproduce. You need breeze ci-image build --python 3.10 for Docker-based mypy. See mypy (type checking) docs.
  • Provider tests: Failing: provider distributions tests / Compat 2.11.1:P3.10:, Postgres tests: providers / DB-prov:Postgres:14:3.10:-amazon,celer...standard, MySQL tests: providers / DB-prov:MySQL:8.0:3.10:-amazon,celer...standard, Sqlite tests: providers / DB-prov:Sqlite:3.10:-amazon,celer...standard, Non-DB tests: providers / Non-DB-prov::3.10:-amazon,celer...standard (+7 more). Run provider tests with breeze run pytest <provider-test-path> -xvs. See Provider tests docs.
  • Other failing CI checks: Failing: CI image checks / Test Python API client, Postgres tests: core / DB-core:Postgres:14:3.10:API...Serialization, MySQL tests: core / DB-core:MySQL:8.0:3.10:API...Serialization, Sqlite tests: core / DB-core:Sqlite:3.10:API...Serialization, Non-DB tests: core / Non-DB-core::3.10:API...Serialization (+8 more). Run prek run --from-ref main locally to reproduce. See static checks docs.

Note: Your branch is 45 commits behind main. Some check failures may be caused by changes in the base branch rather than by your PR. Please rebase your branch and push again to get up-to-date CI results.

What to do next:

  • The comment informs you what you need to do.
  • Fix each issue, then mark the PR as "Ready for review" in the GitHub UI - but only after making sure that all the issues are fixed.
  • Maintainers will then proceed with a normal review.

Converting a PR to draft is not a rejection — it is an invitation to bring the PR up to the project's standards so that maintainer review time is spent productively. If you have questions, feel free to ask on the Airflow Slack.


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

@potiuk
potiuk marked this pull request as draft March 11, 2026 00:01
@sam-dumont
sam-dumontforce-pushed the fix/ti-skip-downstream-state-guard branch 4 times, most recently from 08ece10 to 9dd13fbCompareMarch 13, 2026 14:29
@sam-dumont
sam-dumont marked this pull request as ready for review March 13, 2026 14:43
@sam-dumont

Copy link
Copy Markdown
ContributorAuthor

@potiuk Hi ! fixed the issues following your advice, thanks a lot ! the 3 failures here seems to be transient, not sure if I can relaunch the job myself, does not seems like it. Hope it's going to be okay now :)

In HA deployments, ti_skip_downstream() issues a bulk UPDATE without
a state guard. When a BranchOperator decides to skip downstream tasks,
it can overwrite a task already RUNNING on a worker to SKIPPED, causing
a 409 heartbeat conflict that kills the task mid-execution.
Add a skippable_state_clause to the UPDATE WHERE clause so RUNNING,
SUCCESS, and FAILED tasks are never overwritten to SKIPPED.
QUEUED tasks are intentionally allowed to be skipped: no work has been
done yet and the BranchOperator's decision should take priority. The
worker pod will get a benign 409 on PATCH /run and exit cleanly.
related: apache#59378
@kaxil
kaxilforce-pushed the fix/ti-skip-downstream-state-guard branch from 9dd13fb to 4d57ba4CompareMarch 19, 2026 01:48
@kaxilkaxil added this to the Airflow 3.2.0 milestone Mar 19, 2026
@kaxil
kaxil merged commit 6d1794a into apache:mainMar 19, 2026
130 of 132 checks passed
fat-catTW pushed a commit to fat-catTW/airflow that referenced this pull request Mar 22, 2026
…pache#63266)
In HA deployments, ti_skip_downstream() issues a bulk UPDATE without
a state guard. When a BranchOperator decides to skip downstream tasks,
it can overwrite a task already RUNNING on a worker to SKIPPED, causing
a 409 heartbeat conflict that kills the task mid-execution.
Add a skippable_state_clause to the UPDATE WHERE clause so RUNNING,
SUCCESS, and FAILED tasks are never overwritten to SKIPPED.
QUEUED tasks are intentionally allowed to be skipped: no work has been
done yet and the BranchOperator's decision should take priority. The
worker pod will get a benign 409 on PATCH /run and exit cleanly.
closes: apache#59378closes: apache#57618
Suraj-kumar00 pushed a commit to Suraj-kumar00/airflow that referenced this pull request Apr 7, 2026
…pache#63266)
In HA deployments, ti_skip_downstream() issues a bulk UPDATE without
a state guard. When a BranchOperator decides to skip downstream tasks,
it can overwrite a task already RUNNING on a worker to SKIPPED, causing
a 409 heartbeat conflict that kills the task mid-execution.
Add a skippable_state_clause to the UPDATE WHERE clause so RUNNING,
SUCCESS, and FAILED tasks are never overwritten to SKIPPED.
QUEUED tasks are intentionally allowed to be skipped: no work has been
done yet and the BranchOperator's decision should take priority. The
worker pod will get a benign 409 on PATCH /run and exit cleanly.
closes: apache#59378closes: apache#57618
abhijeets25012-tech pushed a commit to abhijeets25012-tech/airflow that referenced this pull request Apr 9, 2026
…pache#63266)
In HA deployments, ti_skip_downstream() issues a bulk UPDATE without
a state guard. When a BranchOperator decides to skip downstream tasks,
it can overwrite a task already RUNNING on a worker to SKIPPED, causing
a 409 heartbeat conflict that kills the task mid-execution.
Add a skippable_state_clause to the UPDATE WHERE clause so RUNNING,
SUCCESS, and FAILED tasks are never overwritten to SKIPPED.
QUEUED tasks are intentionally allowed to be skipped: no work has been
done yet and the BranchOperator's decision should take priority. The
worker pod will get a benign 409 on PATCH /run and exit cleanly.
closes: apache#59378closes: apache#57618
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:APIAirflow's REST/HTTP APIarea:task-sdk

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@sam-dumont@potiuk@kaxil
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Fix ti_skip_downstream overwriting RUNNING tasks to SKIPPED - #63266

Merged
kaxil merged 1 commit into
apache:mainfrom
sam-dumont:fix/ti-skip-downstream-state-guard
Mar 19, 2026
Merged

Fix ti_skip_downstream overwriting RUNNING tasks to SKIPPED#63266
kaxil merged 1 commit into
apache:mainfrom
sam-dumont:fix/ti-skip-downstream-state-guard

Conversation

@sam-dumont

Copy link
Copy Markdown
Contributor

ti_skip_downstream() issues an UPDATE filtered by (dag_id, run_id, task_id, map_index) without a state guard. When a BranchOperator on one scheduler decides to skip downstream tasks, the UPDATE can overwrite a task already RUNNING on a worker. The worker's next heartbeat returns 409 with current_state: skipped, killing the task mid-execution.

This is a companion fix to #60330, which guards schedule_tis() against the same class of race condition. Different code path (Execution API routes vs dagrun.py), same root cause : unguarded bulk UPDATEs on TI state.

Production data (12 days, 5 schedulers, ~500 concurrent workers)

We deployed both fixes as monkey patches on our prod cluster and monitored 409 heartbeat errors via CloudWatch :

Before any fix 14-169 errors/day
After schedule_tis 3-4/day (all current_state: skipped)
After both fixes 0 errors for 18+ hours
MetricBefore fixesAfter schedule_tis onlyAfter both fixes
Total 409s/day14-1693-40
current_state: scheduledpresent00
current_state: failed47/day00
current_state: skipped8/day2-5/day0

Fix

Add skippable_state_clause to the UPDATE's WHERE clause :

skippable_state_clause=or_(
TI.state.is_(None),
TI.state.not_in([RUNNING, SUCCESS, FAILED]),
)

The or_(IS NULL, NOT IN) pattern handles SQL NULL semantics : NULL NOT IN (...) evaluates to NULL (falsy), so tasks with state=None need an explicit IS NULL check to remain skippable.

QUEUED is intentionally NOT guarded : a QUEUED task hasn't started executing yet, so the BranchOperator's decision should take priority. The worker pod will get a benign 409 on PATCH /run and exit cleanly. Blocking QUEUED would cause a semantic error where the wrong branch executes.

Tests

5 regression tests in TestTISkipDownstreamRaceCondition :

  • RUNNING / SUCCESS / FAILED tasks protected from overwrite (parametrized)
  • QUEUED task correctly skipped (BranchOperator decision wins over queue)
  • None-state task still correctly skipped (happy path)

related: #59378

related: #60330

related: #57618


Was generative AI tooling used to co-author this PR?
  • Yes — Claude Code (claude-opus-4-6)

Generated-by: Claude Code (claude-opus-4-6) following the guidelines

@boring-cyborgboring-cyborgBot added area:API Airflow's REST/HTTP API area:task-sdk labels Mar 10, 2026
@sam-dumont
sam-dumontforce-pushed the fix/ti-skip-downstream-state-guard branch from f4a3ded to 589633aCompareMarch 10, 2026 14:44
@potiuk

potiuk commented Mar 11, 2026

Copy link
Copy Markdown
Member

@sam-dumont This PR has been converted to draft because it does not yet meet our Pull Request quality criteria.

Issues found:

  • Pre-commit / static checks: Failing: CI image checks / Static checks. Run prek run --from-ref main locally to find and fix issues. See Pre-commit / static checks docs.
  • mypy (type checking): Failing: CI image checks / MyPy checks (mypy-airflow-core), CI image checks / MyPy checks (mypy-task-sdk). Run prek --stage manual mypy-airflow-core --all-files && prek --stage manual mypy-task-sdk --all-files locally to reproduce. You need breeze ci-image build --python 3.10 for Docker-based mypy. See mypy (type checking) docs.
  • Provider tests: Failing: provider distributions tests / Compat 2.11.1:P3.10:, Postgres tests: providers / DB-prov:Postgres:14:3.10:-amazon,celer...standard, MySQL tests: providers / DB-prov:MySQL:8.0:3.10:-amazon,celer...standard, Sqlite tests: providers / DB-prov:Sqlite:3.10:-amazon,celer...standard, Non-DB tests: providers / Non-DB-prov::3.10:-amazon,celer...standard (+7 more). Run provider tests with breeze run pytest <provider-test-path> -xvs. See Provider tests docs.
  • Other failing CI checks: Failing: CI image checks / Test Python API client, Postgres tests: core / DB-core:Postgres:14:3.10:API...Serialization, MySQL tests: core / DB-core:MySQL:8.0:3.10:API...Serialization, Sqlite tests: core / DB-core:Sqlite:3.10:API...Serialization, Non-DB tests: core / Non-DB-core::3.10:API...Serialization (+8 more). Run prek run --from-ref main locally to reproduce. See static checks docs.

Note: Your branch is 45 commits behind main. Some check failures may be caused by changes in the base branch rather than by your PR. Please rebase your branch and push again to get up-to-date CI results.

What to do next:

  • The comment informs you what you need to do.
  • Fix each issue, then mark the PR as "Ready for review" in the GitHub UI - but only after making sure that all the issues are fixed.
  • Maintainers will then proceed with a normal review.

Converting a PR to draft is not a rejection — it is an invitation to bring the PR up to the project's standards so that maintainer review time is spent productively. If you have questions, feel free to ask on the Airflow Slack.


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

@potiuk
potiuk marked this pull request as draft March 11, 2026 00:01
@sam-dumont
sam-dumontforce-pushed the fix/ti-skip-downstream-state-guard branch 4 times, most recently from 08ece10 to 9dd13fbCompareMarch 13, 2026 14:29
@sam-dumont
sam-dumont marked this pull request as ready for review March 13, 2026 14:43
@sam-dumont

Copy link
Copy Markdown
ContributorAuthor

@potiuk Hi ! fixed the issues following your advice, thanks a lot ! the 3 failures here seems to be transient, not sure if I can relaunch the job myself, does not seems like it. Hope it's going to be okay now :)

In HA deployments, ti_skip_downstream() issues a bulk UPDATE without
a state guard. When a BranchOperator decides to skip downstream tasks,
it can overwrite a task already RUNNING on a worker to SKIPPED, causing
a 409 heartbeat conflict that kills the task mid-execution.
Add a skippable_state_clause to the UPDATE WHERE clause so RUNNING,
SUCCESS, and FAILED tasks are never overwritten to SKIPPED.
QUEUED tasks are intentionally allowed to be skipped: no work has been
done yet and the BranchOperator's decision should take priority. The
worker pod will get a benign 409 on PATCH /run and exit cleanly.
related: apache#59378
@kaxil
kaxilforce-pushed the fix/ti-skip-downstream-state-guard branch from 9dd13fb to 4d57ba4CompareMarch 19, 2026 01:48
@kaxilkaxil added this to the Airflow 3.2.0 milestone Mar 19, 2026
@kaxil
kaxil merged commit 6d1794a into apache:mainMar 19, 2026
130 of 132 checks passed
fat-catTW pushed a commit to fat-catTW/airflow that referenced this pull request Mar 22, 2026
…pache#63266)
In HA deployments, ti_skip_downstream() issues a bulk UPDATE without
a state guard. When a BranchOperator decides to skip downstream tasks,
it can overwrite a task already RUNNING on a worker to SKIPPED, causing
a 409 heartbeat conflict that kills the task mid-execution.
Add a skippable_state_clause to the UPDATE WHERE clause so RUNNING,
SUCCESS, and FAILED tasks are never overwritten to SKIPPED.
QUEUED tasks are intentionally allowed to be skipped: no work has been
done yet and the BranchOperator's decision should take priority. The
worker pod will get a benign 409 on PATCH /run and exit cleanly.
closes: apache#59378closes: apache#57618
Suraj-kumar00 pushed a commit to Suraj-kumar00/airflow that referenced this pull request Apr 7, 2026
…pache#63266)
In HA deployments, ti_skip_downstream() issues a bulk UPDATE without
a state guard. When a BranchOperator decides to skip downstream tasks,
it can overwrite a task already RUNNING on a worker to SKIPPED, causing
a 409 heartbeat conflict that kills the task mid-execution.
Add a skippable_state_clause to the UPDATE WHERE clause so RUNNING,
SUCCESS, and FAILED tasks are never overwritten to SKIPPED.
QUEUED tasks are intentionally allowed to be skipped: no work has been
done yet and the BranchOperator's decision should take priority. The
worker pod will get a benign 409 on PATCH /run and exit cleanly.
closes: apache#59378closes: apache#57618
abhijeets25012-tech pushed a commit to abhijeets25012-tech/airflow that referenced this pull request Apr 9, 2026
…pache#63266)
In HA deployments, ti_skip_downstream() issues a bulk UPDATE without
a state guard. When a BranchOperator decides to skip downstream tasks,
it can overwrite a task already RUNNING on a worker to SKIPPED, causing
a 409 heartbeat conflict that kills the task mid-execution.
Add a skippable_state_clause to the UPDATE WHERE clause so RUNNING,
SUCCESS, and FAILED tasks are never overwritten to SKIPPED.
QUEUED tasks are intentionally allowed to be skipped: no work has been
done yet and the BranchOperator's decision should take priority. The
worker pod will get a benign 409 on PATCH /run and exit cleanly.
closes: apache#59378closes: apache#57618
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:APIAirflow's REST/HTTP APIarea:task-sdk

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

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

Fix ti_skip_downstream overwriting RUNNING tasks to SKIPPED - #63266

Merged
kaxil merged 1 commit into
apache:mainfrom
sam-dumont:fix/ti-skip-downstream-state-guard
Mar 19, 2026
Merged

Fix ti_skip_downstream overwriting RUNNING tasks to SKIPPED#63266
kaxil merged 1 commit into
apache:mainfrom
sam-dumont:fix/ti-skip-downstream-state-guard

Conversation

@sam-dumont

Copy link
Copy Markdown
Contributor

ti_skip_downstream() issues an UPDATE filtered by (dag_id, run_id, task_id, map_index) without a state guard. When a BranchOperator on one scheduler decides to skip downstream tasks, the UPDATE can overwrite a task already RUNNING on a worker. The worker's next heartbeat returns 409 with current_state: skipped, killing the task mid-execution.

This is a companion fix to #60330, which guards schedule_tis() against the same class of race condition. Different code path (Execution API routes vs dagrun.py), same root cause : unguarded bulk UPDATEs on TI state.

Production data (12 days, 5 schedulers, ~500 concurrent workers)

We deployed both fixes as monkey patches on our prod cluster and monitored 409 heartbeat errors via CloudWatch :

Before any fix 14-169 errors/day
After schedule_tis 3-4/day (all current_state: skipped)
After both fixes 0 errors for 18+ hours
MetricBefore fixesAfter schedule_tis onlyAfter both fixes
Total 409s/day14-1693-40
current_state: scheduledpresent00
current_state: failed47/day00
current_state: skipped8/day2-5/day0

Fix

Add skippable_state_clause to the UPDATE's WHERE clause :

skippable_state_clause=or_(
TI.state.is_(None),
TI.state.not_in([RUNNING, SUCCESS, FAILED]),
)

The or_(IS NULL, NOT IN) pattern handles SQL NULL semantics : NULL NOT IN (...) evaluates to NULL (falsy), so tasks with state=None need an explicit IS NULL check to remain skippable.

QUEUED is intentionally NOT guarded : a QUEUED task hasn't started executing yet, so the BranchOperator's decision should take priority. The worker pod will get a benign 409 on PATCH /run and exit cleanly. Blocking QUEUED would cause a semantic error where the wrong branch executes.

Tests

5 regression tests in TestTISkipDownstreamRaceCondition :

  • RUNNING / SUCCESS / FAILED tasks protected from overwrite (parametrized)
  • QUEUED task correctly skipped (BranchOperator decision wins over queue)
  • None-state task still correctly skipped (happy path)

related: #59378

related: #60330

related: #57618


Was generative AI tooling used to co-author this PR?
  • Yes — Claude Code (claude-opus-4-6)

Generated-by: Claude Code (claude-opus-4-6) following the guidelines

@boring-cyborgboring-cyborgBot added area:API Airflow's REST/HTTP API area:task-sdk labels Mar 10, 2026
@sam-dumont
sam-dumontforce-pushed the fix/ti-skip-downstream-state-guard branch from f4a3ded to 589633aCompareMarch 10, 2026 14:44
@potiuk

potiuk commented Mar 11, 2026

Copy link
Copy Markdown
Member

@sam-dumont This PR has been converted to draft because it does not yet meet our Pull Request quality criteria.

Issues found:

  • Pre-commit / static checks: Failing: CI image checks / Static checks. Run prek run --from-ref main locally to find and fix issues. See Pre-commit / static checks docs.
  • mypy (type checking): Failing: CI image checks / MyPy checks (mypy-airflow-core), CI image checks / MyPy checks (mypy-task-sdk). Run prek --stage manual mypy-airflow-core --all-files && prek --stage manual mypy-task-sdk --all-files locally to reproduce. You need breeze ci-image build --python 3.10 for Docker-based mypy. See mypy (type checking) docs.
  • Provider tests: Failing: provider distributions tests / Compat 2.11.1:P3.10:, Postgres tests: providers / DB-prov:Postgres:14:3.10:-amazon,celer...standard, MySQL tests: providers / DB-prov:MySQL:8.0:3.10:-amazon,celer...standard, Sqlite tests: providers / DB-prov:Sqlite:3.10:-amazon,celer...standard, Non-DB tests: providers / Non-DB-prov::3.10:-amazon,celer...standard (+7 more). Run provider tests with breeze run pytest <provider-test-path> -xvs. See Provider tests docs.
  • Other failing CI checks: Failing: CI image checks / Test Python API client, Postgres tests: core / DB-core:Postgres:14:3.10:API...Serialization, MySQL tests: core / DB-core:MySQL:8.0:3.10:API...Serialization, Sqlite tests: core / DB-core:Sqlite:3.10:API...Serialization, Non-DB tests: core / Non-DB-core::3.10:API...Serialization (+8 more). Run prek run --from-ref main locally to reproduce. See static checks docs.

Note: Your branch is 45 commits behind main. Some check failures may be caused by changes in the base branch rather than by your PR. Please rebase your branch and push again to get up-to-date CI results.

What to do next:

  • The comment informs you what you need to do.
  • Fix each issue, then mark the PR as "Ready for review" in the GitHub UI - but only after making sure that all the issues are fixed.
  • Maintainers will then proceed with a normal review.

Converting a PR to draft is not a rejection — it is an invitation to bring the PR up to the project's standards so that maintainer review time is spent productively. If you have questions, feel free to ask on the Airflow Slack.


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

@potiuk
potiuk marked this pull request as draft March 11, 2026 00:01
@sam-dumont
sam-dumontforce-pushed the fix/ti-skip-downstream-state-guard branch 4 times, most recently from 08ece10 to 9dd13fbCompareMarch 13, 2026 14:29
@sam-dumont
sam-dumont marked this pull request as ready for review March 13, 2026 14:43
@sam-dumont

Copy link
Copy Markdown
ContributorAuthor

@potiuk Hi ! fixed the issues following your advice, thanks a lot ! the 3 failures here seems to be transient, not sure if I can relaunch the job myself, does not seems like it. Hope it's going to be okay now :)

In HA deployments, ti_skip_downstream() issues a bulk UPDATE without
a state guard. When a BranchOperator decides to skip downstream tasks,
it can overwrite a task already RUNNING on a worker to SKIPPED, causing
a 409 heartbeat conflict that kills the task mid-execution.
Add a skippable_state_clause to the UPDATE WHERE clause so RUNNING,
SUCCESS, and FAILED tasks are never overwritten to SKIPPED.
QUEUED tasks are intentionally allowed to be skipped: no work has been
done yet and the BranchOperator's decision should take priority. The
worker pod will get a benign 409 on PATCH /run and exit cleanly.
related: apache#59378
@kaxil
kaxilforce-pushed the fix/ti-skip-downstream-state-guard branch from 9dd13fb to 4d57ba4CompareMarch 19, 2026 01:48
@kaxilkaxil added this to the Airflow 3.2.0 milestone Mar 19, 2026
@kaxil
kaxil merged commit 6d1794a into apache:mainMar 19, 2026
130 of 132 checks passed
fat-catTW pushed a commit to fat-catTW/airflow that referenced this pull request Mar 22, 2026
…pache#63266)
In HA deployments, ti_skip_downstream() issues a bulk UPDATE without
a state guard. When a BranchOperator decides to skip downstream tasks,
it can overwrite a task already RUNNING on a worker to SKIPPED, causing
a 409 heartbeat conflict that kills the task mid-execution.
Add a skippable_state_clause to the UPDATE WHERE clause so RUNNING,
SUCCESS, and FAILED tasks are never overwritten to SKIPPED.
QUEUED tasks are intentionally allowed to be skipped: no work has been
done yet and the BranchOperator's decision should take priority. The
worker pod will get a benign 409 on PATCH /run and exit cleanly.
closes: apache#59378closes: apache#57618
Suraj-kumar00 pushed a commit to Suraj-kumar00/airflow that referenced this pull request Apr 7, 2026
…pache#63266)
In HA deployments, ti_skip_downstream() issues a bulk UPDATE without
a state guard. When a BranchOperator decides to skip downstream tasks,
it can overwrite a task already RUNNING on a worker to SKIPPED, causing
a 409 heartbeat conflict that kills the task mid-execution.
Add a skippable_state_clause to the UPDATE WHERE clause so RUNNING,
SUCCESS, and FAILED tasks are never overwritten to SKIPPED.
QUEUED tasks are intentionally allowed to be skipped: no work has been
done yet and the BranchOperator's decision should take priority. The
worker pod will get a benign 409 on PATCH /run and exit cleanly.
closes: apache#59378closes: apache#57618
abhijeets25012-tech pushed a commit to abhijeets25012-tech/airflow that referenced this pull request Apr 9, 2026
…pache#63266)
In HA deployments, ti_skip_downstream() issues a bulk UPDATE without
a state guard. When a BranchOperator decides to skip downstream tasks,
it can overwrite a task already RUNNING on a worker to SKIPPED, causing
a 409 heartbeat conflict that kills the task mid-execution.
Add a skippable_state_clause to the UPDATE WHERE clause so RUNNING,
SUCCESS, and FAILED tasks are never overwritten to SKIPPED.
QUEUED tasks are intentionally allowed to be skipped: no work has been
done yet and the BranchOperator's decision should take priority. The
worker pod will get a benign 409 on PATCH /run and exit cleanly.
closes: apache#59378closes: apache#57618
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:APIAirflow's REST/HTTP APIarea:task-sdk

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@sam-dumont@potiuk@kaxil