Fix deadlock in ti_update_state: FOR UPDATE OF task_instance only - #67246

Merged
kaxil merged 1 commit into
apache:mainfrom
avolant:fix/ti-state-update-dag-run-deadlock
May 21, 2026
Merged

Fix deadlock in ti_update_state: FOR UPDATE OF task_instance only#67246
kaxil merged 1 commit into
apache:mainfrom
avolant:fix/ti-state-update-dag-run-deadlock

Conversation

@avolant

@avolantavolant commented May 20, 2026

Copy link
Copy Markdown
Contributor

Problem

PATCH /execution/task-instances/{id}/state calls:

ti=session.get(TI, task_instance_id, with_for_update=True)

The TI mapper has dag_run = relationship(..., lazy="joined"), so this emits:

SELECT task_instance.*, dag_run.*FROM task_instance
JOIN dag_run ONdag_run.dag_id=task_instance.dag_idANDdag_run.run_id=task_instance.run_idWHEREtask_instance.id= %s
FOR UPDATE

FOR UPDATE with no OF clause locks every relation in the FROM list, both task_instanceanddag_run. Under concurrent task completions in the same DAG run, all workers serialise on a single dag_run row, deadlocking with the scheduler's TriggerRuleDep queries (which run in transactions that also touch task_instance).

Observed in production: ~2,000 psycopg2.errors.DeadlockDetected + widespread statement_timeout (5 s) errors per hour on the PATCH .../state endpoint, with a mean execution time of 4,297 ms (pure lock-wait; the query does zero disk I/O).

Three other call sites in this file already use with_for_update(of=TI) for exactly this reason (lines 152, 339, 697). The two remaining bare with_for_update=True calls are the ones hit on every normal task completion.

Fix

Scope the lock to task_instance only:

ti=session.get(TI, task_instance_id, with_for_update={"of": TI})

This produces FOR UPDATE OF task_instance, leaving dag_run unlocked and breaking the deadlock cycle.

Testing

Existing unit tests cover the state-update path. No behaviour change — the function still reads the dag_run joinedload; it just no longer holds a row lock on it.

session.get(TI, id, with_for_update=True) emits a SELECT that joins
dag_run (via the lazy="joined" relationship) and applies FOR UPDATE to
both tables. Under concurrent task completions this serialises all
workers on the same dag_run row, producing deadlock cycles with the
scheduler's trigger-rule dependency checks.
Three other callsites in this file already use with_for_update={"of": TI}
for exactly this reason. Apply the same fix to the two remaining callsites
in _create_ti_state_update_query_and_update_state and its error-recovery
path.
@kaxilkaxil added this to the Airflow 3.2.2 milestone May 21, 2026
@kaxil
kaxil merged commit 315d159 into apache:mainMay 21, 2026
138 of 141 checks passed
@github-actions

Copy link
Copy Markdown
Contributor

Backport successfully created: v3-2-test

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-2-testPR Link

@vatsrahul1001vatsrahul1001 added the type:bug-fix Changelog: Bug Fixes label May 21, 2026
vatsrahul1001 pushed a commit that referenced this pull request May 21, 2026
…ing dag_run (#67246) (#67264)
session.get(TI, id, with_for_update=True) emits a SELECT that joins
dag_run (via the lazy="joined" relationship) and applies FOR UPDATE to
both tables. Under concurrent task completions this serialises all
workers on the same dag_run row, producing deadlock cycles with the
scheduler's trigger-rule dependency checks.
Three other callsites in this file already use with_for_update={"of": TI}
for exactly this reason. Apply the same fix to the two remaining callsites
in _create_ti_state_update_query_and_update_state and its error-recovery
path.
(cherry picked from commit 315d159)
Co-authored-by: Arthur <arthur.volant@datadoghq.com>
vatsrahul1001 pushed a commit that referenced this pull request May 21, 2026
…ing dag_run (#67246) (#67264)
session.get(TI, id, with_for_update=True) emits a SELECT that joins
dag_run (via the lazy="joined" relationship) and applies FOR UPDATE to
both tables. Under concurrent task completions this serialises all
workers on the same dag_run row, producing deadlock cycles with the
scheduler's trigger-rule dependency checks.
Three other callsites in this file already use with_for_update={"of": TI}
for exactly this reason. Apply the same fix to the two remaining callsites
in _create_ti_state_update_query_and_update_state and its error-recovery
path.
(cherry picked from commit 315d159)
Co-authored-by: Arthur <arthur.volant@datadoghq.com>
kaxil added a commit that referenced this pull request May 22, 2026
…mit (#67353)
PR #59686 dropped the _handle_fail_fast_for_dag call in the MySQL-TIMESTAMP-limit
branch of the reschedule path based on an incorrect SQLA2 deadlock concern. As a
result, DAGs with fail_fast=True silently fail to stop sibling tasks when a
reschedule date exceeds 2038-01-19 on MySQL.
The actual deadlock that motivated #59686 came from a different path (FOR UPDATE
expanding to the lazy-joined dag_run row), fixed in #67246 by scoping the lock
with with_for_update={"of": TI}. With that scope in place, the fail-fast call is
safe and matches the file's two existing fail-fast sites.
Also drops a second misleading comment in the same function claiming session.get
was avoided to "avoid SQLA2 lock contention issues" -- the code itself is fine;
the rationale was wrong.
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:task-sdktype:bug-fixChangelog: Bug Fixes

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants

@avolant@kaxil@vatsrahul1001@hanxdatadog
, '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 deadlock in ti_update_state: FOR UPDATE OF task_instance only - #67246

Merged
kaxil merged 1 commit into
apache:mainfrom
avolant:fix/ti-state-update-dag-run-deadlock
May 21, 2026
Merged

Fix deadlock in ti_update_state: FOR UPDATE OF task_instance only#67246
kaxil merged 1 commit into
apache:mainfrom
avolant:fix/ti-state-update-dag-run-deadlock

Conversation

@avolant

@avolantavolant commented May 20, 2026

Copy link
Copy Markdown
Contributor

Problem

PATCH /execution/task-instances/{id}/state calls:

ti=session.get(TI, task_instance_id, with_for_update=True)

The TI mapper has dag_run = relationship(..., lazy="joined"), so this emits:

SELECT task_instance.*, dag_run.*FROM task_instance
JOIN dag_run ONdag_run.dag_id=task_instance.dag_idANDdag_run.run_id=task_instance.run_idWHEREtask_instance.id= %s
FOR UPDATE

FOR UPDATE with no OF clause locks every relation in the FROM list, both task_instanceanddag_run. Under concurrent task completions in the same DAG run, all workers serialise on a single dag_run row, deadlocking with the scheduler's TriggerRuleDep queries (which run in transactions that also touch task_instance).

Observed in production: ~2,000 psycopg2.errors.DeadlockDetected + widespread statement_timeout (5 s) errors per hour on the PATCH .../state endpoint, with a mean execution time of 4,297 ms (pure lock-wait; the query does zero disk I/O).

Three other call sites in this file already use with_for_update(of=TI) for exactly this reason (lines 152, 339, 697). The two remaining bare with_for_update=True calls are the ones hit on every normal task completion.

Fix

Scope the lock to task_instance only:

ti=session.get(TI, task_instance_id, with_for_update={"of": TI})

This produces FOR UPDATE OF task_instance, leaving dag_run unlocked and breaking the deadlock cycle.

Testing

Existing unit tests cover the state-update path. No behaviour change — the function still reads the dag_run joinedload; it just no longer holds a row lock on it.

session.get(TI, id, with_for_update=True) emits a SELECT that joins
dag_run (via the lazy="joined" relationship) and applies FOR UPDATE to
both tables. Under concurrent task completions this serialises all
workers on the same dag_run row, producing deadlock cycles with the
scheduler's trigger-rule dependency checks.
Three other callsites in this file already use with_for_update={"of": TI}
for exactly this reason. Apply the same fix to the two remaining callsites
in _create_ti_state_update_query_and_update_state and its error-recovery
path.
@kaxilkaxil added this to the Airflow 3.2.2 milestone May 21, 2026
@kaxil
kaxil merged commit 315d159 into apache:mainMay 21, 2026
138 of 141 checks passed
@github-actions

Copy link
Copy Markdown
Contributor

Backport successfully created: v3-2-test

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-2-testPR Link

@vatsrahul1001vatsrahul1001 added the type:bug-fix Changelog: Bug Fixes label May 21, 2026
vatsrahul1001 pushed a commit that referenced this pull request May 21, 2026
…ing dag_run (#67246) (#67264)
session.get(TI, id, with_for_update=True) emits a SELECT that joins
dag_run (via the lazy="joined" relationship) and applies FOR UPDATE to
both tables. Under concurrent task completions this serialises all
workers on the same dag_run row, producing deadlock cycles with the
scheduler's trigger-rule dependency checks.
Three other callsites in this file already use with_for_update={"of": TI}
for exactly this reason. Apply the same fix to the two remaining callsites
in _create_ti_state_update_query_and_update_state and its error-recovery
path.
(cherry picked from commit 315d159)
Co-authored-by: Arthur <arthur.volant@datadoghq.com>
vatsrahul1001 pushed a commit that referenced this pull request May 21, 2026
…ing dag_run (#67246) (#67264)
session.get(TI, id, with_for_update=True) emits a SELECT that joins
dag_run (via the lazy="joined" relationship) and applies FOR UPDATE to
both tables. Under concurrent task completions this serialises all
workers on the same dag_run row, producing deadlock cycles with the
scheduler's trigger-rule dependency checks.
Three other callsites in this file already use with_for_update={"of": TI}
for exactly this reason. Apply the same fix to the two remaining callsites
in _create_ti_state_update_query_and_update_state and its error-recovery
path.
(cherry picked from commit 315d159)
Co-authored-by: Arthur <arthur.volant@datadoghq.com>
kaxil added a commit that referenced this pull request May 22, 2026
…mit (#67353)
PR #59686 dropped the _handle_fail_fast_for_dag call in the MySQL-TIMESTAMP-limit
branch of the reschedule path based on an incorrect SQLA2 deadlock concern. As a
result, DAGs with fail_fast=True silently fail to stop sibling tasks when a
reschedule date exceeds 2038-01-19 on MySQL.
The actual deadlock that motivated #59686 came from a different path (FOR UPDATE
expanding to the lazy-joined dag_run row), fixed in #67246 by scoping the lock
with with_for_update={"of": TI}. With that scope in place, the fail-fast call is
safe and matches the file's two existing fail-fast sites.
Also drops a second misleading comment in the same function claiming session.get
was avoided to "avoid SQLA2 lock contention issues" -- the code itself is fine;
the rationale was wrong.
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:task-sdktype:bug-fixChangelog: Bug Fixes

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants

@avolant@kaxil@vatsrahul1001@hanxdatadog
, '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 deadlock in ti_update_state: FOR UPDATE OF task_instance only - #67246

Merged
kaxil merged 1 commit into
apache:mainfrom
avolant:fix/ti-state-update-dag-run-deadlock
May 21, 2026
Merged

Fix deadlock in ti_update_state: FOR UPDATE OF task_instance only#67246
kaxil merged 1 commit into
apache:mainfrom
avolant:fix/ti-state-update-dag-run-deadlock

Conversation

@avolant

@avolantavolant commented May 20, 2026

Copy link
Copy Markdown
Contributor

Problem

PATCH /execution/task-instances/{id}/state calls:

ti=session.get(TI, task_instance_id, with_for_update=True)

The TI mapper has dag_run = relationship(..., lazy="joined"), so this emits:

SELECT task_instance.*, dag_run.*FROM task_instance
JOIN dag_run ONdag_run.dag_id=task_instance.dag_idANDdag_run.run_id=task_instance.run_idWHEREtask_instance.id= %s
FOR UPDATE

FOR UPDATE with no OF clause locks every relation in the FROM list, both task_instanceanddag_run. Under concurrent task completions in the same DAG run, all workers serialise on a single dag_run row, deadlocking with the scheduler's TriggerRuleDep queries (which run in transactions that also touch task_instance).

Observed in production: ~2,000 psycopg2.errors.DeadlockDetected + widespread statement_timeout (5 s) errors per hour on the PATCH .../state endpoint, with a mean execution time of 4,297 ms (pure lock-wait; the query does zero disk I/O).

Three other call sites in this file already use with_for_update(of=TI) for exactly this reason (lines 152, 339, 697). The two remaining bare with_for_update=True calls are the ones hit on every normal task completion.

Fix

Scope the lock to task_instance only:

ti=session.get(TI, task_instance_id, with_for_update={"of": TI})

This produces FOR UPDATE OF task_instance, leaving dag_run unlocked and breaking the deadlock cycle.

Testing

Existing unit tests cover the state-update path. No behaviour change — the function still reads the dag_run joinedload; it just no longer holds a row lock on it.

session.get(TI, id, with_for_update=True) emits a SELECT that joins
dag_run (via the lazy="joined" relationship) and applies FOR UPDATE to
both tables. Under concurrent task completions this serialises all
workers on the same dag_run row, producing deadlock cycles with the
scheduler's trigger-rule dependency checks.
Three other callsites in this file already use with_for_update={"of": TI}
for exactly this reason. Apply the same fix to the two remaining callsites
in _create_ti_state_update_query_and_update_state and its error-recovery
path.
@kaxilkaxil added this to the Airflow 3.2.2 milestone May 21, 2026
@kaxil
kaxil merged commit 315d159 into apache:mainMay 21, 2026
138 of 141 checks passed
@github-actions

Copy link
Copy Markdown
Contributor

Backport successfully created: v3-2-test

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-2-testPR Link

@vatsrahul1001vatsrahul1001 added the type:bug-fix Changelog: Bug Fixes label May 21, 2026
vatsrahul1001 pushed a commit that referenced this pull request May 21, 2026
…ing dag_run (#67246) (#67264)
session.get(TI, id, with_for_update=True) emits a SELECT that joins
dag_run (via the lazy="joined" relationship) and applies FOR UPDATE to
both tables. Under concurrent task completions this serialises all
workers on the same dag_run row, producing deadlock cycles with the
scheduler's trigger-rule dependency checks.
Three other callsites in this file already use with_for_update={"of": TI}
for exactly this reason. Apply the same fix to the two remaining callsites
in _create_ti_state_update_query_and_update_state and its error-recovery
path.
(cherry picked from commit 315d159)
Co-authored-by: Arthur <arthur.volant@datadoghq.com>
vatsrahul1001 pushed a commit that referenced this pull request May 21, 2026
…ing dag_run (#67246) (#67264)
session.get(TI, id, with_for_update=True) emits a SELECT that joins
dag_run (via the lazy="joined" relationship) and applies FOR UPDATE to
both tables. Under concurrent task completions this serialises all
workers on the same dag_run row, producing deadlock cycles with the
scheduler's trigger-rule dependency checks.
Three other callsites in this file already use with_for_update={"of": TI}
for exactly this reason. Apply the same fix to the two remaining callsites
in _create_ti_state_update_query_and_update_state and its error-recovery
path.
(cherry picked from commit 315d159)
Co-authored-by: Arthur <arthur.volant@datadoghq.com>
kaxil added a commit that referenced this pull request May 22, 2026
…mit (#67353)
PR #59686 dropped the _handle_fail_fast_for_dag call in the MySQL-TIMESTAMP-limit
branch of the reschedule path based on an incorrect SQLA2 deadlock concern. As a
result, DAGs with fail_fast=True silently fail to stop sibling tasks when a
reschedule date exceeds 2038-01-19 on MySQL.
The actual deadlock that motivated #59686 came from a different path (FOR UPDATE
expanding to the lazy-joined dag_run row), fixed in #67246 by scoping the lock
with with_for_update={"of": TI}. With that scope in place, the fail-fast call is
safe and matches the file's two existing fail-fast sites.
Also drops a second misleading comment in the same function claiming session.get
was avoided to "avoid SQLA2 lock contention issues" -- the code itself is fine;
the rationale was wrong.
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:task-sdktype:bug-fixChangelog: Bug Fixes

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants

@avolant@kaxil@vatsrahul1001@hanxdatadog
, '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 deadlock in ti_update_state: FOR UPDATE OF task_instance only - #67246

Merged
kaxil merged 1 commit into
apache:mainfrom
avolant:fix/ti-state-update-dag-run-deadlock
May 21, 2026
Merged

Fix deadlock in ti_update_state: FOR UPDATE OF task_instance only#67246
kaxil merged 1 commit into
apache:mainfrom
avolant:fix/ti-state-update-dag-run-deadlock

Conversation

@avolant

@avolantavolant commented May 20, 2026

Copy link
Copy Markdown
Contributor

Problem

PATCH /execution/task-instances/{id}/state calls:

ti=session.get(TI, task_instance_id, with_for_update=True)

The TI mapper has dag_run = relationship(..., lazy="joined"), so this emits:

SELECT task_instance.*, dag_run.*FROM task_instance
JOIN dag_run ONdag_run.dag_id=task_instance.dag_idANDdag_run.run_id=task_instance.run_idWHEREtask_instance.id= %s
FOR UPDATE

FOR UPDATE with no OF clause locks every relation in the FROM list, both task_instanceanddag_run. Under concurrent task completions in the same DAG run, all workers serialise on a single dag_run row, deadlocking with the scheduler's TriggerRuleDep queries (which run in transactions that also touch task_instance).

Observed in production: ~2,000 psycopg2.errors.DeadlockDetected + widespread statement_timeout (5 s) errors per hour on the PATCH .../state endpoint, with a mean execution time of 4,297 ms (pure lock-wait; the query does zero disk I/O).

Three other call sites in this file already use with_for_update(of=TI) for exactly this reason (lines 152, 339, 697). The two remaining bare with_for_update=True calls are the ones hit on every normal task completion.

Fix

Scope the lock to task_instance only:

ti=session.get(TI, task_instance_id, with_for_update={"of": TI})

This produces FOR UPDATE OF task_instance, leaving dag_run unlocked and breaking the deadlock cycle.

Testing

Existing unit tests cover the state-update path. No behaviour change — the function still reads the dag_run joinedload; it just no longer holds a row lock on it.

session.get(TI, id, with_for_update=True) emits a SELECT that joins
dag_run (via the lazy="joined" relationship) and applies FOR UPDATE to
both tables. Under concurrent task completions this serialises all
workers on the same dag_run row, producing deadlock cycles with the
scheduler's trigger-rule dependency checks.
Three other callsites in this file already use with_for_update={"of": TI}
for exactly this reason. Apply the same fix to the two remaining callsites
in _create_ti_state_update_query_and_update_state and its error-recovery
path.
@kaxilkaxil added this to the Airflow 3.2.2 milestone May 21, 2026
@kaxil
kaxil merged commit 315d159 into apache:mainMay 21, 2026
138 of 141 checks passed
@github-actions

Copy link
Copy Markdown
Contributor

Backport successfully created: v3-2-test

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-2-testPR Link

@vatsrahul1001vatsrahul1001 added the type:bug-fix Changelog: Bug Fixes label May 21, 2026
vatsrahul1001 pushed a commit that referenced this pull request May 21, 2026
…ing dag_run (#67246) (#67264)
session.get(TI, id, with_for_update=True) emits a SELECT that joins
dag_run (via the lazy="joined" relationship) and applies FOR UPDATE to
both tables. Under concurrent task completions this serialises all
workers on the same dag_run row, producing deadlock cycles with the
scheduler's trigger-rule dependency checks.
Three other callsites in this file already use with_for_update={"of": TI}
for exactly this reason. Apply the same fix to the two remaining callsites
in _create_ti_state_update_query_and_update_state and its error-recovery
path.
(cherry picked from commit 315d159)
Co-authored-by: Arthur <arthur.volant@datadoghq.com>
vatsrahul1001 pushed a commit that referenced this pull request May 21, 2026
…ing dag_run (#67246) (#67264)
session.get(TI, id, with_for_update=True) emits a SELECT that joins
dag_run (via the lazy="joined" relationship) and applies FOR UPDATE to
both tables. Under concurrent task completions this serialises all
workers on the same dag_run row, producing deadlock cycles with the
scheduler's trigger-rule dependency checks.
Three other callsites in this file already use with_for_update={"of": TI}
for exactly this reason. Apply the same fix to the two remaining callsites
in _create_ti_state_update_query_and_update_state and its error-recovery
path.
(cherry picked from commit 315d159)
Co-authored-by: Arthur <arthur.volant@datadoghq.com>
kaxil added a commit that referenced this pull request May 22, 2026
…mit (#67353)
PR #59686 dropped the _handle_fail_fast_for_dag call in the MySQL-TIMESTAMP-limit
branch of the reschedule path based on an incorrect SQLA2 deadlock concern. As a
result, DAGs with fail_fast=True silently fail to stop sibling tasks when a
reschedule date exceeds 2038-01-19 on MySQL.
The actual deadlock that motivated #59686 came from a different path (FOR UPDATE
expanding to the lazy-joined dag_run row), fixed in #67246 by scoping the lock
with with_for_update={"of": TI}. With that scope in place, the fail-fast call is
safe and matches the file's two existing fail-fast sites.
Also drops a second misleading comment in the same function claiming session.get
was avoided to "avoid SQLA2 lock contention issues" -- the code itself is fine;
the rationale was wrong.
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:task-sdktype:bug-fixChangelog: Bug Fixes

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants

@avolant@kaxil@vatsrahul1001@hanxdatadog
, '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 deadlock in ti_update_state: FOR UPDATE OF task_instance only - #67246

Merged
kaxil merged 1 commit into
apache:mainfrom
avolant:fix/ti-state-update-dag-run-deadlock
May 21, 2026
Merged

Fix deadlock in ti_update_state: FOR UPDATE OF task_instance only#67246
kaxil merged 1 commit into
apache:mainfrom
avolant:fix/ti-state-update-dag-run-deadlock

Conversation

@avolant

@avolantavolant commented May 20, 2026

Copy link
Copy Markdown
Contributor

Problem

PATCH /execution/task-instances/{id}/state calls:

ti=session.get(TI, task_instance_id, with_for_update=True)

The TI mapper has dag_run = relationship(..., lazy="joined"), so this emits:

SELECT task_instance.*, dag_run.*FROM task_instance
JOIN dag_run ONdag_run.dag_id=task_instance.dag_idANDdag_run.run_id=task_instance.run_idWHEREtask_instance.id= %s
FOR UPDATE

FOR UPDATE with no OF clause locks every relation in the FROM list, both task_instanceanddag_run. Under concurrent task completions in the same DAG run, all workers serialise on a single dag_run row, deadlocking with the scheduler's TriggerRuleDep queries (which run in transactions that also touch task_instance).

Observed in production: ~2,000 psycopg2.errors.DeadlockDetected + widespread statement_timeout (5 s) errors per hour on the PATCH .../state endpoint, with a mean execution time of 4,297 ms (pure lock-wait; the query does zero disk I/O).

Three other call sites in this file already use with_for_update(of=TI) for exactly this reason (lines 152, 339, 697). The two remaining bare with_for_update=True calls are the ones hit on every normal task completion.

Fix

Scope the lock to task_instance only:

ti=session.get(TI, task_instance_id, with_for_update={"of": TI})

This produces FOR UPDATE OF task_instance, leaving dag_run unlocked and breaking the deadlock cycle.

Testing

Existing unit tests cover the state-update path. No behaviour change — the function still reads the dag_run joinedload; it just no longer holds a row lock on it.

session.get(TI, id, with_for_update=True) emits a SELECT that joins
dag_run (via the lazy="joined" relationship) and applies FOR UPDATE to
both tables. Under concurrent task completions this serialises all
workers on the same dag_run row, producing deadlock cycles with the
scheduler's trigger-rule dependency checks.
Three other callsites in this file already use with_for_update={"of": TI}
for exactly this reason. Apply the same fix to the two remaining callsites
in _create_ti_state_update_query_and_update_state and its error-recovery
path.
@kaxilkaxil added this to the Airflow 3.2.2 milestone May 21, 2026
@kaxil
kaxil merged commit 315d159 into apache:mainMay 21, 2026
138 of 141 checks passed
@github-actions

Copy link
Copy Markdown
Contributor

Backport successfully created: v3-2-test

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-2-testPR Link

@vatsrahul1001vatsrahul1001 added the type:bug-fix Changelog: Bug Fixes label May 21, 2026
vatsrahul1001 pushed a commit that referenced this pull request May 21, 2026
…ing dag_run (#67246) (#67264)
session.get(TI, id, with_for_update=True) emits a SELECT that joins
dag_run (via the lazy="joined" relationship) and applies FOR UPDATE to
both tables. Under concurrent task completions this serialises all
workers on the same dag_run row, producing deadlock cycles with the
scheduler's trigger-rule dependency checks.
Three other callsites in this file already use with_for_update={"of": TI}
for exactly this reason. Apply the same fix to the two remaining callsites
in _create_ti_state_update_query_and_update_state and its error-recovery
path.
(cherry picked from commit 315d159)
Co-authored-by: Arthur <arthur.volant@datadoghq.com>
vatsrahul1001 pushed a commit that referenced this pull request May 21, 2026
…ing dag_run (#67246) (#67264)
session.get(TI, id, with_for_update=True) emits a SELECT that joins
dag_run (via the lazy="joined" relationship) and applies FOR UPDATE to
both tables. Under concurrent task completions this serialises all
workers on the same dag_run row, producing deadlock cycles with the
scheduler's trigger-rule dependency checks.
Three other callsites in this file already use with_for_update={"of": TI}
for exactly this reason. Apply the same fix to the two remaining callsites
in _create_ti_state_update_query_and_update_state and its error-recovery
path.
(cherry picked from commit 315d159)
Co-authored-by: Arthur <arthur.volant@datadoghq.com>
kaxil added a commit that referenced this pull request May 22, 2026
…mit (#67353)
PR #59686 dropped the _handle_fail_fast_for_dag call in the MySQL-TIMESTAMP-limit
branch of the reschedule path based on an incorrect SQLA2 deadlock concern. As a
result, DAGs with fail_fast=True silently fail to stop sibling tasks when a
reschedule date exceeds 2038-01-19 on MySQL.
The actual deadlock that motivated #59686 came from a different path (FOR UPDATE
expanding to the lazy-joined dag_run row), fixed in #67246 by scoping the lock
with with_for_update={"of": TI}. With that scope in place, the fail-fast call is
safe and matches the file's two existing fail-fast sites.
Also drops a second misleading comment in the same function claiming session.get
was avoided to "avoid SQLA2 lock contention issues" -- the code itself is fine;
the rationale was wrong.
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:task-sdktype:bug-fixChangelog: Bug Fixes

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants

@avolant@kaxil@vatsrahul1001@hanxdatadog
, '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 deadlock in ti_update_state: FOR UPDATE OF task_instance only - #67246

Merged
kaxil merged 1 commit into
apache:mainfrom
avolant:fix/ti-state-update-dag-run-deadlock
May 21, 2026
Merged

Fix deadlock in ti_update_state: FOR UPDATE OF task_instance only#67246
kaxil merged 1 commit into
apache:mainfrom
avolant:fix/ti-state-update-dag-run-deadlock

Conversation

@avolant

@avolantavolant commented May 20, 2026

Copy link
Copy Markdown
Contributor

Problem

PATCH /execution/task-instances/{id}/state calls:

ti=session.get(TI, task_instance_id, with_for_update=True)

The TI mapper has dag_run = relationship(..., lazy="joined"), so this emits:

SELECT task_instance.*, dag_run.*FROM task_instance
JOIN dag_run ONdag_run.dag_id=task_instance.dag_idANDdag_run.run_id=task_instance.run_idWHEREtask_instance.id= %s
FOR UPDATE

FOR UPDATE with no OF clause locks every relation in the FROM list, both task_instanceanddag_run. Under concurrent task completions in the same DAG run, all workers serialise on a single dag_run row, deadlocking with the scheduler's TriggerRuleDep queries (which run in transactions that also touch task_instance).

Observed in production: ~2,000 psycopg2.errors.DeadlockDetected + widespread statement_timeout (5 s) errors per hour on the PATCH .../state endpoint, with a mean execution time of 4,297 ms (pure lock-wait; the query does zero disk I/O).

Three other call sites in this file already use with_for_update(of=TI) for exactly this reason (lines 152, 339, 697). The two remaining bare with_for_update=True calls are the ones hit on every normal task completion.

Fix

Scope the lock to task_instance only:

ti=session.get(TI, task_instance_id, with_for_update={"of": TI})

This produces FOR UPDATE OF task_instance, leaving dag_run unlocked and breaking the deadlock cycle.

Testing

Existing unit tests cover the state-update path. No behaviour change — the function still reads the dag_run joinedload; it just no longer holds a row lock on it.

session.get(TI, id, with_for_update=True) emits a SELECT that joins
dag_run (via the lazy="joined" relationship) and applies FOR UPDATE to
both tables. Under concurrent task completions this serialises all
workers on the same dag_run row, producing deadlock cycles with the
scheduler's trigger-rule dependency checks.
Three other callsites in this file already use with_for_update={"of": TI}
for exactly this reason. Apply the same fix to the two remaining callsites
in _create_ti_state_update_query_and_update_state and its error-recovery
path.
@kaxilkaxil added this to the Airflow 3.2.2 milestone May 21, 2026
@kaxil
kaxil merged commit 315d159 into apache:mainMay 21, 2026
138 of 141 checks passed
@github-actions

Copy link
Copy Markdown
Contributor

Backport successfully created: v3-2-test

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-2-testPR Link

@vatsrahul1001vatsrahul1001 added the type:bug-fix Changelog: Bug Fixes label May 21, 2026
vatsrahul1001 pushed a commit that referenced this pull request May 21, 2026
…ing dag_run (#67246) (#67264)
session.get(TI, id, with_for_update=True) emits a SELECT that joins
dag_run (via the lazy="joined" relationship) and applies FOR UPDATE to
both tables. Under concurrent task completions this serialises all
workers on the same dag_run row, producing deadlock cycles with the
scheduler's trigger-rule dependency checks.
Three other callsites in this file already use with_for_update={"of": TI}
for exactly this reason. Apply the same fix to the two remaining callsites
in _create_ti_state_update_query_and_update_state and its error-recovery
path.
(cherry picked from commit 315d159)
Co-authored-by: Arthur <arthur.volant@datadoghq.com>
vatsrahul1001 pushed a commit that referenced this pull request May 21, 2026
…ing dag_run (#67246) (#67264)
session.get(TI, id, with_for_update=True) emits a SELECT that joins
dag_run (via the lazy="joined" relationship) and applies FOR UPDATE to
both tables. Under concurrent task completions this serialises all
workers on the same dag_run row, producing deadlock cycles with the
scheduler's trigger-rule dependency checks.
Three other callsites in this file already use with_for_update={"of": TI}
for exactly this reason. Apply the same fix to the two remaining callsites
in _create_ti_state_update_query_and_update_state and its error-recovery
path.
(cherry picked from commit 315d159)
Co-authored-by: Arthur <arthur.volant@datadoghq.com>
kaxil added a commit that referenced this pull request May 22, 2026
…mit (#67353)
PR #59686 dropped the _handle_fail_fast_for_dag call in the MySQL-TIMESTAMP-limit
branch of the reschedule path based on an incorrect SQLA2 deadlock concern. As a
result, DAGs with fail_fast=True silently fail to stop sibling tasks when a
reschedule date exceeds 2038-01-19 on MySQL.
The actual deadlock that motivated #59686 came from a different path (FOR UPDATE
expanding to the lazy-joined dag_run row), fixed in #67246 by scoping the lock
with with_for_update={"of": TI}. With that scope in place, the fail-fast call is
safe and matches the file's two existing fail-fast sites.
Also drops a second misleading comment in the same function claiming session.get
was avoided to "avoid SQLA2 lock contention issues" -- the code itself is fine;
the rationale was wrong.
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:task-sdktype:bug-fixChangelog: Bug Fixes

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants

@avolant@kaxil@vatsrahul1001@hanxdatadog
, '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 deadlock in ti_update_state: FOR UPDATE OF task_instance only - #67246

Merged
kaxil merged 1 commit into
apache:mainfrom
avolant:fix/ti-state-update-dag-run-deadlock
May 21, 2026
Merged

Fix deadlock in ti_update_state: FOR UPDATE OF task_instance only#67246
kaxil merged 1 commit into
apache:mainfrom
avolant:fix/ti-state-update-dag-run-deadlock

Conversation

@avolant

@avolantavolant commented May 20, 2026

Copy link
Copy Markdown
Contributor

Problem

PATCH /execution/task-instances/{id}/state calls:

ti=session.get(TI, task_instance_id, with_for_update=True)

The TI mapper has dag_run = relationship(..., lazy="joined"), so this emits:

SELECT task_instance.*, dag_run.*FROM task_instance
JOIN dag_run ONdag_run.dag_id=task_instance.dag_idANDdag_run.run_id=task_instance.run_idWHEREtask_instance.id= %s
FOR UPDATE

FOR UPDATE with no OF clause locks every relation in the FROM list, both task_instanceanddag_run. Under concurrent task completions in the same DAG run, all workers serialise on a single dag_run row, deadlocking with the scheduler's TriggerRuleDep queries (which run in transactions that also touch task_instance).

Observed in production: ~2,000 psycopg2.errors.DeadlockDetected + widespread statement_timeout (5 s) errors per hour on the PATCH .../state endpoint, with a mean execution time of 4,297 ms (pure lock-wait; the query does zero disk I/O).

Three other call sites in this file already use with_for_update(of=TI) for exactly this reason (lines 152, 339, 697). The two remaining bare with_for_update=True calls are the ones hit on every normal task completion.

Fix

Scope the lock to task_instance only:

ti=session.get(TI, task_instance_id, with_for_update={"of": TI})

This produces FOR UPDATE OF task_instance, leaving dag_run unlocked and breaking the deadlock cycle.

Testing

Existing unit tests cover the state-update path. No behaviour change — the function still reads the dag_run joinedload; it just no longer holds a row lock on it.

session.get(TI, id, with_for_update=True) emits a SELECT that joins
dag_run (via the lazy="joined" relationship) and applies FOR UPDATE to
both tables. Under concurrent task completions this serialises all
workers on the same dag_run row, producing deadlock cycles with the
scheduler's trigger-rule dependency checks.
Three other callsites in this file already use with_for_update={"of": TI}
for exactly this reason. Apply the same fix to the two remaining callsites
in _create_ti_state_update_query_and_update_state and its error-recovery
path.
@kaxilkaxil added this to the Airflow 3.2.2 milestone May 21, 2026
@kaxil
kaxil merged commit 315d159 into apache:mainMay 21, 2026
138 of 141 checks passed
@github-actions

Copy link
Copy Markdown
Contributor

Backport successfully created: v3-2-test

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-2-testPR Link

@vatsrahul1001vatsrahul1001 added the type:bug-fix Changelog: Bug Fixes label May 21, 2026
vatsrahul1001 pushed a commit that referenced this pull request May 21, 2026
…ing dag_run (#67246) (#67264)
session.get(TI, id, with_for_update=True) emits a SELECT that joins
dag_run (via the lazy="joined" relationship) and applies FOR UPDATE to
both tables. Under concurrent task completions this serialises all
workers on the same dag_run row, producing deadlock cycles with the
scheduler's trigger-rule dependency checks.
Three other callsites in this file already use with_for_update={"of": TI}
for exactly this reason. Apply the same fix to the two remaining callsites
in _create_ti_state_update_query_and_update_state and its error-recovery
path.
(cherry picked from commit 315d159)
Co-authored-by: Arthur <arthur.volant@datadoghq.com>
vatsrahul1001 pushed a commit that referenced this pull request May 21, 2026
…ing dag_run (#67246) (#67264)
session.get(TI, id, with_for_update=True) emits a SELECT that joins
dag_run (via the lazy="joined" relationship) and applies FOR UPDATE to
both tables. Under concurrent task completions this serialises all
workers on the same dag_run row, producing deadlock cycles with the
scheduler's trigger-rule dependency checks.
Three other callsites in this file already use with_for_update={"of": TI}
for exactly this reason. Apply the same fix to the two remaining callsites
in _create_ti_state_update_query_and_update_state and its error-recovery
path.
(cherry picked from commit 315d159)
Co-authored-by: Arthur <arthur.volant@datadoghq.com>
kaxil added a commit that referenced this pull request May 22, 2026
…mit (#67353)
PR #59686 dropped the _handle_fail_fast_for_dag call in the MySQL-TIMESTAMP-limit
branch of the reschedule path based on an incorrect SQLA2 deadlock concern. As a
result, DAGs with fail_fast=True silently fail to stop sibling tasks when a
reschedule date exceeds 2038-01-19 on MySQL.
The actual deadlock that motivated #59686 came from a different path (FOR UPDATE
expanding to the lazy-joined dag_run row), fixed in #67246 by scoping the lock
with with_for_update={"of": TI}. With that scope in place, the fail-fast call is
safe and matches the file's two existing fail-fast sites.
Also drops a second misleading comment in the same function claiming session.get
was avoided to "avoid SQLA2 lock contention issues" -- the code itself is fine;
the rationale was wrong.
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:task-sdktype:bug-fixChangelog: Bug Fixes

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants

@avolant@kaxil@vatsrahul1001@hanxdatadog
, '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 deadlock in ti_update_state: FOR UPDATE OF task_instance only - #67246

Merged
kaxil merged 1 commit into
apache:mainfrom
avolant:fix/ti-state-update-dag-run-deadlock
May 21, 2026
Merged

Fix deadlock in ti_update_state: FOR UPDATE OF task_instance only#67246
kaxil merged 1 commit into
apache:mainfrom
avolant:fix/ti-state-update-dag-run-deadlock

Conversation

@avolant

@avolantavolant commented May 20, 2026

Copy link
Copy Markdown
Contributor

Problem

PATCH /execution/task-instances/{id}/state calls:

ti=session.get(TI, task_instance_id, with_for_update=True)

The TI mapper has dag_run = relationship(..., lazy="joined"), so this emits:

SELECT task_instance.*, dag_run.*FROM task_instance
JOIN dag_run ONdag_run.dag_id=task_instance.dag_idANDdag_run.run_id=task_instance.run_idWHEREtask_instance.id= %s
FOR UPDATE

FOR UPDATE with no OF clause locks every relation in the FROM list, both task_instanceanddag_run. Under concurrent task completions in the same DAG run, all workers serialise on a single dag_run row, deadlocking with the scheduler's TriggerRuleDep queries (which run in transactions that also touch task_instance).

Observed in production: ~2,000 psycopg2.errors.DeadlockDetected + widespread statement_timeout (5 s) errors per hour on the PATCH .../state endpoint, with a mean execution time of 4,297 ms (pure lock-wait; the query does zero disk I/O).

Three other call sites in this file already use with_for_update(of=TI) for exactly this reason (lines 152, 339, 697). The two remaining bare with_for_update=True calls are the ones hit on every normal task completion.

Fix

Scope the lock to task_instance only:

ti=session.get(TI, task_instance_id, with_for_update={"of": TI})

This produces FOR UPDATE OF task_instance, leaving dag_run unlocked and breaking the deadlock cycle.

Testing

Existing unit tests cover the state-update path. No behaviour change — the function still reads the dag_run joinedload; it just no longer holds a row lock on it.

session.get(TI, id, with_for_update=True) emits a SELECT that joins
dag_run (via the lazy="joined" relationship) and applies FOR UPDATE to
both tables. Under concurrent task completions this serialises all
workers on the same dag_run row, producing deadlock cycles with the
scheduler's trigger-rule dependency checks.
Three other callsites in this file already use with_for_update={"of": TI}
for exactly this reason. Apply the same fix to the two remaining callsites
in _create_ti_state_update_query_and_update_state and its error-recovery
path.
@kaxilkaxil added this to the Airflow 3.2.2 milestone May 21, 2026
@kaxil
kaxil merged commit 315d159 into apache:mainMay 21, 2026
138 of 141 checks passed
@github-actions

Copy link
Copy Markdown
Contributor

Backport successfully created: v3-2-test

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-2-testPR Link

@vatsrahul1001vatsrahul1001 added the type:bug-fix Changelog: Bug Fixes label May 21, 2026
vatsrahul1001 pushed a commit that referenced this pull request May 21, 2026
…ing dag_run (#67246) (#67264)
session.get(TI, id, with_for_update=True) emits a SELECT that joins
dag_run (via the lazy="joined" relationship) and applies FOR UPDATE to
both tables. Under concurrent task completions this serialises all
workers on the same dag_run row, producing deadlock cycles with the
scheduler's trigger-rule dependency checks.
Three other callsites in this file already use with_for_update={"of": TI}
for exactly this reason. Apply the same fix to the two remaining callsites
in _create_ti_state_update_query_and_update_state and its error-recovery
path.
(cherry picked from commit 315d159)
Co-authored-by: Arthur <arthur.volant@datadoghq.com>
vatsrahul1001 pushed a commit that referenced this pull request May 21, 2026
…ing dag_run (#67246) (#67264)
session.get(TI, id, with_for_update=True) emits a SELECT that joins
dag_run (via the lazy="joined" relationship) and applies FOR UPDATE to
both tables. Under concurrent task completions this serialises all
workers on the same dag_run row, producing deadlock cycles with the
scheduler's trigger-rule dependency checks.
Three other callsites in this file already use with_for_update={"of": TI}
for exactly this reason. Apply the same fix to the two remaining callsites
in _create_ti_state_update_query_and_update_state and its error-recovery
path.
(cherry picked from commit 315d159)
Co-authored-by: Arthur <arthur.volant@datadoghq.com>
kaxil added a commit that referenced this pull request May 22, 2026
…mit (#67353)
PR #59686 dropped the _handle_fail_fast_for_dag call in the MySQL-TIMESTAMP-limit
branch of the reschedule path based on an incorrect SQLA2 deadlock concern. As a
result, DAGs with fail_fast=True silently fail to stop sibling tasks when a
reschedule date exceeds 2038-01-19 on MySQL.
The actual deadlock that motivated #59686 came from a different path (FOR UPDATE
expanding to the lazy-joined dag_run row), fixed in #67246 by scoping the lock
with with_for_update={"of": TI}. With that scope in place, the fail-fast call is
safe and matches the file's two existing fail-fast sites.
Also drops a second misleading comment in the same function claiming session.get
was avoided to "avoid SQLA2 lock contention issues" -- the code itself is fine;
the rationale was wrong.
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:task-sdktype:bug-fixChangelog: Bug Fixes

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants

@avolant@kaxil@vatsrahul1001@hanxdatadog