kaxil commented on code in PR #67592:
URL: https://github.com/apache/airflow/pull/67592#discussion_r4161621930


##########
airflow-core/src/airflow/api_fastapi/execution_api/routes/task_instances.py:
##########
@@ -244,6 +251,12 @@ def ti_run(
                 extra=json.dumps({"host_name": ti_run_payload.hostname}) if 
ti_run_payload.hostname else None,
             )
         )
+        # One sample per queue wait, not per try: the scheduler refreshes 
queued_dttm on every
+        # queueing, so a retry and a resume from deferral each waited for a 
slot of their own.
+        # task.scheduled_duration counts per try instead, so the two disagree 
on retries by design.
+        # queued_dttm is None only in rare races and test setups.

Review Comment:
   "task.scheduled_duration counts per try instead" isn't accurate. 
`emit_state_change_metric` returns early whenever `end_date` is set, so 
`scheduled_duration` skips retries but does emit on deferral resumes. Maybe: 
"task.scheduled_duration skips retries (emit_state_change_metric returns early 
while end_date is set), so the two disagree on retries by design."
   
   On the next line, per your reply to @samraj2k the real case for a missing 
`queued_dttm` is runs that skip the scheduler's queueing, e.g. `dag.test()`, 
rather than rare races.



##########
airflow-core/newsfragments/67592.bugfix.rst:
##########
@@ -0,0 +1 @@
+Restore the ``task.queued_duration`` metric, which stopped being emitted when 
Airflow 3 workers moved to the Task SDK, and record it once per queue wait so 
that a task resuming from deferral also reports the wait for its worker slot.

Review Comment:
   Airflow 2 already reported the queue wait on deferral resumes: `_defer_task` 
never set `end_date`, so `emit_state_change_metric` fired on resume. What 
actually changes from 2.x is retries and reschedule-mode sensor pokes, both of 
which 2.x skipped because they leave `end_date` set (`_handle_reschedule` sets 
it). A reschedule sensor that pokes every minute for an hour goes from one 
sample to about 60. That will move percentiles for anyone upgrading from 2.x 
who alerts on this timer. It fits "once per queue wait", and I left the 
reschedule path out when I listed the transitions earlier, but the newsfragment 
should say so.



##########
airflow-core/src/airflow/api_fastapi/execution_api/routes/task_instances.py:
##########
@@ -299,6 +324,14 @@ def ti_run(
             or 0
         )
 
+        if emit_queued_duration:
+            # Tags mirror the sibling task.scheduled_duration, which 
emit_state_change_metric sends
+            # as {**ti.stats_tags, "queue": ti.queue}; stats_tags reads the 
team off the transient
+            # _team_name. stats.timing also emits the legacy dotted name from 
the metrics registry.
+            dr._team_name = dr.team_name
+            tags = {**dr.stats_tags, "task_id": ti.task_id, "queue": ti.queue}
+            stats.timing("task.queued_duration", timezone.utcnow() - 
ti.queued_dttm, tags=tags)

Review Comment:
   Optional: this emits before code that can still fail, namely the 
`invalid_arg_bindings` 500, `issue_execution_token`, and the commit. The SDK 
client retries 5xx responses and the rollback puts the TI back in QUEUED, so 
the retried request emits a second sample. Moving this block to just before 
`return context` leaves only the commit able to fail after the emit.



##########
airflow-core/tests/unit/api_fastapi/execution_api/versions/head/test_task_instances.py:
##########
@@ -1245,6 +1245,309 @@ def test_ti_run_creates_audit_log(self, client, 
session, create_task_instance, t
         assert logs[0].owner == ti.task.owner
         assert logs[0].extra == '{"host_name": "random-hostname"}'
 
+    @pytest.mark.parametrize(
+        "scenario",
+        ["first_run", "retry", "deferral_resume"],
+    )
+    def test_ti_run_emits_queued_duration_metric(
+        self, client, session, create_task_instance, time_machine, scenario
+    ):
+        """task.queued_duration is emitted once per queue wait.
+
+        The scheduler refreshes queued_dttm every time it queues a task, so a 
retry and a
+        resume from deferral each waited for a worker slot of their own and 
must emit too. A
+        retry still carries the previous attempt's end_date on the row when 
ti_run is reached,
+        so this also asserts the emit does not depend on end_date being unset.
+        """
+        queued_at = timezone.parse("2024-09-30T12:00:00Z")
+        run_at = queued_at.add(seconds=42)
+
+        ti = create_task_instance(
+            task_id=f"test_ti_run_emits_queued_duration_metric_{scenario}",
+            state=State.QUEUED,
+            dagrun_state=DagRunState.RUNNING,
+            session=session,
+            start_date=queued_at,
+            dag_id=str(uuid4()),
+        )
+        ti.queued_dttm = queued_at
+        ti.queue = "default"
+        if scenario == "retry":
+            # A retried TI still has the previous attempt's end_date set on 
the row until
+            # ti_run clears it; the metric must fire regardless.
+            ti.end_date = queued_at.add(seconds=10)
+        elif scenario == "deferral_resume":
+            ti.next_method = "execute_complete"
+        session.commit()
+
+        # The metric has to stay sliceable the same way as its sibling 
task.scheduled_duration,
+        # which emit_state_change_metric sends as {**ti.stats_tags, "queue": 
ti.queue}. Deriving
+        # the expectation from that same expression makes the two drift apart 
only if this fails.
+        expected_tags = {**ti.stats_tags, "queue": ti.queue}
+        assert "run_type" in expected_tags
+
+        time_machine.move_to(run_at, tick=False)
+
+        with 
mock.patch("airflow.api_fastapi.execution_api.routes.task_instances.stats") as 
mock_stats:

Review Comment:
   Nit: this and the other new `task_instances.stats` patches could take 
`autospec=True`, so the `stats.timing` call is checked against the real 
signature.



##########
airflow-core/tests/unit/api_fastapi/execution_api/versions/head/test_task_instances.py:
##########
@@ -1245,6 +1245,309 @@ def test_ti_run_creates_audit_log(self, client, 
session, create_task_instance, t
         assert logs[0].owner == ti.task.owner
         assert logs[0].extra == '{"host_name": "random-hostname"}'
 
+    @pytest.mark.parametrize(
+        "scenario",
+        ["first_run", "retry", "deferral_resume"],
+    )
+    def test_ti_run_emits_queued_duration_metric(
+        self, client, session, create_task_instance, time_machine, scenario
+    ):
+        """task.queued_duration is emitted once per queue wait.
+
+        The scheduler refreshes queued_dttm every time it queues a task, so a 
retry and a
+        resume from deferral each waited for a worker slot of their own and 
must emit too. A
+        retry still carries the previous attempt's end_date on the row when 
ti_run is reached,
+        so this also asserts the emit does not depend on end_date being unset.
+        """
+        queued_at = timezone.parse("2024-09-30T12:00:00Z")
+        run_at = queued_at.add(seconds=42)
+
+        ti = create_task_instance(
+            task_id=f"test_ti_run_emits_queued_duration_metric_{scenario}",
+            state=State.QUEUED,
+            dagrun_state=DagRunState.RUNNING,
+            session=session,
+            start_date=queued_at,
+            dag_id=str(uuid4()),
+        )
+        ti.queued_dttm = queued_at
+        ti.queue = "default"
+        if scenario == "retry":
+            # A retried TI still has the previous attempt's end_date set on 
the row until
+            # ti_run clears it; the metric must fire regardless.
+            ti.end_date = queued_at.add(seconds=10)
+        elif scenario == "deferral_resume":
+            ti.next_method = "execute_complete"
+        session.commit()
+
+        # The metric has to stay sliceable the same way as its sibling 
task.scheduled_duration,
+        # which emit_state_change_metric sends as {**ti.stats_tags, "queue": 
ti.queue}. Deriving
+        # the expectation from that same expression makes the two drift apart 
only if this fails.
+        expected_tags = {**ti.stats_tags, "queue": ti.queue}
+        assert "run_type" in expected_tags
+
+        time_machine.move_to(run_at, tick=False)
+
+        with 
mock.patch("airflow.api_fastapi.execution_api.routes.task_instances.stats") as 
mock_stats:
+            response = client.patch(
+                f"/execution/task-instances/{ti.id}/run",
+                json={
+                    "state": "running",
+                    "hostname": "random-hostname",
+                    "unixname": "random-unixname",
+                    "pid": 100,
+                    "start_date": run_at.isoformat(),
+                },
+            )
+
+        assert response.status_code == 200
+        mock_stats.timing.assert_called_once_with(
+            "task.queued_duration",
+            run_at - queued_at,
+            tags=expected_tags,
+        )
+
+    def test_ti_run_skips_queued_duration_metric_without_queued_dttm(
+        self, client, session, create_task_instance, time_machine
+    ):
+        """queued_dttm is what the wait is measured from, so a row without one 
(rare race /
+        test setups) has nothing to report."""

Review Comment:
   Same as the route comment: the missing `queued_dttm` case is runs that skip 
the scheduler's queueing (e.g. `dag.test()`), not a rare race.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to