Load the correct DAG version when a task starts from a trigger - #69988

Open
seanghaeli wants to merge 3 commits into
apache:mainfrom
aws-mwaa:fix/triggerer-pin-dag-version
Open

Load the correct DAG version when a task starts from a trigger#69988
seanghaeli wants to merge 3 commits into
apache:mainfrom
aws-mwaa:fix/triggerer-pin-dag-version

Conversation

@seanghaeli

@seanghaeliseanghaeli commented Jul 16, 2026

Copy link
Copy Markdown
Contributor

If a Dag run without a bundle_version and one of its tasks point at different Dag versions, the trigger should load the latest version. Currently, it loads the task's possibly outdated version.

@seanghaeli
seanghaeliforce-pushed the fix/triggerer-pin-dag-version branch 5 times, most recently from f157049 to ecbf158CompareJuly 17, 2026 00:54
@seanghaeli
seanghaeli marked this pull request as ready for review July 17, 2026 00:59
@seanghaeliseanghaeli changed the title Pin triggerer start_from_trigger DAG resolution to the dagrun's versionLoad the correct DAG version when a task starts from a triggerJul 17, 2026
@seanghaeli

Copy link
Copy Markdown
ContributorAuthor

@o-nikolas this one is similar in spirit to your PR #69941

@o-nikolaso-nikolas left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Seems like a reasonable change semantically. But someone who knows the trigger code better should review to see if this is what we want (it seems like it should be though to me at least). Maybe @vincbeck?

@viiccwenviiccwen left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

left comments.

Comment threadairflow-core/tests/unit/jobs/test_triggerer_job.py Outdated
Comment on lines +879 to +880
version_id=trigger.task_instance.get_dagrun(session=session).created_dag_version_id
or trigger.task_instance.dag_version_id,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

created_dag_version_id is populated even for Dags with disable_bundle_versioning=True; dag_run.bundle_version is what determines
whether the run is pinned. This unconditional preference for created_dag_version_id therefore makes an unpinned start_from_trigger task load the Dag version from when the run was created, even after the scheduler has advanced the unfinished TI's dag_version_id following a reparse.

This differs from DBDagBag._version_from_dag_run(), which intentionally resolves the latest version when bundle_version is absent. Please use the run's created_dag_version_id only when the run is pinned, retain the TI version for unpinned runs, and cover both cases in the regression test.

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.

I believe this is addressed now, could you take a look?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

@viiccwen ☝️

Comment threadairflow-core/tests/unit/jobs/test_triggerer_job.py
@eladkaleladkal added this to the Airflow 3.3.1 milestone Jul 19, 2026
@eladkaleladkal added type:bug-fix Changelog: Bug Fixes backport-to-v3-3-test Backport to v3-3-test labels Jul 19, 2026
@potiuk

Copy link
Copy Markdown
Member

@seanghaeli — There are 3 unresolved review thread(s) on this PR from @viiccwen. Could you either push a fix or reply in each thread explaining why the feedback doesn't apply? Once you believe the feedback is addressed, mark the thread as resolved so the reviewer isn't re-pinged needlessly. Thanks!


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

@o-nikolas

Copy link
Copy Markdown
Contributor

@viiccwen do the changes look good to you now?

@vatsrahul1001

Copy link
Copy Markdown
Contributor

LGTM, I think we are good to merge after code owners review @dstandish@hussein-awala

@jason810496jason810496 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Thanks for the fix.

IMO, we should introduce something like get_serialized_dag_model_for_run rather than calling the private method of the DagBag.

# airflow-core/src/airflow/models/dagbag.py, next to get_dag_for_rundefget_serialized_dag_model_for_run(
self, dag_run: DagRun, *, session: Session
) ->SerializedDagModel|None:
"""Return the SerializedDagModel for the version a run executes against."""ifversion_id:=self._version_from_dag_run(dag_run=dag_run, session=session):
returnself.get_serialized_dag_model(version_id=version_id, session=session)
returnNone

defget_dag_for_run(self, dag_run: DagRun, session: Session) ->SerializedDAG|None:
ifversion_id:=self._version_from_dag_run(dag_run=dag_run, session=session):
returnself._get_dag(version_id=version_id, session=session)
returnNone

Comment threadairflow-core/src/airflow/jobs/triggerer_job_runner.py
@vatsrahul1001

Copy link
Copy Markdown
Contributor

Moving to 3.3.2 as this is still pending code owner review and do not want to rush on merging this without that

@vatsrahul1001

Copy link
Copy Markdown
Contributor

@dstandish@hussein-awala can you review this?

@hussein-awalahussein-awala left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Thanks for the fix, resolving the version from the run is the right direction and it
matches what workers already do. One thing worries me though. This applies to every
trigger, not only the start_from_trigger ones, and that makes the get_task call
able to kill the triggerer.

Here is the scenario. Take a DAG in a bundle without versioning, so its runs have
bundle_version = None.

  1. my_dag v1 has a deferrable task wait. A run is created at v1 and wait
    defers. The TI is DEFERRED with dag_version_id = v1 and there is a Trigger
    row for it.
  2. Someone edits the file and renames wait to wait_for_data. The dag processor
    creates v2 and it becomes the latest. The run is still RUNNING, its
    bundle_version is still None and created_dag_version_id is still v1.
  3. The scheduler does not clean the TI up yet.
    _check_for_removed_or_restored_tasks only marks it REMOVED when
    self.state != DagRunState.RUNNING, so nothing happens until the next
    task_instance_scheduling_decisions pass over that run.
  4. In the meantime the triggerer builds the workload for that trigger, which happens
    on a restart, on a rolling deploy, or on HA failover after assign_unassigned.
    _version_from_dag_run sees no bundle_version, returns v2 as the latest, and
    serialized_dag_model.dag.get_task("wait") raises TaskNotFound.
  5. Nobody catches it. Not _create_workload, not build_trigger_workloads, not
    run_once. TriggererJobRunner.run() logs it and re-raises, so the process exits
    and every other trigger it was running goes down with it. After the restart
    assign_unassigned gives it the same trigger back and it crashes again, until the
    scheduler gets around to marking the TI REMOVED and clean_unused deletes the
    trigger.

Steps 2 and 4 are really the same event in practice. A rolling deploy of new DAG
code is exactly when new versions show up and triggerers restart.

Before this change the version came from the TI, so the task was always there and
step 4 just worked. We already guard this pattern elsewhere. In
task_instances.py the same resolution is followed by
with contextlib.suppress(TaskNotFound).

The simplest fix I can think of is to keep the flag check on the TI version and only
re-resolve for the case the title is about:

serialized_dag_model=dag_bag.get_serialized_dag_model(
version_id=trigger.task_instance.dag_version_id, session=session
)
ifserialized_dag_model:
task=serialized_dag_model.dag.get_task(trigger.task_instance.task_id)
iftask.start_from_trigger:
dag_run=trigger.task_instance.get_dagrun(session=session)
run_version_id=DBDagBag._version_from_dag_run(dag_run=dag_run, session=session)
ifrun_version_idandrun_version_id!=trigger.task_instance.dag_version_id:
serialized_dag_model= (
dag_bag.get_serialized_dag_model(version_id=run_version_id, session=session)
orserialized_dag_model
)

Comment threadairflow-core/src/airflow/jobs/triggerer_job_runner.py
@seanghaeli
seanghaeliforce-pushed the fix/triggerer-pin-dag-version branch from 0a7a1ff to 4ce124cCompareAugust 18, 2026 21:47
@seanghaeli

Copy link
Copy Markdown
ContributorAuthor

@hussein-awala good catch, it can crash in this specific sequence. In fact, the current version of main can fail for a similar reason noted in issue #69841. The latest commit introduces a guard so it falls back to running the trigger without Dag context instead of killing the triggerer

For an unpinned run the resolved (latest) Dag version may no longer
contain a renamed or removed deferred task. get_task() then raised
TaskNotFound, which nothing on the load_triggers path caught, killing
the whole triggerer; assign_unassigned re-handed the same trigger to the
restarted triggerer, producing a crash loop. Guard the lookup and fall
through to the plain workload, mirroring the guard in the execution API.
Also introduce DBDagBag.get_serialized_dag_model_for_run so the
triggerer uses a public API instead of the private _version_from_dag_run
helper, as requested in review.
@seanghaeli
seanghaeliforce-pushed the fix/triggerer-pin-dag-version branch from 4ce124c to 05ebadcCompareAugust 18, 2026 22:46
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:Triggererbackport-to-v3-3-testBackport to v3-3-testtype:bug-fixChangelog: Bug Fixes

Projects

None yet

Development

Successfully merging this pull request may close these issues.

9 participants

@seanghaeli@potiuk@o-nikolas@vatsrahul1001@hussein-awala@jason810496@vincbeck@viiccwen@eladkal
, '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

Load the correct DAG version when a task starts from a trigger - #69988

Open
seanghaeli wants to merge 3 commits into
apache:mainfrom
aws-mwaa:fix/triggerer-pin-dag-version
Open

Load the correct DAG version when a task starts from a trigger#69988
seanghaeli wants to merge 3 commits into
apache:mainfrom
aws-mwaa:fix/triggerer-pin-dag-version

Conversation

@seanghaeli

@seanghaeliseanghaeli commented Jul 16, 2026

Copy link
Copy Markdown
Contributor

If a Dag run without a bundle_version and one of its tasks point at different Dag versions, the trigger should load the latest version. Currently, it loads the task's possibly outdated version.

@seanghaeli
seanghaeliforce-pushed the fix/triggerer-pin-dag-version branch 5 times, most recently from f157049 to ecbf158CompareJuly 17, 2026 00:54
@seanghaeli
seanghaeli marked this pull request as ready for review July 17, 2026 00:59
@seanghaeliseanghaeli changed the title Pin triggerer start_from_trigger DAG resolution to the dagrun's versionLoad the correct DAG version when a task starts from a triggerJul 17, 2026
@seanghaeli

Copy link
Copy Markdown
ContributorAuthor

@o-nikolas this one is similar in spirit to your PR #69941

@o-nikolaso-nikolas left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Seems like a reasonable change semantically. But someone who knows the trigger code better should review to see if this is what we want (it seems like it should be though to me at least). Maybe @vincbeck?

@viiccwenviiccwen left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

left comments.

Comment threadairflow-core/tests/unit/jobs/test_triggerer_job.py Outdated
Comment on lines +879 to +880
version_id=trigger.task_instance.get_dagrun(session=session).created_dag_version_id
or trigger.task_instance.dag_version_id,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

created_dag_version_id is populated even for Dags with disable_bundle_versioning=True; dag_run.bundle_version is what determines
whether the run is pinned. This unconditional preference for created_dag_version_id therefore makes an unpinned start_from_trigger task load the Dag version from when the run was created, even after the scheduler has advanced the unfinished TI's dag_version_id following a reparse.

This differs from DBDagBag._version_from_dag_run(), which intentionally resolves the latest version when bundle_version is absent. Please use the run's created_dag_version_id only when the run is pinned, retain the TI version for unpinned runs, and cover both cases in the regression test.

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.

I believe this is addressed now, could you take a look?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

@viiccwen ☝️

Comment threadairflow-core/tests/unit/jobs/test_triggerer_job.py
@eladkaleladkal added this to the Airflow 3.3.1 milestone Jul 19, 2026
@eladkaleladkal added type:bug-fix Changelog: Bug Fixes backport-to-v3-3-test Backport to v3-3-test labels Jul 19, 2026
@potiuk

Copy link
Copy Markdown
Member

@seanghaeli — There are 3 unresolved review thread(s) on this PR from @viiccwen. Could you either push a fix or reply in each thread explaining why the feedback doesn't apply? Once you believe the feedback is addressed, mark the thread as resolved so the reviewer isn't re-pinged needlessly. Thanks!


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

@o-nikolas

Copy link
Copy Markdown
Contributor

@viiccwen do the changes look good to you now?

@vatsrahul1001

Copy link
Copy Markdown
Contributor

LGTM, I think we are good to merge after code owners review @dstandish@hussein-awala

@jason810496jason810496 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Thanks for the fix.

IMO, we should introduce something like get_serialized_dag_model_for_run rather than calling the private method of the DagBag.

# airflow-core/src/airflow/models/dagbag.py, next to get_dag_for_rundefget_serialized_dag_model_for_run(
self, dag_run: DagRun, *, session: Session
) ->SerializedDagModel|None:
"""Return the SerializedDagModel for the version a run executes against."""ifversion_id:=self._version_from_dag_run(dag_run=dag_run, session=session):
returnself.get_serialized_dag_model(version_id=version_id, session=session)
returnNone

defget_dag_for_run(self, dag_run: DagRun, session: Session) ->SerializedDAG|None:
ifversion_id:=self._version_from_dag_run(dag_run=dag_run, session=session):
returnself._get_dag(version_id=version_id, session=session)
returnNone

Comment threadairflow-core/src/airflow/jobs/triggerer_job_runner.py
@vatsrahul1001

Copy link
Copy Markdown
Contributor

Moving to 3.3.2 as this is still pending code owner review and do not want to rush on merging this without that

@vatsrahul1001

Copy link
Copy Markdown
Contributor

@dstandish@hussein-awala can you review this?

@hussein-awalahussein-awala left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Thanks for the fix, resolving the version from the run is the right direction and it
matches what workers already do. One thing worries me though. This applies to every
trigger, not only the start_from_trigger ones, and that makes the get_task call
able to kill the triggerer.

Here is the scenario. Take a DAG in a bundle without versioning, so its runs have
bundle_version = None.

  1. my_dag v1 has a deferrable task wait. A run is created at v1 and wait
    defers. The TI is DEFERRED with dag_version_id = v1 and there is a Trigger
    row for it.
  2. Someone edits the file and renames wait to wait_for_data. The dag processor
    creates v2 and it becomes the latest. The run is still RUNNING, its
    bundle_version is still None and created_dag_version_id is still v1.
  3. The scheduler does not clean the TI up yet.
    _check_for_removed_or_restored_tasks only marks it REMOVED when
    self.state != DagRunState.RUNNING, so nothing happens until the next
    task_instance_scheduling_decisions pass over that run.
  4. In the meantime the triggerer builds the workload for that trigger, which happens
    on a restart, on a rolling deploy, or on HA failover after assign_unassigned.
    _version_from_dag_run sees no bundle_version, returns v2 as the latest, and
    serialized_dag_model.dag.get_task("wait") raises TaskNotFound.
  5. Nobody catches it. Not _create_workload, not build_trigger_workloads, not
    run_once. TriggererJobRunner.run() logs it and re-raises, so the process exits
    and every other trigger it was running goes down with it. After the restart
    assign_unassigned gives it the same trigger back and it crashes again, until the
    scheduler gets around to marking the TI REMOVED and clean_unused deletes the
    trigger.

Steps 2 and 4 are really the same event in practice. A rolling deploy of new DAG
code is exactly when new versions show up and triggerers restart.

Before this change the version came from the TI, so the task was always there and
step 4 just worked. We already guard this pattern elsewhere. In
task_instances.py the same resolution is followed by
with contextlib.suppress(TaskNotFound).

The simplest fix I can think of is to keep the flag check on the TI version and only
re-resolve for the case the title is about:

serialized_dag_model=dag_bag.get_serialized_dag_model(
version_id=trigger.task_instance.dag_version_id, session=session
)
ifserialized_dag_model:
task=serialized_dag_model.dag.get_task(trigger.task_instance.task_id)
iftask.start_from_trigger:
dag_run=trigger.task_instance.get_dagrun(session=session)
run_version_id=DBDagBag._version_from_dag_run(dag_run=dag_run, session=session)
ifrun_version_idandrun_version_id!=trigger.task_instance.dag_version_id:
serialized_dag_model= (
dag_bag.get_serialized_dag_model(version_id=run_version_id, session=session)
orserialized_dag_model
)

Comment threadairflow-core/src/airflow/jobs/triggerer_job_runner.py
@seanghaeli
seanghaeliforce-pushed the fix/triggerer-pin-dag-version branch from 0a7a1ff to 4ce124cCompareAugust 18, 2026 21:47
@seanghaeli

Copy link
Copy Markdown
ContributorAuthor

@hussein-awala good catch, it can crash in this specific sequence. In fact, the current version of main can fail for a similar reason noted in issue #69841. The latest commit introduces a guard so it falls back to running the trigger without Dag context instead of killing the triggerer

For an unpinned run the resolved (latest) Dag version may no longer
contain a renamed or removed deferred task. get_task() then raised
TaskNotFound, which nothing on the load_triggers path caught, killing
the whole triggerer; assign_unassigned re-handed the same trigger to the
restarted triggerer, producing a crash loop. Guard the lookup and fall
through to the plain workload, mirroring the guard in the execution API.
Also introduce DBDagBag.get_serialized_dag_model_for_run so the
triggerer uses a public API instead of the private _version_from_dag_run
helper, as requested in review.
@seanghaeli
seanghaeliforce-pushed the fix/triggerer-pin-dag-version branch from 4ce124c to 05ebadcCompareAugust 18, 2026 22:46
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:Triggererbackport-to-v3-3-testBackport to v3-3-testtype:bug-fixChangelog: Bug Fixes

Projects

None yet

Development

Successfully merging this pull request may close these issues.

9 participants

@seanghaeli@potiuk@o-nikolas@vatsrahul1001@hussein-awala@jason810496@vincbeck@viiccwen@eladkal
, '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

Load the correct DAG version when a task starts from a trigger - #69988

Open
seanghaeli wants to merge 3 commits into
apache:mainfrom
aws-mwaa:fix/triggerer-pin-dag-version
Open

Load the correct DAG version when a task starts from a trigger#69988
seanghaeli wants to merge 3 commits into
apache:mainfrom
aws-mwaa:fix/triggerer-pin-dag-version

Conversation

@seanghaeli

@seanghaeliseanghaeli commented Jul 16, 2026

Copy link
Copy Markdown
Contributor

If a Dag run without a bundle_version and one of its tasks point at different Dag versions, the trigger should load the latest version. Currently, it loads the task's possibly outdated version.

@seanghaeli
seanghaeliforce-pushed the fix/triggerer-pin-dag-version branch 5 times, most recently from f157049 to ecbf158CompareJuly 17, 2026 00:54
@seanghaeli
seanghaeli marked this pull request as ready for review July 17, 2026 00:59
@seanghaeliseanghaeli changed the title Pin triggerer start_from_trigger DAG resolution to the dagrun's versionLoad the correct DAG version when a task starts from a triggerJul 17, 2026
@seanghaeli

Copy link
Copy Markdown
ContributorAuthor

@o-nikolas this one is similar in spirit to your PR #69941

@o-nikolaso-nikolas left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Seems like a reasonable change semantically. But someone who knows the trigger code better should review to see if this is what we want (it seems like it should be though to me at least). Maybe @vincbeck?

@viiccwenviiccwen left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

left comments.

Comment threadairflow-core/tests/unit/jobs/test_triggerer_job.py Outdated
Comment on lines +879 to +880
version_id=trigger.task_instance.get_dagrun(session=session).created_dag_version_id
or trigger.task_instance.dag_version_id,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

created_dag_version_id is populated even for Dags with disable_bundle_versioning=True; dag_run.bundle_version is what determines
whether the run is pinned. This unconditional preference for created_dag_version_id therefore makes an unpinned start_from_trigger task load the Dag version from when the run was created, even after the scheduler has advanced the unfinished TI's dag_version_id following a reparse.

This differs from DBDagBag._version_from_dag_run(), which intentionally resolves the latest version when bundle_version is absent. Please use the run's created_dag_version_id only when the run is pinned, retain the TI version for unpinned runs, and cover both cases in the regression test.

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.

I believe this is addressed now, could you take a look?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

@viiccwen ☝️

Comment threadairflow-core/tests/unit/jobs/test_triggerer_job.py
@eladkaleladkal added this to the Airflow 3.3.1 milestone Jul 19, 2026
@eladkaleladkal added type:bug-fix Changelog: Bug Fixes backport-to-v3-3-test Backport to v3-3-test labels Jul 19, 2026
@potiuk

Copy link
Copy Markdown
Member

@seanghaeli — There are 3 unresolved review thread(s) on this PR from @viiccwen. Could you either push a fix or reply in each thread explaining why the feedback doesn't apply? Once you believe the feedback is addressed, mark the thread as resolved so the reviewer isn't re-pinged needlessly. Thanks!


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

@o-nikolas

Copy link
Copy Markdown
Contributor

@viiccwen do the changes look good to you now?

@vatsrahul1001

Copy link
Copy Markdown
Contributor

LGTM, I think we are good to merge after code owners review @dstandish@hussein-awala

@jason810496jason810496 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Thanks for the fix.

IMO, we should introduce something like get_serialized_dag_model_for_run rather than calling the private method of the DagBag.

# airflow-core/src/airflow/models/dagbag.py, next to get_dag_for_rundefget_serialized_dag_model_for_run(
self, dag_run: DagRun, *, session: Session
) ->SerializedDagModel|None:
"""Return the SerializedDagModel for the version a run executes against."""ifversion_id:=self._version_from_dag_run(dag_run=dag_run, session=session):
returnself.get_serialized_dag_model(version_id=version_id, session=session)
returnNone

defget_dag_for_run(self, dag_run: DagRun, session: Session) ->SerializedDAG|None:
ifversion_id:=self._version_from_dag_run(dag_run=dag_run, session=session):
returnself._get_dag(version_id=version_id, session=session)
returnNone

Comment threadairflow-core/src/airflow/jobs/triggerer_job_runner.py
@vatsrahul1001

Copy link
Copy Markdown
Contributor

Moving to 3.3.2 as this is still pending code owner review and do not want to rush on merging this without that

@vatsrahul1001

Copy link
Copy Markdown
Contributor

@dstandish@hussein-awala can you review this?

@hussein-awalahussein-awala left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Thanks for the fix, resolving the version from the run is the right direction and it
matches what workers already do. One thing worries me though. This applies to every
trigger, not only the start_from_trigger ones, and that makes the get_task call
able to kill the triggerer.

Here is the scenario. Take a DAG in a bundle without versioning, so its runs have
bundle_version = None.

  1. my_dag v1 has a deferrable task wait. A run is created at v1 and wait
    defers. The TI is DEFERRED with dag_version_id = v1 and there is a Trigger
    row for it.
  2. Someone edits the file and renames wait to wait_for_data. The dag processor
    creates v2 and it becomes the latest. The run is still RUNNING, its
    bundle_version is still None and created_dag_version_id is still v1.
  3. The scheduler does not clean the TI up yet.
    _check_for_removed_or_restored_tasks only marks it REMOVED when
    self.state != DagRunState.RUNNING, so nothing happens until the next
    task_instance_scheduling_decisions pass over that run.
  4. In the meantime the triggerer builds the workload for that trigger, which happens
    on a restart, on a rolling deploy, or on HA failover after assign_unassigned.
    _version_from_dag_run sees no bundle_version, returns v2 as the latest, and
    serialized_dag_model.dag.get_task("wait") raises TaskNotFound.
  5. Nobody catches it. Not _create_workload, not build_trigger_workloads, not
    run_once. TriggererJobRunner.run() logs it and re-raises, so the process exits
    and every other trigger it was running goes down with it. After the restart
    assign_unassigned gives it the same trigger back and it crashes again, until the
    scheduler gets around to marking the TI REMOVED and clean_unused deletes the
    trigger.

Steps 2 and 4 are really the same event in practice. A rolling deploy of new DAG
code is exactly when new versions show up and triggerers restart.

Before this change the version came from the TI, so the task was always there and
step 4 just worked. We already guard this pattern elsewhere. In
task_instances.py the same resolution is followed by
with contextlib.suppress(TaskNotFound).

The simplest fix I can think of is to keep the flag check on the TI version and only
re-resolve for the case the title is about:

serialized_dag_model=dag_bag.get_serialized_dag_model(
version_id=trigger.task_instance.dag_version_id, session=session
)
ifserialized_dag_model:
task=serialized_dag_model.dag.get_task(trigger.task_instance.task_id)
iftask.start_from_trigger:
dag_run=trigger.task_instance.get_dagrun(session=session)
run_version_id=DBDagBag._version_from_dag_run(dag_run=dag_run, session=session)
ifrun_version_idandrun_version_id!=trigger.task_instance.dag_version_id:
serialized_dag_model= (
dag_bag.get_serialized_dag_model(version_id=run_version_id, session=session)
orserialized_dag_model
)

Comment threadairflow-core/src/airflow/jobs/triggerer_job_runner.py
@seanghaeli
seanghaeliforce-pushed the fix/triggerer-pin-dag-version branch from 0a7a1ff to 4ce124cCompareAugust 18, 2026 21:47
@seanghaeli

Copy link
Copy Markdown
ContributorAuthor

@hussein-awala good catch, it can crash in this specific sequence. In fact, the current version of main can fail for a similar reason noted in issue #69841. The latest commit introduces a guard so it falls back to running the trigger without Dag context instead of killing the triggerer

For an unpinned run the resolved (latest) Dag version may no longer
contain a renamed or removed deferred task. get_task() then raised
TaskNotFound, which nothing on the load_triggers path caught, killing
the whole triggerer; assign_unassigned re-handed the same trigger to the
restarted triggerer, producing a crash loop. Guard the lookup and fall
through to the plain workload, mirroring the guard in the execution API.
Also introduce DBDagBag.get_serialized_dag_model_for_run so the
triggerer uses a public API instead of the private _version_from_dag_run
helper, as requested in review.
@seanghaeli
seanghaeliforce-pushed the fix/triggerer-pin-dag-version branch from 4ce124c to 05ebadcCompareAugust 18, 2026 22:46
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:Triggererbackport-to-v3-3-testBackport to v3-3-testtype:bug-fixChangelog: Bug Fixes

Projects

None yet

Development

Successfully merging this pull request may close these issues.

9 participants

@seanghaeli@potiuk@o-nikolas@vatsrahul1001@hussein-awala@jason810496@vincbeck@viiccwen@eladkal
, '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

Load the correct DAG version when a task starts from a trigger - #69988

Open
seanghaeli wants to merge 3 commits into
apache:mainfrom
aws-mwaa:fix/triggerer-pin-dag-version
Open

Load the correct DAG version when a task starts from a trigger#69988
seanghaeli wants to merge 3 commits into
apache:mainfrom
aws-mwaa:fix/triggerer-pin-dag-version

Conversation

@seanghaeli

@seanghaeliseanghaeli commented Jul 16, 2026

Copy link
Copy Markdown
Contributor

If a Dag run without a bundle_version and one of its tasks point at different Dag versions, the trigger should load the latest version. Currently, it loads the task's possibly outdated version.

@seanghaeli
seanghaeliforce-pushed the fix/triggerer-pin-dag-version branch 5 times, most recently from f157049 to ecbf158CompareJuly 17, 2026 00:54
@seanghaeli
seanghaeli marked this pull request as ready for review July 17, 2026 00:59
@seanghaeliseanghaeli changed the title Pin triggerer start_from_trigger DAG resolution to the dagrun's versionLoad the correct DAG version when a task starts from a triggerJul 17, 2026
@seanghaeli

Copy link
Copy Markdown
ContributorAuthor

@o-nikolas this one is similar in spirit to your PR #69941

@o-nikolaso-nikolas left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Seems like a reasonable change semantically. But someone who knows the trigger code better should review to see if this is what we want (it seems like it should be though to me at least). Maybe @vincbeck?

@viiccwenviiccwen left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

left comments.

Comment threadairflow-core/tests/unit/jobs/test_triggerer_job.py Outdated
Comment on lines +879 to +880
version_id=trigger.task_instance.get_dagrun(session=session).created_dag_version_id
or trigger.task_instance.dag_version_id,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

created_dag_version_id is populated even for Dags with disable_bundle_versioning=True; dag_run.bundle_version is what determines
whether the run is pinned. This unconditional preference for created_dag_version_id therefore makes an unpinned start_from_trigger task load the Dag version from when the run was created, even after the scheduler has advanced the unfinished TI's dag_version_id following a reparse.

This differs from DBDagBag._version_from_dag_run(), which intentionally resolves the latest version when bundle_version is absent. Please use the run's created_dag_version_id only when the run is pinned, retain the TI version for unpinned runs, and cover both cases in the regression test.

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.

I believe this is addressed now, could you take a look?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

@viiccwen ☝️

Comment threadairflow-core/tests/unit/jobs/test_triggerer_job.py
@eladkaleladkal added this to the Airflow 3.3.1 milestone Jul 19, 2026
@eladkaleladkal added type:bug-fix Changelog: Bug Fixes backport-to-v3-3-test Backport to v3-3-test labels Jul 19, 2026
@potiuk

Copy link
Copy Markdown
Member

@seanghaeli — There are 3 unresolved review thread(s) on this PR from @viiccwen. Could you either push a fix or reply in each thread explaining why the feedback doesn't apply? Once you believe the feedback is addressed, mark the thread as resolved so the reviewer isn't re-pinged needlessly. Thanks!


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

@o-nikolas

Copy link
Copy Markdown
Contributor

@viiccwen do the changes look good to you now?

@vatsrahul1001

Copy link
Copy Markdown
Contributor

LGTM, I think we are good to merge after code owners review @dstandish@hussein-awala

@jason810496jason810496 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Thanks for the fix.

IMO, we should introduce something like get_serialized_dag_model_for_run rather than calling the private method of the DagBag.

# airflow-core/src/airflow/models/dagbag.py, next to get_dag_for_rundefget_serialized_dag_model_for_run(
self, dag_run: DagRun, *, session: Session
) ->SerializedDagModel|None:
"""Return the SerializedDagModel for the version a run executes against."""ifversion_id:=self._version_from_dag_run(dag_run=dag_run, session=session):
returnself.get_serialized_dag_model(version_id=version_id, session=session)
returnNone

defget_dag_for_run(self, dag_run: DagRun, session: Session) ->SerializedDAG|None:
ifversion_id:=self._version_from_dag_run(dag_run=dag_run, session=session):
returnself._get_dag(version_id=version_id, session=session)
returnNone

Comment threadairflow-core/src/airflow/jobs/triggerer_job_runner.py
@vatsrahul1001

Copy link
Copy Markdown
Contributor

Moving to 3.3.2 as this is still pending code owner review and do not want to rush on merging this without that

@vatsrahul1001

Copy link
Copy Markdown
Contributor

@dstandish@hussein-awala can you review this?

@hussein-awalahussein-awala left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Thanks for the fix, resolving the version from the run is the right direction and it
matches what workers already do. One thing worries me though. This applies to every
trigger, not only the start_from_trigger ones, and that makes the get_task call
able to kill the triggerer.

Here is the scenario. Take a DAG in a bundle without versioning, so its runs have
bundle_version = None.

  1. my_dag v1 has a deferrable task wait. A run is created at v1 and wait
    defers. The TI is DEFERRED with dag_version_id = v1 and there is a Trigger
    row for it.
  2. Someone edits the file and renames wait to wait_for_data. The dag processor
    creates v2 and it becomes the latest. The run is still RUNNING, its
    bundle_version is still None and created_dag_version_id is still v1.
  3. The scheduler does not clean the TI up yet.
    _check_for_removed_or_restored_tasks only marks it REMOVED when
    self.state != DagRunState.RUNNING, so nothing happens until the next
    task_instance_scheduling_decisions pass over that run.
  4. In the meantime the triggerer builds the workload for that trigger, which happens
    on a restart, on a rolling deploy, or on HA failover after assign_unassigned.
    _version_from_dag_run sees no bundle_version, returns v2 as the latest, and
    serialized_dag_model.dag.get_task("wait") raises TaskNotFound.
  5. Nobody catches it. Not _create_workload, not build_trigger_workloads, not
    run_once. TriggererJobRunner.run() logs it and re-raises, so the process exits
    and every other trigger it was running goes down with it. After the restart
    assign_unassigned gives it the same trigger back and it crashes again, until the
    scheduler gets around to marking the TI REMOVED and clean_unused deletes the
    trigger.

Steps 2 and 4 are really the same event in practice. A rolling deploy of new DAG
code is exactly when new versions show up and triggerers restart.

Before this change the version came from the TI, so the task was always there and
step 4 just worked. We already guard this pattern elsewhere. In
task_instances.py the same resolution is followed by
with contextlib.suppress(TaskNotFound).

The simplest fix I can think of is to keep the flag check on the TI version and only
re-resolve for the case the title is about:

serialized_dag_model=dag_bag.get_serialized_dag_model(
version_id=trigger.task_instance.dag_version_id, session=session
)
ifserialized_dag_model:
task=serialized_dag_model.dag.get_task(trigger.task_instance.task_id)
iftask.start_from_trigger:
dag_run=trigger.task_instance.get_dagrun(session=session)
run_version_id=DBDagBag._version_from_dag_run(dag_run=dag_run, session=session)
ifrun_version_idandrun_version_id!=trigger.task_instance.dag_version_id:
serialized_dag_model= (
dag_bag.get_serialized_dag_model(version_id=run_version_id, session=session)
orserialized_dag_model
)

Comment threadairflow-core/src/airflow/jobs/triggerer_job_runner.py
@seanghaeli
seanghaeliforce-pushed the fix/triggerer-pin-dag-version branch from 0a7a1ff to 4ce124cCompareAugust 18, 2026 21:47
@seanghaeli

Copy link
Copy Markdown
ContributorAuthor

@hussein-awala good catch, it can crash in this specific sequence. In fact, the current version of main can fail for a similar reason noted in issue #69841. The latest commit introduces a guard so it falls back to running the trigger without Dag context instead of killing the triggerer

For an unpinned run the resolved (latest) Dag version may no longer
contain a renamed or removed deferred task. get_task() then raised
TaskNotFound, which nothing on the load_triggers path caught, killing
the whole triggerer; assign_unassigned re-handed the same trigger to the
restarted triggerer, producing a crash loop. Guard the lookup and fall
through to the plain workload, mirroring the guard in the execution API.
Also introduce DBDagBag.get_serialized_dag_model_for_run so the
triggerer uses a public API instead of the private _version_from_dag_run
helper, as requested in review.
@seanghaeli
seanghaeliforce-pushed the fix/triggerer-pin-dag-version branch from 4ce124c to 05ebadcCompareAugust 18, 2026 22:46
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:Triggererbackport-to-v3-3-testBackport to v3-3-testtype:bug-fixChangelog: Bug Fixes

Projects

None yet

Development

Successfully merging this pull request may close these issues.

9 participants

@seanghaeli@potiuk@o-nikolas@vatsrahul1001@hussein-awala@jason810496@vincbeck@viiccwen@eladkal
, '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

Load the correct DAG version when a task starts from a trigger - #69988

Open
seanghaeli wants to merge 3 commits into
apache:mainfrom
aws-mwaa:fix/triggerer-pin-dag-version
Open

Load the correct DAG version when a task starts from a trigger#69988
seanghaeli wants to merge 3 commits into
apache:mainfrom
aws-mwaa:fix/triggerer-pin-dag-version

Conversation

@seanghaeli

@seanghaeliseanghaeli commented Jul 16, 2026

Copy link
Copy Markdown
Contributor

If a Dag run without a bundle_version and one of its tasks point at different Dag versions, the trigger should load the latest version. Currently, it loads the task's possibly outdated version.

@seanghaeli
seanghaeliforce-pushed the fix/triggerer-pin-dag-version branch 5 times, most recently from f157049 to ecbf158CompareJuly 17, 2026 00:54
@seanghaeli
seanghaeli marked this pull request as ready for review July 17, 2026 00:59
@seanghaeliseanghaeli changed the title Pin triggerer start_from_trigger DAG resolution to the dagrun's versionLoad the correct DAG version when a task starts from a triggerJul 17, 2026
@seanghaeli

Copy link
Copy Markdown
ContributorAuthor

@o-nikolas this one is similar in spirit to your PR #69941

@o-nikolaso-nikolas left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Seems like a reasonable change semantically. But someone who knows the trigger code better should review to see if this is what we want (it seems like it should be though to me at least). Maybe @vincbeck?

@viiccwenviiccwen left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

left comments.

Comment threadairflow-core/tests/unit/jobs/test_triggerer_job.py Outdated
Comment on lines +879 to +880
version_id=trigger.task_instance.get_dagrun(session=session).created_dag_version_id
or trigger.task_instance.dag_version_id,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

created_dag_version_id is populated even for Dags with disable_bundle_versioning=True; dag_run.bundle_version is what determines
whether the run is pinned. This unconditional preference for created_dag_version_id therefore makes an unpinned start_from_trigger task load the Dag version from when the run was created, even after the scheduler has advanced the unfinished TI's dag_version_id following a reparse.

This differs from DBDagBag._version_from_dag_run(), which intentionally resolves the latest version when bundle_version is absent. Please use the run's created_dag_version_id only when the run is pinned, retain the TI version for unpinned runs, and cover both cases in the regression test.

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.

I believe this is addressed now, could you take a look?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

@viiccwen ☝️

Comment threadairflow-core/tests/unit/jobs/test_triggerer_job.py
@eladkaleladkal added this to the Airflow 3.3.1 milestone Jul 19, 2026
@eladkaleladkal added type:bug-fix Changelog: Bug Fixes backport-to-v3-3-test Backport to v3-3-test labels Jul 19, 2026
@potiuk

Copy link
Copy Markdown
Member

@seanghaeli — There are 3 unresolved review thread(s) on this PR from @viiccwen. Could you either push a fix or reply in each thread explaining why the feedback doesn't apply? Once you believe the feedback is addressed, mark the thread as resolved so the reviewer isn't re-pinged needlessly. Thanks!


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

@o-nikolas

Copy link
Copy Markdown
Contributor

@viiccwen do the changes look good to you now?

@vatsrahul1001

Copy link
Copy Markdown
Contributor

LGTM, I think we are good to merge after code owners review @dstandish@hussein-awala

@jason810496jason810496 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Thanks for the fix.

IMO, we should introduce something like get_serialized_dag_model_for_run rather than calling the private method of the DagBag.

# airflow-core/src/airflow/models/dagbag.py, next to get_dag_for_rundefget_serialized_dag_model_for_run(
self, dag_run: DagRun, *, session: Session
) ->SerializedDagModel|None:
"""Return the SerializedDagModel for the version a run executes against."""ifversion_id:=self._version_from_dag_run(dag_run=dag_run, session=session):
returnself.get_serialized_dag_model(version_id=version_id, session=session)
returnNone

defget_dag_for_run(self, dag_run: DagRun, session: Session) ->SerializedDAG|None:
ifversion_id:=self._version_from_dag_run(dag_run=dag_run, session=session):
returnself._get_dag(version_id=version_id, session=session)
returnNone

Comment threadairflow-core/src/airflow/jobs/triggerer_job_runner.py
@vatsrahul1001

Copy link
Copy Markdown
Contributor

Moving to 3.3.2 as this is still pending code owner review and do not want to rush on merging this without that

@vatsrahul1001

Copy link
Copy Markdown
Contributor

@dstandish@hussein-awala can you review this?

@hussein-awalahussein-awala left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Thanks for the fix, resolving the version from the run is the right direction and it
matches what workers already do. One thing worries me though. This applies to every
trigger, not only the start_from_trigger ones, and that makes the get_task call
able to kill the triggerer.

Here is the scenario. Take a DAG in a bundle without versioning, so its runs have
bundle_version = None.

  1. my_dag v1 has a deferrable task wait. A run is created at v1 and wait
    defers. The TI is DEFERRED with dag_version_id = v1 and there is a Trigger
    row for it.
  2. Someone edits the file and renames wait to wait_for_data. The dag processor
    creates v2 and it becomes the latest. The run is still RUNNING, its
    bundle_version is still None and created_dag_version_id is still v1.
  3. The scheduler does not clean the TI up yet.
    _check_for_removed_or_restored_tasks only marks it REMOVED when
    self.state != DagRunState.RUNNING, so nothing happens until the next
    task_instance_scheduling_decisions pass over that run.
  4. In the meantime the triggerer builds the workload for that trigger, which happens
    on a restart, on a rolling deploy, or on HA failover after assign_unassigned.
    _version_from_dag_run sees no bundle_version, returns v2 as the latest, and
    serialized_dag_model.dag.get_task("wait") raises TaskNotFound.
  5. Nobody catches it. Not _create_workload, not build_trigger_workloads, not
    run_once. TriggererJobRunner.run() logs it and re-raises, so the process exits
    and every other trigger it was running goes down with it. After the restart
    assign_unassigned gives it the same trigger back and it crashes again, until the
    scheduler gets around to marking the TI REMOVED and clean_unused deletes the
    trigger.

Steps 2 and 4 are really the same event in practice. A rolling deploy of new DAG
code is exactly when new versions show up and triggerers restart.

Before this change the version came from the TI, so the task was always there and
step 4 just worked. We already guard this pattern elsewhere. In
task_instances.py the same resolution is followed by
with contextlib.suppress(TaskNotFound).

The simplest fix I can think of is to keep the flag check on the TI version and only
re-resolve for the case the title is about:

serialized_dag_model=dag_bag.get_serialized_dag_model(
version_id=trigger.task_instance.dag_version_id, session=session
)
ifserialized_dag_model:
task=serialized_dag_model.dag.get_task(trigger.task_instance.task_id)
iftask.start_from_trigger:
dag_run=trigger.task_instance.get_dagrun(session=session)
run_version_id=DBDagBag._version_from_dag_run(dag_run=dag_run, session=session)
ifrun_version_idandrun_version_id!=trigger.task_instance.dag_version_id:
serialized_dag_model= (
dag_bag.get_serialized_dag_model(version_id=run_version_id, session=session)
orserialized_dag_model
)

Comment threadairflow-core/src/airflow/jobs/triggerer_job_runner.py
@seanghaeli
seanghaeliforce-pushed the fix/triggerer-pin-dag-version branch from 0a7a1ff to 4ce124cCompareAugust 18, 2026 21:47
@seanghaeli

Copy link
Copy Markdown
ContributorAuthor

@hussein-awala good catch, it can crash in this specific sequence. In fact, the current version of main can fail for a similar reason noted in issue #69841. The latest commit introduces a guard so it falls back to running the trigger without Dag context instead of killing the triggerer

For an unpinned run the resolved (latest) Dag version may no longer
contain a renamed or removed deferred task. get_task() then raised
TaskNotFound, which nothing on the load_triggers path caught, killing
the whole triggerer; assign_unassigned re-handed the same trigger to the
restarted triggerer, producing a crash loop. Guard the lookup and fall
through to the plain workload, mirroring the guard in the execution API.
Also introduce DBDagBag.get_serialized_dag_model_for_run so the
triggerer uses a public API instead of the private _version_from_dag_run
helper, as requested in review.
@seanghaeli
seanghaeliforce-pushed the fix/triggerer-pin-dag-version branch from 4ce124c to 05ebadcCompareAugust 18, 2026 22:46
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:Triggererbackport-to-v3-3-testBackport to v3-3-testtype:bug-fixChangelog: Bug Fixes

Projects

None yet

Development

Successfully merging this pull request may close these issues.

9 participants

@seanghaeli@potiuk@o-nikolas@vatsrahul1001@hussein-awala@jason810496@vincbeck@viiccwen@eladkal
, '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

Load the correct DAG version when a task starts from a trigger - #69988

Open
seanghaeli wants to merge 3 commits into
apache:mainfrom
aws-mwaa:fix/triggerer-pin-dag-version
Open

Load the correct DAG version when a task starts from a trigger#69988
seanghaeli wants to merge 3 commits into
apache:mainfrom
aws-mwaa:fix/triggerer-pin-dag-version

Conversation

@seanghaeli

@seanghaeliseanghaeli commented Jul 16, 2026

Copy link
Copy Markdown
Contributor

If a Dag run without a bundle_version and one of its tasks point at different Dag versions, the trigger should load the latest version. Currently, it loads the task's possibly outdated version.

@seanghaeli
seanghaeliforce-pushed the fix/triggerer-pin-dag-version branch 5 times, most recently from f157049 to ecbf158CompareJuly 17, 2026 00:54
@seanghaeli
seanghaeli marked this pull request as ready for review July 17, 2026 00:59
@seanghaeliseanghaeli changed the title Pin triggerer start_from_trigger DAG resolution to the dagrun's versionLoad the correct DAG version when a task starts from a triggerJul 17, 2026
@seanghaeli

Copy link
Copy Markdown
ContributorAuthor

@o-nikolas this one is similar in spirit to your PR #69941

@o-nikolaso-nikolas left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Seems like a reasonable change semantically. But someone who knows the trigger code better should review to see if this is what we want (it seems like it should be though to me at least). Maybe @vincbeck?

@viiccwenviiccwen left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

left comments.

Comment threadairflow-core/tests/unit/jobs/test_triggerer_job.py Outdated
Comment on lines +879 to +880
version_id=trigger.task_instance.get_dagrun(session=session).created_dag_version_id
or trigger.task_instance.dag_version_id,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

created_dag_version_id is populated even for Dags with disable_bundle_versioning=True; dag_run.bundle_version is what determines
whether the run is pinned. This unconditional preference for created_dag_version_id therefore makes an unpinned start_from_trigger task load the Dag version from when the run was created, even after the scheduler has advanced the unfinished TI's dag_version_id following a reparse.

This differs from DBDagBag._version_from_dag_run(), which intentionally resolves the latest version when bundle_version is absent. Please use the run's created_dag_version_id only when the run is pinned, retain the TI version for unpinned runs, and cover both cases in the regression test.

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.

I believe this is addressed now, could you take a look?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

@viiccwen ☝️

Comment threadairflow-core/tests/unit/jobs/test_triggerer_job.py
@eladkaleladkal added this to the Airflow 3.3.1 milestone Jul 19, 2026
@eladkaleladkal added type:bug-fix Changelog: Bug Fixes backport-to-v3-3-test Backport to v3-3-test labels Jul 19, 2026
@potiuk

Copy link
Copy Markdown
Member

@seanghaeli — There are 3 unresolved review thread(s) on this PR from @viiccwen. Could you either push a fix or reply in each thread explaining why the feedback doesn't apply? Once you believe the feedback is addressed, mark the thread as resolved so the reviewer isn't re-pinged needlessly. Thanks!


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

@o-nikolas

Copy link
Copy Markdown
Contributor

@viiccwen do the changes look good to you now?

@vatsrahul1001

Copy link
Copy Markdown
Contributor

LGTM, I think we are good to merge after code owners review @dstandish@hussein-awala

@jason810496jason810496 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Thanks for the fix.

IMO, we should introduce something like get_serialized_dag_model_for_run rather than calling the private method of the DagBag.

# airflow-core/src/airflow/models/dagbag.py, next to get_dag_for_rundefget_serialized_dag_model_for_run(
self, dag_run: DagRun, *, session: Session
) ->SerializedDagModel|None:
"""Return the SerializedDagModel for the version a run executes against."""ifversion_id:=self._version_from_dag_run(dag_run=dag_run, session=session):
returnself.get_serialized_dag_model(version_id=version_id, session=session)
returnNone

defget_dag_for_run(self, dag_run: DagRun, session: Session) ->SerializedDAG|None:
ifversion_id:=self._version_from_dag_run(dag_run=dag_run, session=session):
returnself._get_dag(version_id=version_id, session=session)
returnNone

Comment threadairflow-core/src/airflow/jobs/triggerer_job_runner.py
@vatsrahul1001

Copy link
Copy Markdown
Contributor

Moving to 3.3.2 as this is still pending code owner review and do not want to rush on merging this without that

@vatsrahul1001

Copy link
Copy Markdown
Contributor

@dstandish@hussein-awala can you review this?

@hussein-awalahussein-awala left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Thanks for the fix, resolving the version from the run is the right direction and it
matches what workers already do. One thing worries me though. This applies to every
trigger, not only the start_from_trigger ones, and that makes the get_task call
able to kill the triggerer.

Here is the scenario. Take a DAG in a bundle without versioning, so its runs have
bundle_version = None.

  1. my_dag v1 has a deferrable task wait. A run is created at v1 and wait
    defers. The TI is DEFERRED with dag_version_id = v1 and there is a Trigger
    row for it.
  2. Someone edits the file and renames wait to wait_for_data. The dag processor
    creates v2 and it becomes the latest. The run is still RUNNING, its
    bundle_version is still None and created_dag_version_id is still v1.
  3. The scheduler does not clean the TI up yet.
    _check_for_removed_or_restored_tasks only marks it REMOVED when
    self.state != DagRunState.RUNNING, so nothing happens until the next
    task_instance_scheduling_decisions pass over that run.
  4. In the meantime the triggerer builds the workload for that trigger, which happens
    on a restart, on a rolling deploy, or on HA failover after assign_unassigned.
    _version_from_dag_run sees no bundle_version, returns v2 as the latest, and
    serialized_dag_model.dag.get_task("wait") raises TaskNotFound.
  5. Nobody catches it. Not _create_workload, not build_trigger_workloads, not
    run_once. TriggererJobRunner.run() logs it and re-raises, so the process exits
    and every other trigger it was running goes down with it. After the restart
    assign_unassigned gives it the same trigger back and it crashes again, until the
    scheduler gets around to marking the TI REMOVED and clean_unused deletes the
    trigger.

Steps 2 and 4 are really the same event in practice. A rolling deploy of new DAG
code is exactly when new versions show up and triggerers restart.

Before this change the version came from the TI, so the task was always there and
step 4 just worked. We already guard this pattern elsewhere. In
task_instances.py the same resolution is followed by
with contextlib.suppress(TaskNotFound).

The simplest fix I can think of is to keep the flag check on the TI version and only
re-resolve for the case the title is about:

serialized_dag_model=dag_bag.get_serialized_dag_model(
version_id=trigger.task_instance.dag_version_id, session=session
)
ifserialized_dag_model:
task=serialized_dag_model.dag.get_task(trigger.task_instance.task_id)
iftask.start_from_trigger:
dag_run=trigger.task_instance.get_dagrun(session=session)
run_version_id=DBDagBag._version_from_dag_run(dag_run=dag_run, session=session)
ifrun_version_idandrun_version_id!=trigger.task_instance.dag_version_id:
serialized_dag_model= (
dag_bag.get_serialized_dag_model(version_id=run_version_id, session=session)
orserialized_dag_model
)

Comment threadairflow-core/src/airflow/jobs/triggerer_job_runner.py
@seanghaeli
seanghaeliforce-pushed the fix/triggerer-pin-dag-version branch from 0a7a1ff to 4ce124cCompareAugust 18, 2026 21:47
@seanghaeli

Copy link
Copy Markdown
ContributorAuthor

@hussein-awala good catch, it can crash in this specific sequence. In fact, the current version of main can fail for a similar reason noted in issue #69841. The latest commit introduces a guard so it falls back to running the trigger without Dag context instead of killing the triggerer

For an unpinned run the resolved (latest) Dag version may no longer
contain a renamed or removed deferred task. get_task() then raised
TaskNotFound, which nothing on the load_triggers path caught, killing
the whole triggerer; assign_unassigned re-handed the same trigger to the
restarted triggerer, producing a crash loop. Guard the lookup and fall
through to the plain workload, mirroring the guard in the execution API.
Also introduce DBDagBag.get_serialized_dag_model_for_run so the
triggerer uses a public API instead of the private _version_from_dag_run
helper, as requested in review.
@seanghaeli
seanghaeliforce-pushed the fix/triggerer-pin-dag-version branch from 4ce124c to 05ebadcCompareAugust 18, 2026 22:46
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:Triggererbackport-to-v3-3-testBackport to v3-3-testtype:bug-fixChangelog: Bug Fixes

Projects

None yet

Development

Successfully merging this pull request may close these issues.

9 participants

@seanghaeli@potiuk@o-nikolas@vatsrahul1001@hussein-awala@jason810496@vincbeck@viiccwen@eladkal
, '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

Load the correct DAG version when a task starts from a trigger - #69988

Open
seanghaeli wants to merge 3 commits into
apache:mainfrom
aws-mwaa:fix/triggerer-pin-dag-version
Open

Load the correct DAG version when a task starts from a trigger#69988
seanghaeli wants to merge 3 commits into
apache:mainfrom
aws-mwaa:fix/triggerer-pin-dag-version

Conversation

@seanghaeli

@seanghaeliseanghaeli commented Jul 16, 2026

Copy link
Copy Markdown
Contributor

If a Dag run without a bundle_version and one of its tasks point at different Dag versions, the trigger should load the latest version. Currently, it loads the task's possibly outdated version.

@seanghaeli
seanghaeliforce-pushed the fix/triggerer-pin-dag-version branch 5 times, most recently from f157049 to ecbf158CompareJuly 17, 2026 00:54
@seanghaeli
seanghaeli marked this pull request as ready for review July 17, 2026 00:59
@seanghaeliseanghaeli changed the title Pin triggerer start_from_trigger DAG resolution to the dagrun's versionLoad the correct DAG version when a task starts from a triggerJul 17, 2026
@seanghaeli

Copy link
Copy Markdown
ContributorAuthor

@o-nikolas this one is similar in spirit to your PR #69941

@o-nikolaso-nikolas left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Seems like a reasonable change semantically. But someone who knows the trigger code better should review to see if this is what we want (it seems like it should be though to me at least). Maybe @vincbeck?

@viiccwenviiccwen left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

left comments.

Comment threadairflow-core/tests/unit/jobs/test_triggerer_job.py Outdated
Comment on lines +879 to +880
version_id=trigger.task_instance.get_dagrun(session=session).created_dag_version_id
or trigger.task_instance.dag_version_id,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

created_dag_version_id is populated even for Dags with disable_bundle_versioning=True; dag_run.bundle_version is what determines
whether the run is pinned. This unconditional preference for created_dag_version_id therefore makes an unpinned start_from_trigger task load the Dag version from when the run was created, even after the scheduler has advanced the unfinished TI's dag_version_id following a reparse.

This differs from DBDagBag._version_from_dag_run(), which intentionally resolves the latest version when bundle_version is absent. Please use the run's created_dag_version_id only when the run is pinned, retain the TI version for unpinned runs, and cover both cases in the regression test.

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.

I believe this is addressed now, could you take a look?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

@viiccwen ☝️

Comment threadairflow-core/tests/unit/jobs/test_triggerer_job.py
@eladkaleladkal added this to the Airflow 3.3.1 milestone Jul 19, 2026
@eladkaleladkal added type:bug-fix Changelog: Bug Fixes backport-to-v3-3-test Backport to v3-3-test labels Jul 19, 2026
@potiuk

Copy link
Copy Markdown
Member

@seanghaeli — There are 3 unresolved review thread(s) on this PR from @viiccwen. Could you either push a fix or reply in each thread explaining why the feedback doesn't apply? Once you believe the feedback is addressed, mark the thread as resolved so the reviewer isn't re-pinged needlessly. Thanks!


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

@o-nikolas

Copy link
Copy Markdown
Contributor

@viiccwen do the changes look good to you now?

@vatsrahul1001

Copy link
Copy Markdown
Contributor

LGTM, I think we are good to merge after code owners review @dstandish@hussein-awala

@jason810496jason810496 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Thanks for the fix.

IMO, we should introduce something like get_serialized_dag_model_for_run rather than calling the private method of the DagBag.

# airflow-core/src/airflow/models/dagbag.py, next to get_dag_for_rundefget_serialized_dag_model_for_run(
self, dag_run: DagRun, *, session: Session
) ->SerializedDagModel|None:
"""Return the SerializedDagModel for the version a run executes against."""ifversion_id:=self._version_from_dag_run(dag_run=dag_run, session=session):
returnself.get_serialized_dag_model(version_id=version_id, session=session)
returnNone

defget_dag_for_run(self, dag_run: DagRun, session: Session) ->SerializedDAG|None:
ifversion_id:=self._version_from_dag_run(dag_run=dag_run, session=session):
returnself._get_dag(version_id=version_id, session=session)
returnNone

Comment threadairflow-core/src/airflow/jobs/triggerer_job_runner.py
@vatsrahul1001

Copy link
Copy Markdown
Contributor

Moving to 3.3.2 as this is still pending code owner review and do not want to rush on merging this without that

@vatsrahul1001

Copy link
Copy Markdown
Contributor

@dstandish@hussein-awala can you review this?

@hussein-awalahussein-awala left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Thanks for the fix, resolving the version from the run is the right direction and it
matches what workers already do. One thing worries me though. This applies to every
trigger, not only the start_from_trigger ones, and that makes the get_task call
able to kill the triggerer.

Here is the scenario. Take a DAG in a bundle without versioning, so its runs have
bundle_version = None.

  1. my_dag v1 has a deferrable task wait. A run is created at v1 and wait
    defers. The TI is DEFERRED with dag_version_id = v1 and there is a Trigger
    row for it.
  2. Someone edits the file and renames wait to wait_for_data. The dag processor
    creates v2 and it becomes the latest. The run is still RUNNING, its
    bundle_version is still None and created_dag_version_id is still v1.
  3. The scheduler does not clean the TI up yet.
    _check_for_removed_or_restored_tasks only marks it REMOVED when
    self.state != DagRunState.RUNNING, so nothing happens until the next
    task_instance_scheduling_decisions pass over that run.
  4. In the meantime the triggerer builds the workload for that trigger, which happens
    on a restart, on a rolling deploy, or on HA failover after assign_unassigned.
    _version_from_dag_run sees no bundle_version, returns v2 as the latest, and
    serialized_dag_model.dag.get_task("wait") raises TaskNotFound.
  5. Nobody catches it. Not _create_workload, not build_trigger_workloads, not
    run_once. TriggererJobRunner.run() logs it and re-raises, so the process exits
    and every other trigger it was running goes down with it. After the restart
    assign_unassigned gives it the same trigger back and it crashes again, until the
    scheduler gets around to marking the TI REMOVED and clean_unused deletes the
    trigger.

Steps 2 and 4 are really the same event in practice. A rolling deploy of new DAG
code is exactly when new versions show up and triggerers restart.

Before this change the version came from the TI, so the task was always there and
step 4 just worked. We already guard this pattern elsewhere. In
task_instances.py the same resolution is followed by
with contextlib.suppress(TaskNotFound).

The simplest fix I can think of is to keep the flag check on the TI version and only
re-resolve for the case the title is about:

serialized_dag_model=dag_bag.get_serialized_dag_model(
version_id=trigger.task_instance.dag_version_id, session=session
)
ifserialized_dag_model:
task=serialized_dag_model.dag.get_task(trigger.task_instance.task_id)
iftask.start_from_trigger:
dag_run=trigger.task_instance.get_dagrun(session=session)
run_version_id=DBDagBag._version_from_dag_run(dag_run=dag_run, session=session)
ifrun_version_idandrun_version_id!=trigger.task_instance.dag_version_id:
serialized_dag_model= (
dag_bag.get_serialized_dag_model(version_id=run_version_id, session=session)
orserialized_dag_model
)

Comment threadairflow-core/src/airflow/jobs/triggerer_job_runner.py
@seanghaeli
seanghaeliforce-pushed the fix/triggerer-pin-dag-version branch from 0a7a1ff to 4ce124cCompareAugust 18, 2026 21:47
@seanghaeli

Copy link
Copy Markdown
ContributorAuthor

@hussein-awala good catch, it can crash in this specific sequence. In fact, the current version of main can fail for a similar reason noted in issue #69841. The latest commit introduces a guard so it falls back to running the trigger without Dag context instead of killing the triggerer

For an unpinned run the resolved (latest) Dag version may no longer
contain a renamed or removed deferred task. get_task() then raised
TaskNotFound, which nothing on the load_triggers path caught, killing
the whole triggerer; assign_unassigned re-handed the same trigger to the
restarted triggerer, producing a crash loop. Guard the lookup and fall
through to the plain workload, mirroring the guard in the execution API.
Also introduce DBDagBag.get_serialized_dag_model_for_run so the
triggerer uses a public API instead of the private _version_from_dag_run
helper, as requested in review.
@seanghaeli
seanghaeliforce-pushed the fix/triggerer-pin-dag-version branch from 4ce124c to 05ebadcCompareAugust 18, 2026 22:46
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:Triggererbackport-to-v3-3-testBackport to v3-3-testtype:bug-fixChangelog: Bug Fixes

Projects

None yet

Development

Successfully merging this pull request may close these issues.

9 participants

@seanghaeli@potiuk@o-nikolas@vatsrahul1001@hussein-awala@jason810496@vincbeck@viiccwen@eladkal
, '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

Load the correct DAG version when a task starts from a trigger - #69988

Open
seanghaeli wants to merge 3 commits into
apache:mainfrom
aws-mwaa:fix/triggerer-pin-dag-version
Open

Load the correct DAG version when a task starts from a trigger#69988
seanghaeli wants to merge 3 commits into
apache:mainfrom
aws-mwaa:fix/triggerer-pin-dag-version

Conversation

@seanghaeli

@seanghaeliseanghaeli commented Jul 16, 2026

Copy link
Copy Markdown
Contributor

If a Dag run without a bundle_version and one of its tasks point at different Dag versions, the trigger should load the latest version. Currently, it loads the task's possibly outdated version.

@seanghaeli
seanghaeliforce-pushed the fix/triggerer-pin-dag-version branch 5 times, most recently from f157049 to ecbf158CompareJuly 17, 2026 00:54
@seanghaeli
seanghaeli marked this pull request as ready for review July 17, 2026 00:59
@seanghaeliseanghaeli changed the title Pin triggerer start_from_trigger DAG resolution to the dagrun's versionLoad the correct DAG version when a task starts from a triggerJul 17, 2026
@seanghaeli

Copy link
Copy Markdown
ContributorAuthor

@o-nikolas this one is similar in spirit to your PR #69941

@o-nikolaso-nikolas left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Seems like a reasonable change semantically. But someone who knows the trigger code better should review to see if this is what we want (it seems like it should be though to me at least). Maybe @vincbeck?

@viiccwenviiccwen left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

left comments.

Comment threadairflow-core/tests/unit/jobs/test_triggerer_job.py Outdated
Comment on lines +879 to +880
version_id=trigger.task_instance.get_dagrun(session=session).created_dag_version_id
or trigger.task_instance.dag_version_id,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

created_dag_version_id is populated even for Dags with disable_bundle_versioning=True; dag_run.bundle_version is what determines
whether the run is pinned. This unconditional preference for created_dag_version_id therefore makes an unpinned start_from_trigger task load the Dag version from when the run was created, even after the scheduler has advanced the unfinished TI's dag_version_id following a reparse.

This differs from DBDagBag._version_from_dag_run(), which intentionally resolves the latest version when bundle_version is absent. Please use the run's created_dag_version_id only when the run is pinned, retain the TI version for unpinned runs, and cover both cases in the regression test.

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.

I believe this is addressed now, could you take a look?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

@viiccwen ☝️

Comment threadairflow-core/tests/unit/jobs/test_triggerer_job.py
@eladkaleladkal added this to the Airflow 3.3.1 milestone Jul 19, 2026
@eladkaleladkal added type:bug-fix Changelog: Bug Fixes backport-to-v3-3-test Backport to v3-3-test labels Jul 19, 2026
@potiuk

Copy link
Copy Markdown
Member

@seanghaeli — There are 3 unresolved review thread(s) on this PR from @viiccwen. Could you either push a fix or reply in each thread explaining why the feedback doesn't apply? Once you believe the feedback is addressed, mark the thread as resolved so the reviewer isn't re-pinged needlessly. Thanks!


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

@o-nikolas

Copy link
Copy Markdown
Contributor

@viiccwen do the changes look good to you now?

@vatsrahul1001

Copy link
Copy Markdown
Contributor

LGTM, I think we are good to merge after code owners review @dstandish@hussein-awala

@jason810496jason810496 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Thanks for the fix.

IMO, we should introduce something like get_serialized_dag_model_for_run rather than calling the private method of the DagBag.

# airflow-core/src/airflow/models/dagbag.py, next to get_dag_for_rundefget_serialized_dag_model_for_run(
self, dag_run: DagRun, *, session: Session
) ->SerializedDagModel|None:
"""Return the SerializedDagModel for the version a run executes against."""ifversion_id:=self._version_from_dag_run(dag_run=dag_run, session=session):
returnself.get_serialized_dag_model(version_id=version_id, session=session)
returnNone

defget_dag_for_run(self, dag_run: DagRun, session: Session) ->SerializedDAG|None:
ifversion_id:=self._version_from_dag_run(dag_run=dag_run, session=session):
returnself._get_dag(version_id=version_id, session=session)
returnNone

Comment threadairflow-core/src/airflow/jobs/triggerer_job_runner.py
@vatsrahul1001

Copy link
Copy Markdown
Contributor

Moving to 3.3.2 as this is still pending code owner review and do not want to rush on merging this without that

@vatsrahul1001

Copy link
Copy Markdown
Contributor

@dstandish@hussein-awala can you review this?

@hussein-awalahussein-awala left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Thanks for the fix, resolving the version from the run is the right direction and it
matches what workers already do. One thing worries me though. This applies to every
trigger, not only the start_from_trigger ones, and that makes the get_task call
able to kill the triggerer.

Here is the scenario. Take a DAG in a bundle without versioning, so its runs have
bundle_version = None.

  1. my_dag v1 has a deferrable task wait. A run is created at v1 and wait
    defers. The TI is DEFERRED with dag_version_id = v1 and there is a Trigger
    row for it.
  2. Someone edits the file and renames wait to wait_for_data. The dag processor
    creates v2 and it becomes the latest. The run is still RUNNING, its
    bundle_version is still None and created_dag_version_id is still v1.
  3. The scheduler does not clean the TI up yet.
    _check_for_removed_or_restored_tasks only marks it REMOVED when
    self.state != DagRunState.RUNNING, so nothing happens until the next
    task_instance_scheduling_decisions pass over that run.
  4. In the meantime the triggerer builds the workload for that trigger, which happens
    on a restart, on a rolling deploy, or on HA failover after assign_unassigned.
    _version_from_dag_run sees no bundle_version, returns v2 as the latest, and
    serialized_dag_model.dag.get_task("wait") raises TaskNotFound.
  5. Nobody catches it. Not _create_workload, not build_trigger_workloads, not
    run_once. TriggererJobRunner.run() logs it and re-raises, so the process exits
    and every other trigger it was running goes down with it. After the restart
    assign_unassigned gives it the same trigger back and it crashes again, until the
    scheduler gets around to marking the TI REMOVED and clean_unused deletes the
    trigger.

Steps 2 and 4 are really the same event in practice. A rolling deploy of new DAG
code is exactly when new versions show up and triggerers restart.

Before this change the version came from the TI, so the task was always there and
step 4 just worked. We already guard this pattern elsewhere. In
task_instances.py the same resolution is followed by
with contextlib.suppress(TaskNotFound).

The simplest fix I can think of is to keep the flag check on the TI version and only
re-resolve for the case the title is about:

serialized_dag_model=dag_bag.get_serialized_dag_model(
version_id=trigger.task_instance.dag_version_id, session=session
)
ifserialized_dag_model:
task=serialized_dag_model.dag.get_task(trigger.task_instance.task_id)
iftask.start_from_trigger:
dag_run=trigger.task_instance.get_dagrun(session=session)
run_version_id=DBDagBag._version_from_dag_run(dag_run=dag_run, session=session)
ifrun_version_idandrun_version_id!=trigger.task_instance.dag_version_id:
serialized_dag_model= (
dag_bag.get_serialized_dag_model(version_id=run_version_id, session=session)
orserialized_dag_model
)

Comment threadairflow-core/src/airflow/jobs/triggerer_job_runner.py
@seanghaeli
seanghaeliforce-pushed the fix/triggerer-pin-dag-version branch from 0a7a1ff to 4ce124cCompareAugust 18, 2026 21:47
@seanghaeli

Copy link
Copy Markdown
ContributorAuthor

@hussein-awala good catch, it can crash in this specific sequence. In fact, the current version of main can fail for a similar reason noted in issue #69841. The latest commit introduces a guard so it falls back to running the trigger without Dag context instead of killing the triggerer

For an unpinned run the resolved (latest) Dag version may no longer
contain a renamed or removed deferred task. get_task() then raised
TaskNotFound, which nothing on the load_triggers path caught, killing
the whole triggerer; assign_unassigned re-handed the same trigger to the
restarted triggerer, producing a crash loop. Guard the lookup and fall
through to the plain workload, mirroring the guard in the execution API.
Also introduce DBDagBag.get_serialized_dag_model_for_run so the
triggerer uses a public API instead of the private _version_from_dag_run
helper, as requested in review.
@seanghaeli
seanghaeliforce-pushed the fix/triggerer-pin-dag-version branch from 4ce124c to 05ebadcCompareAugust 18, 2026 22:46
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:Triggererbackport-to-v3-3-testBackport to v3-3-testtype:bug-fixChangelog: Bug Fixes

Projects

None yet

Development

Successfully merging this pull request may close these issues.

9 participants

@seanghaeli@potiuk@o-nikolas@vatsrahul1001@hussein-awala@jason810496@vincbeck@viiccwen@eladkal