Register asset events before taking the task_instance row lock - #70078

Closed
Andrushika wants to merge 1 commit into
apache:mainfrom
Andrushika:reorder-asset-registration-before-ti-lock
Closed

Register asset events before taking the task_instance row lock#70078
Andrushika wants to merge 1 commit into
apache:mainfrom
Andrushika:reorder-asset-registration-before-ti-lock

Conversation

@Andrushika

@AndrushikaAndrushika commented Jul 18, 2026

Copy link
Copy Markdown
Contributor

Register asset events before taking the task_instance row lock

Why

When a task succeeds, ti_update_state currently takes the task_instance row lock and then registers asset events while holding the lock.

Asset registration can be slow because it performs asset lookups, event inserts, and Dag-run queueing. When many tasks finish at the same time, concurrent requests can wait for the row lock while each request continues to use an API-server thread. In the case reported in #66853, this caused the API server to run out of memory and get OOMKilled.

What

This change registers the asset events before taking the task_instance row lock:
image
The asset events and the task state are still written in the same database transaction.

In this PR, we first read the state without the row lock, then perform the asset registration. This can hit a TOCTOU problem: the state may change in between, for example, when the worker retries. So before we actually change the state, we take the row lock and read the state one more time. The lock makes sure no one else can change it after this point. Then we decide what to do.

There are three cases:

  1. The state is still RUNNING: the normal path. We change the state to SUCCESS and commit. The state and the events are saved together in one transaction.

  2. The state is already SUCCESS: another request won the race and already saved its own events. We roll back so the events are not saved twice, and return 200.

  3. The state is something else, such as FAILED: the task is no longer running, so the update is rejected with 409 and the transaction is rolled back. No events from this request are saved.

This locked re-check keeps the task state and asset events atomic: the events are committed only together with the SUCCESS state.

Benchmark

The task_instance row lock is held during one successful task completion, on Postgres. The task emits to N assets to make the registration slow. Median of 14 runs (lock holding time):

outletsbeforeafter
300426 ms2.0 ms
600974 ms2.1 ms

Before, the lock is held for the whole registration, so it grows with the number of outlets. After, the registration runs before the lock, so the lock is held only for the state update and stays flat.

This is an alternative to #66854.

related: #66853


Was generative AI tooling used to co-author this PR?
  • Yes (please specify the tool below)

Generated-by: Claude Code Opus 4.8 following the guidelines


  • Read the Pull Request Guidelines for more information. Note: commit author/co-author name and email in commits become permanently public when merged.
  • For fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
  • When adding dependency, check compliance with the ASF 3rd Party License Policy.
  • For significant user-facing changes create newsfragment: {pr_number}.significant.rst, in airflow-core/newsfragments. You can add this file in a follow-up commit after the PR is created so you know the PR number.

Under high fan-out, ti_update_state held the task_instance row lock for
the whole asset-event registration (asset lookups, event inserts,
dag-run queueing). Concurrent completions piled up on the lock, each
occupying an API-server thread, until the server was OOMKilled (apache#66853).
Register the events before acquiring the lock instead, still inside the
same transaction. The locked SELECT then re-checks the state: on the
normal path the state flip and the events commit atomically; if a
duplicate or a concurrent failure wins the race, the transaction rolls
back and the events are discarded with it. Nothing is ever registered
without its SUCCESS state, so no queue or reconciliation is needed.
@Andrushika

Copy link
Copy Markdown
ContributorAuthor

Close this since #70951 is merged.

@Andrushika
Andrushika deleted the reorder-asset-registration-before-ti-lock branch August 4, 2026 09:16
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.

1 participant

@Andrushika
, '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

Register asset events before taking the task_instance row lock - #70078

Closed
Andrushika wants to merge 1 commit into
apache:mainfrom
Andrushika:reorder-asset-registration-before-ti-lock
Closed

Register asset events before taking the task_instance row lock#70078
Andrushika wants to merge 1 commit into
apache:mainfrom
Andrushika:reorder-asset-registration-before-ti-lock

Conversation

@Andrushika

@AndrushikaAndrushika commented Jul 18, 2026

Copy link
Copy Markdown
Contributor

Register asset events before taking the task_instance row lock

Why

When a task succeeds, ti_update_state currently takes the task_instance row lock and then registers asset events while holding the lock.

Asset registration can be slow because it performs asset lookups, event inserts, and Dag-run queueing. When many tasks finish at the same time, concurrent requests can wait for the row lock while each request continues to use an API-server thread. In the case reported in #66853, this caused the API server to run out of memory and get OOMKilled.

What

This change registers the asset events before taking the task_instance row lock:
image
The asset events and the task state are still written in the same database transaction.

In this PR, we first read the state without the row lock, then perform the asset registration. This can hit a TOCTOU problem: the state may change in between, for example, when the worker retries. So before we actually change the state, we take the row lock and read the state one more time. The lock makes sure no one else can change it after this point. Then we decide what to do.

There are three cases:

  1. The state is still RUNNING: the normal path. We change the state to SUCCESS and commit. The state and the events are saved together in one transaction.

  2. The state is already SUCCESS: another request won the race and already saved its own events. We roll back so the events are not saved twice, and return 200.

  3. The state is something else, such as FAILED: the task is no longer running, so the update is rejected with 409 and the transaction is rolled back. No events from this request are saved.

This locked re-check keeps the task state and asset events atomic: the events are committed only together with the SUCCESS state.

Benchmark

The task_instance row lock is held during one successful task completion, on Postgres. The task emits to N assets to make the registration slow. Median of 14 runs (lock holding time):

outletsbeforeafter
300426 ms2.0 ms
600974 ms2.1 ms

Before, the lock is held for the whole registration, so it grows with the number of outlets. After, the registration runs before the lock, so the lock is held only for the state update and stays flat.

This is an alternative to #66854.

related: #66853


Was generative AI tooling used to co-author this PR?
  • Yes (please specify the tool below)

Generated-by: Claude Code Opus 4.8 following the guidelines


  • Read the Pull Request Guidelines for more information. Note: commit author/co-author name and email in commits become permanently public when merged.
  • For fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
  • When adding dependency, check compliance with the ASF 3rd Party License Policy.
  • For significant user-facing changes create newsfragment: {pr_number}.significant.rst, in airflow-core/newsfragments. You can add this file in a follow-up commit after the PR is created so you know the PR number.

Under high fan-out, ti_update_state held the task_instance row lock for
the whole asset-event registration (asset lookups, event inserts,
dag-run queueing). Concurrent completions piled up on the lock, each
occupying an API-server thread, until the server was OOMKilled (apache#66853).
Register the events before acquiring the lock instead, still inside the
same transaction. The locked SELECT then re-checks the state: on the
normal path the state flip and the events commit atomically; if a
duplicate or a concurrent failure wins the race, the transaction rolls
back and the events are discarded with it. Nothing is ever registered
without its SUCCESS state, so no queue or reconciliation is needed.
@Andrushika

Copy link
Copy Markdown
ContributorAuthor

Close this since #70951 is merged.

@Andrushika
Andrushika deleted the reorder-asset-registration-before-ti-lock branch August 4, 2026 09:16
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.

1 participant

@Andrushika
, '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

Register asset events before taking the task_instance row lock - #70078

Closed
Andrushika wants to merge 1 commit into
apache:mainfrom
Andrushika:reorder-asset-registration-before-ti-lock
Closed

Register asset events before taking the task_instance row lock#70078
Andrushika wants to merge 1 commit into
apache:mainfrom
Andrushika:reorder-asset-registration-before-ti-lock

Conversation

@Andrushika

@AndrushikaAndrushika commented Jul 18, 2026

Copy link
Copy Markdown
Contributor

Register asset events before taking the task_instance row lock

Why

When a task succeeds, ti_update_state currently takes the task_instance row lock and then registers asset events while holding the lock.

Asset registration can be slow because it performs asset lookups, event inserts, and Dag-run queueing. When many tasks finish at the same time, concurrent requests can wait for the row lock while each request continues to use an API-server thread. In the case reported in #66853, this caused the API server to run out of memory and get OOMKilled.

What

This change registers the asset events before taking the task_instance row lock:
image
The asset events and the task state are still written in the same database transaction.

In this PR, we first read the state without the row lock, then perform the asset registration. This can hit a TOCTOU problem: the state may change in between, for example, when the worker retries. So before we actually change the state, we take the row lock and read the state one more time. The lock makes sure no one else can change it after this point. Then we decide what to do.

There are three cases:

  1. The state is still RUNNING: the normal path. We change the state to SUCCESS and commit. The state and the events are saved together in one transaction.

  2. The state is already SUCCESS: another request won the race and already saved its own events. We roll back so the events are not saved twice, and return 200.

  3. The state is something else, such as FAILED: the task is no longer running, so the update is rejected with 409 and the transaction is rolled back. No events from this request are saved.

This locked re-check keeps the task state and asset events atomic: the events are committed only together with the SUCCESS state.

Benchmark

The task_instance row lock is held during one successful task completion, on Postgres. The task emits to N assets to make the registration slow. Median of 14 runs (lock holding time):

outletsbeforeafter
300426 ms2.0 ms
600974 ms2.1 ms

Before, the lock is held for the whole registration, so it grows with the number of outlets. After, the registration runs before the lock, so the lock is held only for the state update and stays flat.

This is an alternative to #66854.

related: #66853


Was generative AI tooling used to co-author this PR?
  • Yes (please specify the tool below)

Generated-by: Claude Code Opus 4.8 following the guidelines


  • Read the Pull Request Guidelines for more information. Note: commit author/co-author name and email in commits become permanently public when merged.
  • For fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
  • When adding dependency, check compliance with the ASF 3rd Party License Policy.
  • For significant user-facing changes create newsfragment: {pr_number}.significant.rst, in airflow-core/newsfragments. You can add this file in a follow-up commit after the PR is created so you know the PR number.

Under high fan-out, ti_update_state held the task_instance row lock for
the whole asset-event registration (asset lookups, event inserts,
dag-run queueing). Concurrent completions piled up on the lock, each
occupying an API-server thread, until the server was OOMKilled (apache#66853).
Register the events before acquiring the lock instead, still inside the
same transaction. The locked SELECT then re-checks the state: on the
normal path the state flip and the events commit atomically; if a
duplicate or a concurrent failure wins the race, the transaction rolls
back and the events are discarded with it. Nothing is ever registered
without its SUCCESS state, so no queue or reconciliation is needed.
@Andrushika

Copy link
Copy Markdown
ContributorAuthor

Close this since #70951 is merged.

@Andrushika
Andrushika deleted the reorder-asset-registration-before-ti-lock branch August 4, 2026 09:16
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.

1 participant

@Andrushika
, '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

Register asset events before taking the task_instance row lock - #70078

Closed
Andrushika wants to merge 1 commit into
apache:mainfrom
Andrushika:reorder-asset-registration-before-ti-lock
Closed

Register asset events before taking the task_instance row lock#70078
Andrushika wants to merge 1 commit into
apache:mainfrom
Andrushika:reorder-asset-registration-before-ti-lock

Conversation

@Andrushika

@AndrushikaAndrushika commented Jul 18, 2026

Copy link
Copy Markdown
Contributor

Register asset events before taking the task_instance row lock

Why

When a task succeeds, ti_update_state currently takes the task_instance row lock and then registers asset events while holding the lock.

Asset registration can be slow because it performs asset lookups, event inserts, and Dag-run queueing. When many tasks finish at the same time, concurrent requests can wait for the row lock while each request continues to use an API-server thread. In the case reported in #66853, this caused the API server to run out of memory and get OOMKilled.

What

This change registers the asset events before taking the task_instance row lock:
image
The asset events and the task state are still written in the same database transaction.

In this PR, we first read the state without the row lock, then perform the asset registration. This can hit a TOCTOU problem: the state may change in between, for example, when the worker retries. So before we actually change the state, we take the row lock and read the state one more time. The lock makes sure no one else can change it after this point. Then we decide what to do.

There are three cases:

  1. The state is still RUNNING: the normal path. We change the state to SUCCESS and commit. The state and the events are saved together in one transaction.

  2. The state is already SUCCESS: another request won the race and already saved its own events. We roll back so the events are not saved twice, and return 200.

  3. The state is something else, such as FAILED: the task is no longer running, so the update is rejected with 409 and the transaction is rolled back. No events from this request are saved.

This locked re-check keeps the task state and asset events atomic: the events are committed only together with the SUCCESS state.

Benchmark

The task_instance row lock is held during one successful task completion, on Postgres. The task emits to N assets to make the registration slow. Median of 14 runs (lock holding time):

outletsbeforeafter
300426 ms2.0 ms
600974 ms2.1 ms

Before, the lock is held for the whole registration, so it grows with the number of outlets. After, the registration runs before the lock, so the lock is held only for the state update and stays flat.

This is an alternative to #66854.

related: #66853


Was generative AI tooling used to co-author this PR?
  • Yes (please specify the tool below)

Generated-by: Claude Code Opus 4.8 following the guidelines


  • Read the Pull Request Guidelines for more information. Note: commit author/co-author name and email in commits become permanently public when merged.
  • For fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
  • When adding dependency, check compliance with the ASF 3rd Party License Policy.
  • For significant user-facing changes create newsfragment: {pr_number}.significant.rst, in airflow-core/newsfragments. You can add this file in a follow-up commit after the PR is created so you know the PR number.

Under high fan-out, ti_update_state held the task_instance row lock for
the whole asset-event registration (asset lookups, event inserts,
dag-run queueing). Concurrent completions piled up on the lock, each
occupying an API-server thread, until the server was OOMKilled (apache#66853).
Register the events before acquiring the lock instead, still inside the
same transaction. The locked SELECT then re-checks the state: on the
normal path the state flip and the events commit atomically; if a
duplicate or a concurrent failure wins the race, the transaction rolls
back and the events are discarded with it. Nothing is ever registered
without its SUCCESS state, so no queue or reconciliation is needed.
@Andrushika

Copy link
Copy Markdown
ContributorAuthor

Close this since #70951 is merged.

@Andrushika
Andrushika deleted the reorder-asset-registration-before-ti-lock branch August 4, 2026 09:16
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.

1 participant

@Andrushika
, '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

Register asset events before taking the task_instance row lock - #70078

Closed
Andrushika wants to merge 1 commit into
apache:mainfrom
Andrushika:reorder-asset-registration-before-ti-lock
Closed

Register asset events before taking the task_instance row lock#70078
Andrushika wants to merge 1 commit into
apache:mainfrom
Andrushika:reorder-asset-registration-before-ti-lock

Conversation

@Andrushika

@AndrushikaAndrushika commented Jul 18, 2026

Copy link
Copy Markdown
Contributor

Register asset events before taking the task_instance row lock

Why

When a task succeeds, ti_update_state currently takes the task_instance row lock and then registers asset events while holding the lock.

Asset registration can be slow because it performs asset lookups, event inserts, and Dag-run queueing. When many tasks finish at the same time, concurrent requests can wait for the row lock while each request continues to use an API-server thread. In the case reported in #66853, this caused the API server to run out of memory and get OOMKilled.

What

This change registers the asset events before taking the task_instance row lock:
image
The asset events and the task state are still written in the same database transaction.

In this PR, we first read the state without the row lock, then perform the asset registration. This can hit a TOCTOU problem: the state may change in between, for example, when the worker retries. So before we actually change the state, we take the row lock and read the state one more time. The lock makes sure no one else can change it after this point. Then we decide what to do.

There are three cases:

  1. The state is still RUNNING: the normal path. We change the state to SUCCESS and commit. The state and the events are saved together in one transaction.

  2. The state is already SUCCESS: another request won the race and already saved its own events. We roll back so the events are not saved twice, and return 200.

  3. The state is something else, such as FAILED: the task is no longer running, so the update is rejected with 409 and the transaction is rolled back. No events from this request are saved.

This locked re-check keeps the task state and asset events atomic: the events are committed only together with the SUCCESS state.

Benchmark

The task_instance row lock is held during one successful task completion, on Postgres. The task emits to N assets to make the registration slow. Median of 14 runs (lock holding time):

outletsbeforeafter
300426 ms2.0 ms
600974 ms2.1 ms

Before, the lock is held for the whole registration, so it grows with the number of outlets. After, the registration runs before the lock, so the lock is held only for the state update and stays flat.

This is an alternative to #66854.

related: #66853


Was generative AI tooling used to co-author this PR?
  • Yes (please specify the tool below)

Generated-by: Claude Code Opus 4.8 following the guidelines


  • Read the Pull Request Guidelines for more information. Note: commit author/co-author name and email in commits become permanently public when merged.
  • For fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
  • When adding dependency, check compliance with the ASF 3rd Party License Policy.
  • For significant user-facing changes create newsfragment: {pr_number}.significant.rst, in airflow-core/newsfragments. You can add this file in a follow-up commit after the PR is created so you know the PR number.

Under high fan-out, ti_update_state held the task_instance row lock for
the whole asset-event registration (asset lookups, event inserts,
dag-run queueing). Concurrent completions piled up on the lock, each
occupying an API-server thread, until the server was OOMKilled (apache#66853).
Register the events before acquiring the lock instead, still inside the
same transaction. The locked SELECT then re-checks the state: on the
normal path the state flip and the events commit atomically; if a
duplicate or a concurrent failure wins the race, the transaction rolls
back and the events are discarded with it. Nothing is ever registered
without its SUCCESS state, so no queue or reconciliation is needed.
@Andrushika

Copy link
Copy Markdown
ContributorAuthor

Close this since #70951 is merged.

@Andrushika
Andrushika deleted the reorder-asset-registration-before-ti-lock branch August 4, 2026 09:16
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.

1 participant

@Andrushika
, '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

Register asset events before taking the task_instance row lock - #70078

Closed
Andrushika wants to merge 1 commit into
apache:mainfrom
Andrushika:reorder-asset-registration-before-ti-lock
Closed

Register asset events before taking the task_instance row lock#70078
Andrushika wants to merge 1 commit into
apache:mainfrom
Andrushika:reorder-asset-registration-before-ti-lock

Conversation

@Andrushika

@AndrushikaAndrushika commented Jul 18, 2026

Copy link
Copy Markdown
Contributor

Register asset events before taking the task_instance row lock

Why

When a task succeeds, ti_update_state currently takes the task_instance row lock and then registers asset events while holding the lock.

Asset registration can be slow because it performs asset lookups, event inserts, and Dag-run queueing. When many tasks finish at the same time, concurrent requests can wait for the row lock while each request continues to use an API-server thread. In the case reported in #66853, this caused the API server to run out of memory and get OOMKilled.

What

This change registers the asset events before taking the task_instance row lock:
image
The asset events and the task state are still written in the same database transaction.

In this PR, we first read the state without the row lock, then perform the asset registration. This can hit a TOCTOU problem: the state may change in between, for example, when the worker retries. So before we actually change the state, we take the row lock and read the state one more time. The lock makes sure no one else can change it after this point. Then we decide what to do.

There are three cases:

  1. The state is still RUNNING: the normal path. We change the state to SUCCESS and commit. The state and the events are saved together in one transaction.

  2. The state is already SUCCESS: another request won the race and already saved its own events. We roll back so the events are not saved twice, and return 200.

  3. The state is something else, such as FAILED: the task is no longer running, so the update is rejected with 409 and the transaction is rolled back. No events from this request are saved.

This locked re-check keeps the task state and asset events atomic: the events are committed only together with the SUCCESS state.

Benchmark

The task_instance row lock is held during one successful task completion, on Postgres. The task emits to N assets to make the registration slow. Median of 14 runs (lock holding time):

outletsbeforeafter
300426 ms2.0 ms
600974 ms2.1 ms

Before, the lock is held for the whole registration, so it grows with the number of outlets. After, the registration runs before the lock, so the lock is held only for the state update and stays flat.

This is an alternative to #66854.

related: #66853


Was generative AI tooling used to co-author this PR?
  • Yes (please specify the tool below)

Generated-by: Claude Code Opus 4.8 following the guidelines


  • Read the Pull Request Guidelines for more information. Note: commit author/co-author name and email in commits become permanently public when merged.
  • For fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
  • When adding dependency, check compliance with the ASF 3rd Party License Policy.
  • For significant user-facing changes create newsfragment: {pr_number}.significant.rst, in airflow-core/newsfragments. You can add this file in a follow-up commit after the PR is created so you know the PR number.

Under high fan-out, ti_update_state held the task_instance row lock for
the whole asset-event registration (asset lookups, event inserts,
dag-run queueing). Concurrent completions piled up on the lock, each
occupying an API-server thread, until the server was OOMKilled (apache#66853).
Register the events before acquiring the lock instead, still inside the
same transaction. The locked SELECT then re-checks the state: on the
normal path the state flip and the events commit atomically; if a
duplicate or a concurrent failure wins the race, the transaction rolls
back and the events are discarded with it. Nothing is ever registered
without its SUCCESS state, so no queue or reconciliation is needed.
@Andrushika

Copy link
Copy Markdown
ContributorAuthor

Close this since #70951 is merged.

@Andrushika
Andrushika deleted the reorder-asset-registration-before-ti-lock branch August 4, 2026 09:16
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.

1 participant

@Andrushika
, '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

Register asset events before taking the task_instance row lock - #70078

Closed
Andrushika wants to merge 1 commit into
apache:mainfrom
Andrushika:reorder-asset-registration-before-ti-lock
Closed

Register asset events before taking the task_instance row lock#70078
Andrushika wants to merge 1 commit into
apache:mainfrom
Andrushika:reorder-asset-registration-before-ti-lock

Conversation

@Andrushika

@AndrushikaAndrushika commented Jul 18, 2026

Copy link
Copy Markdown
Contributor

Register asset events before taking the task_instance row lock

Why

When a task succeeds, ti_update_state currently takes the task_instance row lock and then registers asset events while holding the lock.

Asset registration can be slow because it performs asset lookups, event inserts, and Dag-run queueing. When many tasks finish at the same time, concurrent requests can wait for the row lock while each request continues to use an API-server thread. In the case reported in #66853, this caused the API server to run out of memory and get OOMKilled.

What

This change registers the asset events before taking the task_instance row lock:
image
The asset events and the task state are still written in the same database transaction.

In this PR, we first read the state without the row lock, then perform the asset registration. This can hit a TOCTOU problem: the state may change in between, for example, when the worker retries. So before we actually change the state, we take the row lock and read the state one more time. The lock makes sure no one else can change it after this point. Then we decide what to do.

There are three cases:

  1. The state is still RUNNING: the normal path. We change the state to SUCCESS and commit. The state and the events are saved together in one transaction.

  2. The state is already SUCCESS: another request won the race and already saved its own events. We roll back so the events are not saved twice, and return 200.

  3. The state is something else, such as FAILED: the task is no longer running, so the update is rejected with 409 and the transaction is rolled back. No events from this request are saved.

This locked re-check keeps the task state and asset events atomic: the events are committed only together with the SUCCESS state.

Benchmark

The task_instance row lock is held during one successful task completion, on Postgres. The task emits to N assets to make the registration slow. Median of 14 runs (lock holding time):

outletsbeforeafter
300426 ms2.0 ms
600974 ms2.1 ms

Before, the lock is held for the whole registration, so it grows with the number of outlets. After, the registration runs before the lock, so the lock is held only for the state update and stays flat.

This is an alternative to #66854.

related: #66853


Was generative AI tooling used to co-author this PR?
  • Yes (please specify the tool below)

Generated-by: Claude Code Opus 4.8 following the guidelines


  • Read the Pull Request Guidelines for more information. Note: commit author/co-author name and email in commits become permanently public when merged.
  • For fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
  • When adding dependency, check compliance with the ASF 3rd Party License Policy.
  • For significant user-facing changes create newsfragment: {pr_number}.significant.rst, in airflow-core/newsfragments. You can add this file in a follow-up commit after the PR is created so you know the PR number.

Under high fan-out, ti_update_state held the task_instance row lock for
the whole asset-event registration (asset lookups, event inserts,
dag-run queueing). Concurrent completions piled up on the lock, each
occupying an API-server thread, until the server was OOMKilled (apache#66853).
Register the events before acquiring the lock instead, still inside the
same transaction. The locked SELECT then re-checks the state: on the
normal path the state flip and the events commit atomically; if a
duplicate or a concurrent failure wins the race, the transaction rolls
back and the events are discarded with it. Nothing is ever registered
without its SUCCESS state, so no queue or reconciliation is needed.
@Andrushika

Copy link
Copy Markdown
ContributorAuthor

Close this since #70951 is merged.

@Andrushika
Andrushika deleted the reorder-asset-registration-before-ti-lock branch August 4, 2026 09:16
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.

1 participant

@Andrushika
, '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

Register asset events before taking the task_instance row lock - #70078

Closed
Andrushika wants to merge 1 commit into
apache:mainfrom
Andrushika:reorder-asset-registration-before-ti-lock
Closed

Register asset events before taking the task_instance row lock#70078
Andrushika wants to merge 1 commit into
apache:mainfrom
Andrushika:reorder-asset-registration-before-ti-lock

Conversation

@Andrushika

@AndrushikaAndrushika commented Jul 18, 2026

Copy link
Copy Markdown
Contributor

Register asset events before taking the task_instance row lock

Why

When a task succeeds, ti_update_state currently takes the task_instance row lock and then registers asset events while holding the lock.

Asset registration can be slow because it performs asset lookups, event inserts, and Dag-run queueing. When many tasks finish at the same time, concurrent requests can wait for the row lock while each request continues to use an API-server thread. In the case reported in #66853, this caused the API server to run out of memory and get OOMKilled.

What

This change registers the asset events before taking the task_instance row lock:
image
The asset events and the task state are still written in the same database transaction.

In this PR, we first read the state without the row lock, then perform the asset registration. This can hit a TOCTOU problem: the state may change in between, for example, when the worker retries. So before we actually change the state, we take the row lock and read the state one more time. The lock makes sure no one else can change it after this point. Then we decide what to do.

There are three cases:

  1. The state is still RUNNING: the normal path. We change the state to SUCCESS and commit. The state and the events are saved together in one transaction.

  2. The state is already SUCCESS: another request won the race and already saved its own events. We roll back so the events are not saved twice, and return 200.

  3. The state is something else, such as FAILED: the task is no longer running, so the update is rejected with 409 and the transaction is rolled back. No events from this request are saved.

This locked re-check keeps the task state and asset events atomic: the events are committed only together with the SUCCESS state.

Benchmark

The task_instance row lock is held during one successful task completion, on Postgres. The task emits to N assets to make the registration slow. Median of 14 runs (lock holding time):

outletsbeforeafter
300426 ms2.0 ms
600974 ms2.1 ms

Before, the lock is held for the whole registration, so it grows with the number of outlets. After, the registration runs before the lock, so the lock is held only for the state update and stays flat.

This is an alternative to #66854.

related: #66853


Was generative AI tooling used to co-author this PR?
  • Yes (please specify the tool below)

Generated-by: Claude Code Opus 4.8 following the guidelines


  • Read the Pull Request Guidelines for more information. Note: commit author/co-author name and email in commits become permanently public when merged.
  • For fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
  • When adding dependency, check compliance with the ASF 3rd Party License Policy.
  • For significant user-facing changes create newsfragment: {pr_number}.significant.rst, in airflow-core/newsfragments. You can add this file in a follow-up commit after the PR is created so you know the PR number.

Under high fan-out, ti_update_state held the task_instance row lock for
the whole asset-event registration (asset lookups, event inserts,
dag-run queueing). Concurrent completions piled up on the lock, each
occupying an API-server thread, until the server was OOMKilled (apache#66853).
Register the events before acquiring the lock instead, still inside the
same transaction. The locked SELECT then re-checks the state: on the
normal path the state flip and the events commit atomically; if a
duplicate or a concurrent failure wins the race, the transaction rolls
back and the events are discarded with it. Nothing is ever registered
without its SUCCESS state, so no queue or reconciliation is needed.
@Andrushika

Copy link
Copy Markdown
ContributorAuthor

Close this since #70951 is merged.

@Andrushika
Andrushika deleted the reorder-asset-registration-before-ti-lock branch August 4, 2026 09:16
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.

1 participant

@Andrushika