From a8f38c45addb5ecc80b76592f0d05ae680b82daa Mon Sep 17 00:00:00 2001 From: Sameer Mesiah Date: Wed, 3 Jun 2026 00:49:11 +0100 Subject: [PATCH 1/2] Add trigger queue delay metric to TriggerRunner 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/executors/workloads/trigger.py | 1 + .../src/airflow/jobs/triggerer_job_runner.py | 27 +++++++++-- .../tests/unit/jobs/test_triggerer_job.py | 48 +++++++++++++++++++ .../metrics/metrics_template.yaml | 7 +++ 4 files changed, 80 insertions(+), 3 deletions(-) diff --git a/airflow-core/src/airflow/executors/workloads/trigger.py b/airflow-core/src/airflow/executors/workloads/trigger.py index d3b2d0627a7ec..94b372e3cf7e4 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 355763aed5f13..3b8abc4573316 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 StartTriggerer(BaseModel): """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 @@ def start( # type: ignore[override] **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 @@ def update_triggers(self, requested_trigger_ids: set[int]): # 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.monotonic() + + 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 @@ def __init__(self): 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 @@ async def init_comms(self): 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 @@ async def create_triggers(self): 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 @@ async def create_triggers(self): 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.monotonic() - 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 0f307fc150520..1a3fcd580e215 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 send_msg_spy(self, msg, *args, **kwargs): 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 @@ async def test_trigger_kwargs_serialization_cleanup(self, session): 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.monotonic", + 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 1855f737b30ce..a3e5b6bdd7555 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" From 1d70e40d0cc4f1f43a3e2c961fde42c9c3a01ebc Mon Sep 17 00:00:00 2001 From: Sameer Mesiah Date: Sun, 12 Jul 2026 14:52:19 +0100 Subject: [PATCH 2/2] Add newsfragment and change time.monotonic() to time.time(). --- airflow-core/newsfragments/67927.feature.rst | 1 + airflow-core/src/airflow/jobs/triggerer_job_runner.py | 4 ++-- airflow-core/tests/unit/jobs/test_triggerer_job.py | 2 +- 3 files changed, 4 insertions(+), 3 deletions(-) create mode 100644 airflow-core/newsfragments/67927.feature.rst diff --git a/airflow-core/newsfragments/67927.feature.rst b/airflow-core/newsfragments/67927.feature.rst new file mode 100644 index 0000000000000..ca7500343a965 --- /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/jobs/triggerer_job_runner.py b/airflow-core/src/airflow/jobs/triggerer_job_runner.py index 3b8abc4573316..12f3f1d60685a 100644 --- a/airflow-core/src/airflow/jobs/triggerer_job_runner.py +++ b/airflow-core/src/airflow/jobs/triggerer_job_runner.py @@ -979,7 +979,7 @@ def update_triggers(self, requested_trigger_ids: set[int]): if new_trigger_ids: workloads_to_create = self.build_trigger_workloads(new_trigger_ids) - queued_at = time.monotonic() + queued_at = time.time() for workload in workloads_to_create: workload.queued_at = queued_at @@ -1403,7 +1403,7 @@ async def create_triggers(self): if workload.queued_at is not None: stats.timing( "triggerer.trigger_queue_delay", - int((time.monotonic() - workload.queued_at) * 1000), + int((time.time() - workload.queued_at) * 1000), tags=prune_dict({"team_name": self.team_name}), ) diff --git a/airflow-core/tests/unit/jobs/test_triggerer_job.py b/airflow-core/tests/unit/jobs/test_triggerer_job.py index 1a3fcd580e215..c9704c53c23a8 100644 --- a/airflow-core/tests/unit/jobs/test_triggerer_job.py +++ b/airflow-core/tests/unit/jobs/test_triggerer_job.py @@ -1583,7 +1583,7 @@ async def test_create_triggers_emits_queue_delay_metric( runner.to_create.append(workload) with patch( - "airflow.jobs.triggerer_job_runner.time.monotonic", + "airflow.jobs.triggerer_job_runner.time.time", return_value=101.5, ): await runner.create_triggers()