This is an automated email from the ASF dual-hosted git repository.

shahar1 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git


The following commit(s) were added to refs/heads/main by this push:
     new 5930c02ad0f TriggerRunnerSupervisor: Add trigger queue delay metric 
(#67927)
5930c02ad0f is described below

commit 5930c02ad0fc058a6a40ad37b2436105fd7fb51f
Author: SameerMesiah97 <[email protected]>
AuthorDate: Wed Jul 15 21:30:49 2026 +0100

    TriggerRunnerSupervisor: Add trigger queue delay metric (#67927)
    
    Add a trigger_queue_delay timing metric to measure the time between a 
trigger workload being queued by the TriggerRunnerSupervisor and being 
scheduled by the TriggerRunner. Also propagate team_name to TriggerRunner so 
the metric is emitted with the expected tags, and add unit tests.
---
 airflow-core/newsfragments/67927.feature.rst       |  1 +
 .../src/airflow/executors/workloads/trigger.py     |  1 +
 .../src/airflow/jobs/triggerer_job_runner.py       | 27 ++++++++++--
 airflow-core/tests/unit/jobs/test_triggerer_job.py | 48 ++++++++++++++++++++++
 .../observability/metrics/metrics_template.yaml    |  7 ++++
 5 files changed, 81 insertions(+), 3 deletions(-)

diff --git a/airflow-core/newsfragments/67927.feature.rst 
b/airflow-core/newsfragments/67927.feature.rst
new file mode 100644
index 00000000000..ca7500343a9
--- /dev/null
+++ b/airflow-core/newsfragments/67927.feature.rst
@@ -0,0 +1 @@
+Added the ``triggerer.trigger_queue_delay`` metric, which measures the time a 
trigger workload spends queued before being processed by the trigger runner.
diff --git a/airflow-core/src/airflow/executors/workloads/trigger.py 
b/airflow-core/src/airflow/executors/workloads/trigger.py
index d3b2d0627a7..94b372e3cf7 100644
--- a/airflow-core/src/airflow/executors/workloads/trigger.py
+++ b/airflow-core/src/airflow/executors/workloads/trigger.py
@@ -49,3 +49,4 @@ class RunTrigger(BaseModel):
 
     # name: uri of all "watched" Assets
     watched_assets: dict[str, str] | None = None  # Set for BaseEventTrigger 
asset watchers only
+    queued_at: float | None = None
diff --git a/airflow-core/src/airflow/jobs/triggerer_job_runner.py 
b/airflow-core/src/airflow/jobs/triggerer_job_runner.py
index 355763aed5f..12f3f1d6068 100644
--- a/airflow-core/src/airflow/jobs/triggerer_job_runner.py
+++ b/airflow-core/src/airflow/jobs/triggerer_job_runner.py
@@ -313,6 +313,7 @@ class messages:
         """Tell the async trigger runner process to start, and where to send 
status update messages."""
 
         type: Literal["StartTriggerer"] = "StartTriggerer"
+        team_name: str | None = None
 
     class TriggerStateChanges(BaseModel):
         """
@@ -552,7 +553,8 @@ class TriggerRunnerSupervisor(WatchedSubprocess):
             **kwargs,
         )
 
-        msg = messages.StartTriggerer()
+        team_name = kwargs.get("team_name")
+        msg = messages.StartTriggerer(team_name=team_name)
         proc.send_msg(msg, request_id=0)
         return proc
 
@@ -973,8 +975,16 @@ class TriggerRunnerSupervisor(WatchedSubprocess):
         # Work out the two difference sets
         new_trigger_ids = requested_trigger_ids - known_trigger_ids
         cancel_trigger_ids = self.running_triggers - requested_trigger_ids
+
         if new_trigger_ids:
-            
self.creating_triggers.extend(self.build_trigger_workloads(new_trigger_ids))
+            workloads_to_create = self.build_trigger_workloads(new_trigger_ids)
+
+            queued_at = time.time()
+
+            for workload in workloads_to_create:
+                workload.queued_at = queued_at
+
+            self.creating_triggers.extend(workloads_to_create)
 
         if cancel_trigger_ids:
             # Enqueue orphaned triggers for cancellation
@@ -1171,6 +1181,9 @@ class TriggerRunner:
     # Outbound queue of failed triggers
     failed_triggers: deque[tuple[int, BaseException | None]]
 
+    # Team associated with this triggerer instance.
+    team_name: str | None
+
     # Should-we-stop flag
     stop: bool = False
     _stop_event: anyio.Event | None = None
@@ -1188,6 +1201,7 @@ class TriggerRunner:
         self.to_cancel = deque()
         self.events = deque()
         self.failed_triggers = deque()
+        self.team_name = None
         self.job_id = None
         self._stop_event = None
         self._shared_streams = SharedStreamManager(
@@ -1306,6 +1320,8 @@ class TriggerRunner:
         if not isinstance(msg, messages.StartTriggerer):
             raise RuntimeError(f"Required first message to be a 
messages.StartTriggerer, it was {msg}")
 
+        self.team_name = msg.team_name
+
         await self.comms_decoder.start_reader()
 
     @classmethod
@@ -1337,7 +1353,6 @@ class TriggerRunner:
             if trigger_id in self.triggers:
                 self.log.warning("Trigger %s had insertion attempted twice", 
trigger_id)
                 continue
-
             try:
                 trigger_class = 
self.get_trigger_by_classpath(workload.classpath)
             except BaseException as e:
@@ -1385,6 +1400,12 @@ class TriggerRunner:
                 trigger_instance.asset_state_store = AssetStateStoreAccessors(
                     inlets=[Asset(name=name, uri=uri) for name, uri in 
workload.watched_assets.items()]
                 )
+            if workload.queued_at is not None:
+                stats.timing(
+                    "triggerer.trigger_queue_delay",
+                    int((time.time() - workload.queued_at) * 1000),
+                    tags=prune_dict({"team_name": self.team_name}),
+                )
 
             self.triggers[trigger_id] = {
                 "task": asyncio.create_task(
diff --git a/airflow-core/tests/unit/jobs/test_triggerer_job.py 
b/airflow-core/tests/unit/jobs/test_triggerer_job.py
index 0f307fc1505..c9704c53c23 100644
--- a/airflow-core/tests/unit/jobs/test_triggerer_job.py
+++ b/airflow-core/tests/unit/jobs/test_triggerer_job.py
@@ -1003,6 +1003,7 @@ def test_trigger_lifecycle(spy_agency: SpyAgency, 
session, testing_dag_bundle):
                 encrypted_kwargs=trigger_orm.encrypted_kwargs,
                 kind="RunTrigger",
                 dag_data=ANY,
+                queued_at=ANY,
             )
         )
         # OK, now remove it from the DB
@@ -1546,6 +1547,53 @@ class TestTriggerRunner:
         trigger_instance.cancel()
         await runner.cleanup_finished_triggers()
 
+    @pytest.mark.asyncio
+    @pytest.mark.parametrize(
+        ("team_name", "expected_tags"),
+        [
+            pytest.param("team_a", {"team_name": "team_a"}, id="with_team"),
+            pytest.param(None, {}, id="without_team"),
+        ],
+    )
+    @patch("airflow.jobs.triggerer_job_runner.stats.timing")
+    @patch("airflow.jobs.triggerer_job_runner.Trigger._decrypt_kwargs")
+    @patch(
+        
"airflow.jobs.triggerer_job_runner.TriggerRunner.get_trigger_by_classpath",
+        return_value=DateTimeTrigger,
+    )
+    async def test_create_triggers_emits_queue_delay_metric(
+        self,
+        mock_get_trigger_by_classpath,
+        mock_decrypt_kwargs,
+        mock_timing,
+        team_name,
+        expected_tags,
+    ):
+        mock_decrypt_kwargs.return_value = {"moment": timezone.utcnow() + 
datetime.timedelta(hours=1)}
+
+        workload = workloads.RunTrigger.model_construct(
+            id=1,
+            classpath="abc",
+            encrypted_kwargs="fake",
+            queued_at=100.0,
+        )
+
+        runner = TriggerRunner()
+        runner.team_name = team_name
+        runner.to_create.append(workload)
+
+        with patch(
+            "airflow.jobs.triggerer_job_runner.time.time",
+            return_value=101.5,
+        ):
+            await runner.create_triggers()
+
+        mock_timing.assert_called_once_with(
+            "triggerer.trigger_queue_delay",
+            1500,
+            tags=expected_tags,
+        )
+
     @pytest.mark.asyncio
     @patch("airflow.sdk.execution_time.task_runner.SUPERVISOR_COMMS", 
create=True)
     async def test_sync_state_to_supervisor(self, supervisor_builder):
diff --git 
a/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml
 
b/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml
index 1855f737b30..a3e5b6bdd75 100644
--- 
a/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml
+++ 
b/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml
@@ -723,6 +723,13 @@ metrics:
     legacy_name: "-"
     name_variables: []
 
+  - name: "triggerer.trigger_queue_delay"
+    description: "Time in milliseconds between a trigger workload being queued 
and being processed by
+      the TriggerRunner."
+    type: "timer"
+    legacy_name: "-"
+    name_variables: []
+
   - name: "dagrun.first_task_scheduling_delay"
     description: "Milliseconds elapsed between first task start_date and 
dagrun expected start"
     type: "timer"

Reply via email to