This is an automated email from the ASF dual-hosted git repository.
potiuk pushed a commit to branch v3-3-test
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/v3-3-test by this push:
new e9116a93238 [v3-3-test] Fix N+1 queries in trigger asset event
submission (#65367) (#70738)
e9116a93238 is described below
commit e9116a93238e7e40a765e2e0a4ddbe32c9d638df
Author: Rahul Vats <[email protected]>
AuthorDate: Fri Jul 31 00:56:36 2026 +0530
[v3-3-test] Fix N+1 queries in trigger asset event submission (#65367)
(#70738)
---
airflow-core/src/airflow/models/trigger.py | 6 +++++-
airflow-core/tests/unit/models/test_trigger.py | 22 ++++++++++++++++++++++
2 files changed, 27 insertions(+), 1 deletion(-)
diff --git a/airflow-core/src/airflow/models/trigger.py
b/airflow-core/src/airflow/models/trigger.py
index a7e7b8ee293..597513a3d07 100644
--- a/airflow-core/src/airflow/models/trigger.py
+++ b/airflow-core/src/airflow/models/trigger.py
@@ -285,7 +285,11 @@ class Trigger(Base):
handle_event_submit(event, task_instance=task_instance,
session=session)
# Send an event to assets
- trigger = session.scalars(select(cls).where(cls.id ==
trigger_id)).one_or_none()
+ trigger = session.scalars(
+ select(cls)
+ .where(cls.id == trigger_id)
+
.options(selectinload(cls.asset_watchers).selectinload(AssetWatcherModel.asset))
+ ).one_or_none()
if trigger is None:
# Already deleted for some reason
return
diff --git a/airflow-core/tests/unit/models/test_trigger.py
b/airflow-core/tests/unit/models/test_trigger.py
index 1243a9112f9..669d98a7704 100644
--- a/airflow-core/tests/unit/models/test_trigger.py
+++ b/airflow-core/tests/unit/models/test_trigger.py
@@ -50,6 +50,7 @@ from airflow.triggers.base import (
from airflow.utils.session import create_session
from airflow.utils.state import State
+from tests_common.test_utils.asserts import assert_queries_count
from tests_common.test_utils.config import conf_vars
if TYPE_CHECKING:
@@ -234,6 +235,27 @@ def test_submit_event(mock_callback_handle_event, session,
create_task_instance)
mock_callback_handle_event.assert_called_once_with(event, session)
[email protected](("asset_count", "expected_query_count"), [(1, 6), (5,
6)])
+@patch("airflow.models.trigger.AssetManager.register_asset_change")
+def test_submit_event_no_n_plus_one_for_assets(_, session, asset_count,
expected_query_count):
+ """Ensure asset notifications do not trigger per-asset lazy-load
queries."""
+ trigger = Trigger(classpath="airflow.triggers.testing.SuccessTrigger",
kwargs={})
+ session.add(trigger)
+ session.flush()
+ trigger_id = trigger.id
+
+ for i in range(asset_count):
+ asset = AssetModel(name=f"asset_{asset_count}_{i}")
+ asset.add_trigger(trigger, f"watcher_{i}")
+ session.add(asset)
+
+ session.commit()
+ session.expire_all()
+
+ with assert_queries_count(expected_query_count, session=session):
+ Trigger.submit_event(trigger_id, TriggerEvent("payload"),
session=session)
+
+
def test_submit_failure(session, create_task_instance):
"""
Tests that failures submitted to a trigger fail their dependent