Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions airflow-core/newsfragments/67927.feature.rst
Original file line number Diff line number Diff line change
@@ -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.
1 change: 1 addition & 0 deletions airflow-core/src/airflow/executors/workloads/trigger.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
27 changes: 24 additions & 3 deletions airflow-core/src/airflow/jobs/triggerer_job_runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
"""
Expand Down Expand Up @@ -552,7 +553,8 @@ def start( # type: ignore[override]
**kwargs,
)

msg = messages.StartTriggerer()
team_name = kwargs.get("team_name")
Comment thread
ferruzzi marked this conversation as resolved.
msg = messages.StartTriggerer(team_name=team_name)
proc.send_msg(msg, request_id=0)
return proc

Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand All @@ -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(
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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(
Expand Down
48 changes: 48 additions & 0 deletions airflow-core/tests/unit/jobs/test_triggerer_job.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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):
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down