Consume every queued asset event from concurrent mapped outlets - #71994

Closed
Vamsi-klu wants to merge 2 commits into
apache:mainfrom
Vamsi-klu:pr1/54659-mapped-outlet-events
Closed

Consume every queued asset event from concurrent mapped outlets#71994
Vamsi-klu wants to merge 2 commits into
apache:mainfrom
Vamsi-klu:pr1/54659-mapped-outlet-events

Conversation

@Vamsi-klu

@Vamsi-kluVamsi-klu commented Aug 23, 2026

Copy link
Copy Markdown
Contributor

closes: #54659

GitHub still lists a few other pull requests under This was referenced. I named them in an early draft of this writeup and then deleted those lines. GitHub keeps the timeline row even after the text is gone. Those pull requests are not part of this change. The only real link is the issue this PR is meant to fix.

What I did

Tests only. Mapped producer @task(outlets=[asset]).expand(...) succeeds, then one scheduler tick. Assert all N events land on the consumer run and in triggering_asset_events. No production change. No newsfragment.

Why I did

Three mapped outlets were finishing together and the consumer only saw a subset in triggering_asset_events. Consume-by-event-id is already on main. Existing tests still insert ADRQ by hand, so they never covered this emit path.

How I did

Emit: dag_maker.run_ti(..., map_index=N)register_asset_changes_in_db.
Consume: SchedulerJobRunner._create_dagruns_for_dags.
Context: get_template_context after loading consumed_asset_events.

One tick batches all visible events. That is the default, not a bug. Leftovers stay for the next tick.

What's the impact

None at runtime. If the mapped emit path regresses, these tests fail instead of silently dropping events.

What's the testing

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

  • test_mapped_outlet_asset_events_consumed_in_one_tick
  • test_mapped_outlet_asset_events_consumed_across_staggered_ticks
  • test_mapped_outlet_asset_events_same_timestamp_are_all_consumed
  • test_mapped_outlet_asset_events_and_condition_waits_for_all_assets
  • test_mapped_outlet_asset_alias_events_are_all_consumed
uv run --project airflow-core pytest airflow-core/tests/unit/jobs/test_scheduler_job.py \
-k test_mapped_outlet -q

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

Generated-by: Grok 4.6 following the guidelines

apache#70972 already consumes queued asset events by id, so concurrent mapped
outlets no longer strand events behind a timestamp watermark. These
tests lock that contract for mapped producers so a later consume-path
change cannot silently drop events from triggering_asset_events.
closes: apache#54659
The same-timestamp and AND cases only pinned consume if events were
already queued by hand. Running mapped TIs through register/queue
keeps those pins on the user path, including alias context lookup.
closes: apache#54659

@Vamsi-kluVamsi-klu left a comment

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Description review:

  • closes: #54659 is correct.
  • Tests-only is stated; no newsfragment is correct.
  • Default batching in one tick is called out as intended, not a bug.
  • Test names and the pytest command match the diff.

Ready for maintainer review.


Drafted-by: Cursor Grok 4.6 (no human review before posting)

@Vamsi-klu

Vamsi-klu commented Aug 23, 2026

Copy link
Copy Markdown
ContributorAuthor

The extra pull requests under This was referenced are leftover from an earlier draft of the description. I put those names in, then took them out. GitHub keeps the timeline row anyway. Nothing in this PR depends on them. The only real link is the issue this PR is meant to fix.


Drafted-by: Cursor Grok 4.6

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

Labels

area:Schedulerincluding HA (high availability) scheduler

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Asset-triggered DAGs miss events from concurrently completing dynamic mapped tasks

1 participant

@Vamsi-klu
, '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

Consume every queued asset event from concurrent mapped outlets - #71994

Closed
Vamsi-klu wants to merge 2 commits into
apache:mainfrom
Vamsi-klu:pr1/54659-mapped-outlet-events
Closed

Consume every queued asset event from concurrent mapped outlets#71994
Vamsi-klu wants to merge 2 commits into
apache:mainfrom
Vamsi-klu:pr1/54659-mapped-outlet-events

Conversation

@Vamsi-klu

@Vamsi-kluVamsi-klu commented Aug 23, 2026

Copy link
Copy Markdown
Contributor

closes: #54659

GitHub still lists a few other pull requests under This was referenced. I named them in an early draft of this writeup and then deleted those lines. GitHub keeps the timeline row even after the text is gone. Those pull requests are not part of this change. The only real link is the issue this PR is meant to fix.

What I did

Tests only. Mapped producer @task(outlets=[asset]).expand(...) succeeds, then one scheduler tick. Assert all N events land on the consumer run and in triggering_asset_events. No production change. No newsfragment.

Why I did

Three mapped outlets were finishing together and the consumer only saw a subset in triggering_asset_events. Consume-by-event-id is already on main. Existing tests still insert ADRQ by hand, so they never covered this emit path.

How I did

Emit: dag_maker.run_ti(..., map_index=N)register_asset_changes_in_db.
Consume: SchedulerJobRunner._create_dagruns_for_dags.
Context: get_template_context after loading consumed_asset_events.

One tick batches all visible events. That is the default, not a bug. Leftovers stay for the next tick.

What's the impact

None at runtime. If the mapped emit path regresses, these tests fail instead of silently dropping events.

What's the testing

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

  • test_mapped_outlet_asset_events_consumed_in_one_tick
  • test_mapped_outlet_asset_events_consumed_across_staggered_ticks
  • test_mapped_outlet_asset_events_same_timestamp_are_all_consumed
  • test_mapped_outlet_asset_events_and_condition_waits_for_all_assets
  • test_mapped_outlet_asset_alias_events_are_all_consumed
uv run --project airflow-core pytest airflow-core/tests/unit/jobs/test_scheduler_job.py \
-k test_mapped_outlet -q

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

Generated-by: Grok 4.6 following the guidelines

apache#70972 already consumes queued asset events by id, so concurrent mapped
outlets no longer strand events behind a timestamp watermark. These
tests lock that contract for mapped producers so a later consume-path
change cannot silently drop events from triggering_asset_events.
closes: apache#54659
The same-timestamp and AND cases only pinned consume if events were
already queued by hand. Running mapped TIs through register/queue
keeps those pins on the user path, including alias context lookup.
closes: apache#54659

@Vamsi-kluVamsi-klu left a comment

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Description review:

  • closes: #54659 is correct.
  • Tests-only is stated; no newsfragment is correct.
  • Default batching in one tick is called out as intended, not a bug.
  • Test names and the pytest command match the diff.

Ready for maintainer review.


Drafted-by: Cursor Grok 4.6 (no human review before posting)

@Vamsi-klu

Vamsi-klu commented Aug 23, 2026

Copy link
Copy Markdown
ContributorAuthor

The extra pull requests under This was referenced are leftover from an earlier draft of the description. I put those names in, then took them out. GitHub keeps the timeline row anyway. Nothing in this PR depends on them. The only real link is the issue this PR is meant to fix.


Drafted-by: Cursor Grok 4.6

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

Labels

area:Schedulerincluding HA (high availability) scheduler

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Asset-triggered DAGs miss events from concurrently completing dynamic mapped tasks

1 participant

@Vamsi-klu
, '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

Consume every queued asset event from concurrent mapped outlets - #71994

Closed
Vamsi-klu wants to merge 2 commits into
apache:mainfrom
Vamsi-klu:pr1/54659-mapped-outlet-events
Closed

Consume every queued asset event from concurrent mapped outlets#71994
Vamsi-klu wants to merge 2 commits into
apache:mainfrom
Vamsi-klu:pr1/54659-mapped-outlet-events

Conversation

@Vamsi-klu

@Vamsi-kluVamsi-klu commented Aug 23, 2026

Copy link
Copy Markdown
Contributor

closes: #54659

GitHub still lists a few other pull requests under This was referenced. I named them in an early draft of this writeup and then deleted those lines. GitHub keeps the timeline row even after the text is gone. Those pull requests are not part of this change. The only real link is the issue this PR is meant to fix.

What I did

Tests only. Mapped producer @task(outlets=[asset]).expand(...) succeeds, then one scheduler tick. Assert all N events land on the consumer run and in triggering_asset_events. No production change. No newsfragment.

Why I did

Three mapped outlets were finishing together and the consumer only saw a subset in triggering_asset_events. Consume-by-event-id is already on main. Existing tests still insert ADRQ by hand, so they never covered this emit path.

How I did

Emit: dag_maker.run_ti(..., map_index=N)register_asset_changes_in_db.
Consume: SchedulerJobRunner._create_dagruns_for_dags.
Context: get_template_context after loading consumed_asset_events.

One tick batches all visible events. That is the default, not a bug. Leftovers stay for the next tick.

What's the impact

None at runtime. If the mapped emit path regresses, these tests fail instead of silently dropping events.

What's the testing

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

  • test_mapped_outlet_asset_events_consumed_in_one_tick
  • test_mapped_outlet_asset_events_consumed_across_staggered_ticks
  • test_mapped_outlet_asset_events_same_timestamp_are_all_consumed
  • test_mapped_outlet_asset_events_and_condition_waits_for_all_assets
  • test_mapped_outlet_asset_alias_events_are_all_consumed
uv run --project airflow-core pytest airflow-core/tests/unit/jobs/test_scheduler_job.py \
-k test_mapped_outlet -q

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

Generated-by: Grok 4.6 following the guidelines

apache#70972 already consumes queued asset events by id, so concurrent mapped
outlets no longer strand events behind a timestamp watermark. These
tests lock that contract for mapped producers so a later consume-path
change cannot silently drop events from triggering_asset_events.
closes: apache#54659
The same-timestamp and AND cases only pinned consume if events were
already queued by hand. Running mapped TIs through register/queue
keeps those pins on the user path, including alias context lookup.
closes: apache#54659

@Vamsi-kluVamsi-klu left a comment

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Description review:

  • closes: #54659 is correct.
  • Tests-only is stated; no newsfragment is correct.
  • Default batching in one tick is called out as intended, not a bug.
  • Test names and the pytest command match the diff.

Ready for maintainer review.


Drafted-by: Cursor Grok 4.6 (no human review before posting)

@Vamsi-klu

Vamsi-klu commented Aug 23, 2026

Copy link
Copy Markdown
ContributorAuthor

The extra pull requests under This was referenced are leftover from an earlier draft of the description. I put those names in, then took them out. GitHub keeps the timeline row anyway. Nothing in this PR depends on them. The only real link is the issue this PR is meant to fix.


Drafted-by: Cursor Grok 4.6

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

Labels

area:Schedulerincluding HA (high availability) scheduler

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Asset-triggered DAGs miss events from concurrently completing dynamic mapped tasks

1 participant

@Vamsi-klu
, '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

Consume every queued asset event from concurrent mapped outlets - #71994

Closed
Vamsi-klu wants to merge 2 commits into
apache:mainfrom
Vamsi-klu:pr1/54659-mapped-outlet-events
Closed

Consume every queued asset event from concurrent mapped outlets#71994
Vamsi-klu wants to merge 2 commits into
apache:mainfrom
Vamsi-klu:pr1/54659-mapped-outlet-events

Conversation

@Vamsi-klu

@Vamsi-kluVamsi-klu commented Aug 23, 2026

Copy link
Copy Markdown
Contributor

closes: #54659

GitHub still lists a few other pull requests under This was referenced. I named them in an early draft of this writeup and then deleted those lines. GitHub keeps the timeline row even after the text is gone. Those pull requests are not part of this change. The only real link is the issue this PR is meant to fix.

What I did

Tests only. Mapped producer @task(outlets=[asset]).expand(...) succeeds, then one scheduler tick. Assert all N events land on the consumer run and in triggering_asset_events. No production change. No newsfragment.

Why I did

Three mapped outlets were finishing together and the consumer only saw a subset in triggering_asset_events. Consume-by-event-id is already on main. Existing tests still insert ADRQ by hand, so they never covered this emit path.

How I did

Emit: dag_maker.run_ti(..., map_index=N)register_asset_changes_in_db.
Consume: SchedulerJobRunner._create_dagruns_for_dags.
Context: get_template_context after loading consumed_asset_events.

One tick batches all visible events. That is the default, not a bug. Leftovers stay for the next tick.

What's the impact

None at runtime. If the mapped emit path regresses, these tests fail instead of silently dropping events.

What's the testing

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

  • test_mapped_outlet_asset_events_consumed_in_one_tick
  • test_mapped_outlet_asset_events_consumed_across_staggered_ticks
  • test_mapped_outlet_asset_events_same_timestamp_are_all_consumed
  • test_mapped_outlet_asset_events_and_condition_waits_for_all_assets
  • test_mapped_outlet_asset_alias_events_are_all_consumed
uv run --project airflow-core pytest airflow-core/tests/unit/jobs/test_scheduler_job.py \
-k test_mapped_outlet -q

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

Generated-by: Grok 4.6 following the guidelines

apache#70972 already consumes queued asset events by id, so concurrent mapped
outlets no longer strand events behind a timestamp watermark. These
tests lock that contract for mapped producers so a later consume-path
change cannot silently drop events from triggering_asset_events.
closes: apache#54659
The same-timestamp and AND cases only pinned consume if events were
already queued by hand. Running mapped TIs through register/queue
keeps those pins on the user path, including alias context lookup.
closes: apache#54659

@Vamsi-kluVamsi-klu left a comment

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Description review:

  • closes: #54659 is correct.
  • Tests-only is stated; no newsfragment is correct.
  • Default batching in one tick is called out as intended, not a bug.
  • Test names and the pytest command match the diff.

Ready for maintainer review.


Drafted-by: Cursor Grok 4.6 (no human review before posting)

@Vamsi-klu

Vamsi-klu commented Aug 23, 2026

Copy link
Copy Markdown
ContributorAuthor

The extra pull requests under This was referenced are leftover from an earlier draft of the description. I put those names in, then took them out. GitHub keeps the timeline row anyway. Nothing in this PR depends on them. The only real link is the issue this PR is meant to fix.


Drafted-by: Cursor Grok 4.6

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

Labels

area:Schedulerincluding HA (high availability) scheduler

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Asset-triggered DAGs miss events from concurrently completing dynamic mapped tasks

1 participant

@Vamsi-klu
, '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

Consume every queued asset event from concurrent mapped outlets - #71994

Closed
Vamsi-klu wants to merge 2 commits into
apache:mainfrom
Vamsi-klu:pr1/54659-mapped-outlet-events
Closed

Consume every queued asset event from concurrent mapped outlets#71994
Vamsi-klu wants to merge 2 commits into
apache:mainfrom
Vamsi-klu:pr1/54659-mapped-outlet-events

Conversation

@Vamsi-klu

@Vamsi-kluVamsi-klu commented Aug 23, 2026

Copy link
Copy Markdown
Contributor

closes: #54659

GitHub still lists a few other pull requests under This was referenced. I named them in an early draft of this writeup and then deleted those lines. GitHub keeps the timeline row even after the text is gone. Those pull requests are not part of this change. The only real link is the issue this PR is meant to fix.

What I did

Tests only. Mapped producer @task(outlets=[asset]).expand(...) succeeds, then one scheduler tick. Assert all N events land on the consumer run and in triggering_asset_events. No production change. No newsfragment.

Why I did

Three mapped outlets were finishing together and the consumer only saw a subset in triggering_asset_events. Consume-by-event-id is already on main. Existing tests still insert ADRQ by hand, so they never covered this emit path.

How I did

Emit: dag_maker.run_ti(..., map_index=N)register_asset_changes_in_db.
Consume: SchedulerJobRunner._create_dagruns_for_dags.
Context: get_template_context after loading consumed_asset_events.

One tick batches all visible events. That is the default, not a bug. Leftovers stay for the next tick.

What's the impact

None at runtime. If the mapped emit path regresses, these tests fail instead of silently dropping events.

What's the testing

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

  • test_mapped_outlet_asset_events_consumed_in_one_tick
  • test_mapped_outlet_asset_events_consumed_across_staggered_ticks
  • test_mapped_outlet_asset_events_same_timestamp_are_all_consumed
  • test_mapped_outlet_asset_events_and_condition_waits_for_all_assets
  • test_mapped_outlet_asset_alias_events_are_all_consumed
uv run --project airflow-core pytest airflow-core/tests/unit/jobs/test_scheduler_job.py \
-k test_mapped_outlet -q

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

Generated-by: Grok 4.6 following the guidelines

apache#70972 already consumes queued asset events by id, so concurrent mapped
outlets no longer strand events behind a timestamp watermark. These
tests lock that contract for mapped producers so a later consume-path
change cannot silently drop events from triggering_asset_events.
closes: apache#54659
The same-timestamp and AND cases only pinned consume if events were
already queued by hand. Running mapped TIs through register/queue
keeps those pins on the user path, including alias context lookup.
closes: apache#54659

@Vamsi-kluVamsi-klu left a comment

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Description review:

  • closes: #54659 is correct.
  • Tests-only is stated; no newsfragment is correct.
  • Default batching in one tick is called out as intended, not a bug.
  • Test names and the pytest command match the diff.

Ready for maintainer review.


Drafted-by: Cursor Grok 4.6 (no human review before posting)

@Vamsi-klu

Vamsi-klu commented Aug 23, 2026

Copy link
Copy Markdown
ContributorAuthor

The extra pull requests under This was referenced are leftover from an earlier draft of the description. I put those names in, then took them out. GitHub keeps the timeline row anyway. Nothing in this PR depends on them. The only real link is the issue this PR is meant to fix.


Drafted-by: Cursor Grok 4.6

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

Labels

area:Schedulerincluding HA (high availability) scheduler

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Asset-triggered DAGs miss events from concurrently completing dynamic mapped tasks

1 participant

@Vamsi-klu
, '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

Consume every queued asset event from concurrent mapped outlets - #71994

Closed
Vamsi-klu wants to merge 2 commits into
apache:mainfrom
Vamsi-klu:pr1/54659-mapped-outlet-events
Closed

Consume every queued asset event from concurrent mapped outlets#71994
Vamsi-klu wants to merge 2 commits into
apache:mainfrom
Vamsi-klu:pr1/54659-mapped-outlet-events

Conversation

@Vamsi-klu

@Vamsi-kluVamsi-klu commented Aug 23, 2026

Copy link
Copy Markdown
Contributor

closes: #54659

GitHub still lists a few other pull requests under This was referenced. I named them in an early draft of this writeup and then deleted those lines. GitHub keeps the timeline row even after the text is gone. Those pull requests are not part of this change. The only real link is the issue this PR is meant to fix.

What I did

Tests only. Mapped producer @task(outlets=[asset]).expand(...) succeeds, then one scheduler tick. Assert all N events land on the consumer run and in triggering_asset_events. No production change. No newsfragment.

Why I did

Three mapped outlets were finishing together and the consumer only saw a subset in triggering_asset_events. Consume-by-event-id is already on main. Existing tests still insert ADRQ by hand, so they never covered this emit path.

How I did

Emit: dag_maker.run_ti(..., map_index=N)register_asset_changes_in_db.
Consume: SchedulerJobRunner._create_dagruns_for_dags.
Context: get_template_context after loading consumed_asset_events.

One tick batches all visible events. That is the default, not a bug. Leftovers stay for the next tick.

What's the impact

None at runtime. If the mapped emit path regresses, these tests fail instead of silently dropping events.

What's the testing

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

  • test_mapped_outlet_asset_events_consumed_in_one_tick
  • test_mapped_outlet_asset_events_consumed_across_staggered_ticks
  • test_mapped_outlet_asset_events_same_timestamp_are_all_consumed
  • test_mapped_outlet_asset_events_and_condition_waits_for_all_assets
  • test_mapped_outlet_asset_alias_events_are_all_consumed
uv run --project airflow-core pytest airflow-core/tests/unit/jobs/test_scheduler_job.py \
-k test_mapped_outlet -q

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

Generated-by: Grok 4.6 following the guidelines

apache#70972 already consumes queued asset events by id, so concurrent mapped
outlets no longer strand events behind a timestamp watermark. These
tests lock that contract for mapped producers so a later consume-path
change cannot silently drop events from triggering_asset_events.
closes: apache#54659
The same-timestamp and AND cases only pinned consume if events were
already queued by hand. Running mapped TIs through register/queue
keeps those pins on the user path, including alias context lookup.
closes: apache#54659

@Vamsi-kluVamsi-klu left a comment

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Description review:

  • closes: #54659 is correct.
  • Tests-only is stated; no newsfragment is correct.
  • Default batching in one tick is called out as intended, not a bug.
  • Test names and the pytest command match the diff.

Ready for maintainer review.


Drafted-by: Cursor Grok 4.6 (no human review before posting)

@Vamsi-klu

Vamsi-klu commented Aug 23, 2026

Copy link
Copy Markdown
ContributorAuthor

The extra pull requests under This was referenced are leftover from an earlier draft of the description. I put those names in, then took them out. GitHub keeps the timeline row anyway. Nothing in this PR depends on them. The only real link is the issue this PR is meant to fix.


Drafted-by: Cursor Grok 4.6

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

Labels

area:Schedulerincluding HA (high availability) scheduler

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Asset-triggered DAGs miss events from concurrently completing dynamic mapped tasks

1 participant

@Vamsi-klu
, '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

Consume every queued asset event from concurrent mapped outlets - #71994

Closed
Vamsi-klu wants to merge 2 commits into
apache:mainfrom
Vamsi-klu:pr1/54659-mapped-outlet-events
Closed

Consume every queued asset event from concurrent mapped outlets#71994
Vamsi-klu wants to merge 2 commits into
apache:mainfrom
Vamsi-klu:pr1/54659-mapped-outlet-events

Conversation

@Vamsi-klu

@Vamsi-kluVamsi-klu commented Aug 23, 2026

Copy link
Copy Markdown
Contributor

closes: #54659

GitHub still lists a few other pull requests under This was referenced. I named them in an early draft of this writeup and then deleted those lines. GitHub keeps the timeline row even after the text is gone. Those pull requests are not part of this change. The only real link is the issue this PR is meant to fix.

What I did

Tests only. Mapped producer @task(outlets=[asset]).expand(...) succeeds, then one scheduler tick. Assert all N events land on the consumer run and in triggering_asset_events. No production change. No newsfragment.

Why I did

Three mapped outlets were finishing together and the consumer only saw a subset in triggering_asset_events. Consume-by-event-id is already on main. Existing tests still insert ADRQ by hand, so they never covered this emit path.

How I did

Emit: dag_maker.run_ti(..., map_index=N)register_asset_changes_in_db.
Consume: SchedulerJobRunner._create_dagruns_for_dags.
Context: get_template_context after loading consumed_asset_events.

One tick batches all visible events. That is the default, not a bug. Leftovers stay for the next tick.

What's the impact

None at runtime. If the mapped emit path regresses, these tests fail instead of silently dropping events.

What's the testing

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

  • test_mapped_outlet_asset_events_consumed_in_one_tick
  • test_mapped_outlet_asset_events_consumed_across_staggered_ticks
  • test_mapped_outlet_asset_events_same_timestamp_are_all_consumed
  • test_mapped_outlet_asset_events_and_condition_waits_for_all_assets
  • test_mapped_outlet_asset_alias_events_are_all_consumed
uv run --project airflow-core pytest airflow-core/tests/unit/jobs/test_scheduler_job.py \
-k test_mapped_outlet -q

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

Generated-by: Grok 4.6 following the guidelines

apache#70972 already consumes queued asset events by id, so concurrent mapped
outlets no longer strand events behind a timestamp watermark. These
tests lock that contract for mapped producers so a later consume-path
change cannot silently drop events from triggering_asset_events.
closes: apache#54659
The same-timestamp and AND cases only pinned consume if events were
already queued by hand. Running mapped TIs through register/queue
keeps those pins on the user path, including alias context lookup.
closes: apache#54659

@Vamsi-kluVamsi-klu left a comment

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Description review:

  • closes: #54659 is correct.
  • Tests-only is stated; no newsfragment is correct.
  • Default batching in one tick is called out as intended, not a bug.
  • Test names and the pytest command match the diff.

Ready for maintainer review.


Drafted-by: Cursor Grok 4.6 (no human review before posting)

@Vamsi-klu

Vamsi-klu commented Aug 23, 2026

Copy link
Copy Markdown
ContributorAuthor

The extra pull requests under This was referenced are leftover from an earlier draft of the description. I put those names in, then took them out. GitHub keeps the timeline row anyway. Nothing in this PR depends on them. The only real link is the issue this PR is meant to fix.


Drafted-by: Cursor Grok 4.6

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

Labels

area:Schedulerincluding HA (high availability) scheduler

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Asset-triggered DAGs miss events from concurrently completing dynamic mapped tasks

1 participant

@Vamsi-klu
, '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

Consume every queued asset event from concurrent mapped outlets - #71994

Closed
Vamsi-klu wants to merge 2 commits into
apache:mainfrom
Vamsi-klu:pr1/54659-mapped-outlet-events
Closed

Consume every queued asset event from concurrent mapped outlets#71994
Vamsi-klu wants to merge 2 commits into
apache:mainfrom
Vamsi-klu:pr1/54659-mapped-outlet-events

Conversation

@Vamsi-klu

@Vamsi-kluVamsi-klu commented Aug 23, 2026

Copy link
Copy Markdown
Contributor

closes: #54659

GitHub still lists a few other pull requests under This was referenced. I named them in an early draft of this writeup and then deleted those lines. GitHub keeps the timeline row even after the text is gone. Those pull requests are not part of this change. The only real link is the issue this PR is meant to fix.

What I did

Tests only. Mapped producer @task(outlets=[asset]).expand(...) succeeds, then one scheduler tick. Assert all N events land on the consumer run and in triggering_asset_events. No production change. No newsfragment.

Why I did

Three mapped outlets were finishing together and the consumer only saw a subset in triggering_asset_events. Consume-by-event-id is already on main. Existing tests still insert ADRQ by hand, so they never covered this emit path.

How I did

Emit: dag_maker.run_ti(..., map_index=N)register_asset_changes_in_db.
Consume: SchedulerJobRunner._create_dagruns_for_dags.
Context: get_template_context after loading consumed_asset_events.

One tick batches all visible events. That is the default, not a bug. Leftovers stay for the next tick.

What's the impact

None at runtime. If the mapped emit path regresses, these tests fail instead of silently dropping events.

What's the testing

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

  • test_mapped_outlet_asset_events_consumed_in_one_tick
  • test_mapped_outlet_asset_events_consumed_across_staggered_ticks
  • test_mapped_outlet_asset_events_same_timestamp_are_all_consumed
  • test_mapped_outlet_asset_events_and_condition_waits_for_all_assets
  • test_mapped_outlet_asset_alias_events_are_all_consumed
uv run --project airflow-core pytest airflow-core/tests/unit/jobs/test_scheduler_job.py \
-k test_mapped_outlet -q

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

Generated-by: Grok 4.6 following the guidelines

apache#70972 already consumes queued asset events by id, so concurrent mapped
outlets no longer strand events behind a timestamp watermark. These
tests lock that contract for mapped producers so a later consume-path
change cannot silently drop events from triggering_asset_events.
closes: apache#54659
The same-timestamp and AND cases only pinned consume if events were
already queued by hand. Running mapped TIs through register/queue
keeps those pins on the user path, including alias context lookup.
closes: apache#54659

@Vamsi-kluVamsi-klu left a comment

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Description review:

  • closes: #54659 is correct.
  • Tests-only is stated; no newsfragment is correct.
  • Default batching in one tick is called out as intended, not a bug.
  • Test names and the pytest command match the diff.

Ready for maintainer review.


Drafted-by: Cursor Grok 4.6 (no human review before posting)

@Vamsi-klu

Vamsi-klu commented Aug 23, 2026

Copy link
Copy Markdown
ContributorAuthor

The extra pull requests under This was referenced are leftover from an earlier draft of the description. I put those names in, then took them out. GitHub keeps the timeline row anyway. Nothing in this PR depends on them. The only real link is the issue this PR is meant to fix.


Drafted-by: Cursor Grok 4.6

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

Labels

area:Schedulerincluding HA (high availability) scheduler

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Asset-triggered DAGs miss events from concurrently completing dynamic mapped tasks

1 participant

@Vamsi-klu