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/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..12f3f1d60685a 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.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 @@ 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.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 0f307fc150520..c9704c53c23a8 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.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 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"