Fix missing task.queued_duration metric in Airflow 3 - #67592
Conversation
|
Congratulations on your first Pull Request and welcome to the Apache Airflow community! If you have any issues or are unsure about any anything please check our Contributors' Guide
|
|
@eladkal gentle ping! Could you help trigger CI? |
0576c5e to
2988b04
Compare
|
@ashb Gentle ping! This has approval from @henry3260, would appreciate a look when you have a moment. The 4 failing checks are unrelated flaky test_wasb failures. |
2988b04 to
5f57b80
Compare
|
@ashb thanks for the review! On the On the test: |
|
@ashb Gentle ping — I've addressed both your comments: dropped the |
5f57b80 to
948cbcd
Compare
task.queued_duration (and its registry-derived legacy name dag.<dag_id>.<task_id>.queued_duration) stopped firing after the Airflow 3 worker switched to the Task SDK / supervisor / Execution API. The only emit site was TaskInstance.emit_state_change_metric, reachable solely from the legacy LocalTaskJob path, while Airflow 3 workers flip the TI to RUNNING through the ti_run Execution API endpoint instead. Emit it from ti_run on the genuine QUEUED -> RUNNING transition, tagged via DagRun.stats_tags so the restored metric stays sliceable the same way as its sibling task.scheduled_duration (dag_id, task_id, queue, run_type, team_name under multi-team, and Dag tags when configured). Resumes from deferral are skipped so a deferral cycle does not add a second sample inside the same try. closes: apache#63503 closes: apache#66067
948cbcd to
0f3fb97
Compare
|
Moving to 3.3.2 as this is still pending code owner review and do not want to rush on merging this as it not critical |
| # which fired on every transition to RUNNING. Only resumes from deferral are skipped, | ||
| # identified by next_method (the trigger sets it on resume), to avoid re-emitting within the | ||
| # same try. queued_dttm is None only in rare races and test setups. | ||
| emit_queued_duration = ti.queued_dttm is not None and ti.next_method is None |
There was a problem hiding this comment.
This guard came in at @henry3260's request and it does what it says, but it leaves task.queued_duration and its sibling task.scheduled_duration disagreeing about which transitions count, in opposite directions. Taking the three ways a TI reaches RUNNING:
- First run: both metrics emit.
- Retry: only
queued_duration.schedule_tissetsSCHEDULED,scheduled_dttmandtry_numberwithout clearingend_date, so the previous attempt'send_dateis still on the row when the scheduler queues the TI, andemit_state_change_metricreturns early. - Deferral resume: only
scheduled_duration. TheDEFERREDupdate never touchesend_dateandti_runalready set it toNone, so that guard passes, and the trigger refreshesscheduled_dttm, so the sample is a real measurement.
The resume's queue wait is equally real: queued_dttm is refreshed on every queueing and the critical section selects SCHEDULED with no next_method filter, so a deferrable task sits in QUEUED again waiting for a worker slot before execute_complete runs. Skipping it means that wait is never measured, and for sensor-heavy deployments the resume leg is where most of the queue time lives.
I would drop and ti.next_method is None and emit per queue wait, which is also where @ashb's end_date change pointed: one sample per real wait rather than one per try. The per-try reading is defensible too, but then the retry case should not emit either, and the comment should name the axis so the asymmetry reads as deliberate.
Either answer works for me, and this is the only thing I would like settled before merge. Adding these samples after release shifts percentiles for anyone alerting on the timer, which is cheap to decide now and awkward to change later.
On why this did not come up in my last pass: I was looking at the tag set then, and only walked the resume path through the scheduler this time.
There was a problem hiding this comment.
start_from_trigger=True operators are deferred straight from SCHEDULED, never queued, so
the resume this guard skips is their only queue wait — they emit nothing at all today.
Test added.
Also fixed my comment: legacy skipped retries and emitted on resume (it runs before
end_date is cleared), the inverse of what I claimed.
Per queue wait is the only axis that measures those, so that's my vote — your call.
| # it here; add the value looked up above instead. Falsy values are pruned from | ||
| # stats_tags, so only set it when there is a team. The registry-derived legacy name | ||
| # dag.<dag_id>.<task_id>.queued_duration is emitted by stats.timing automatically. | ||
| tags = {**dr.stats_tags, "task_id": ti.task_id, "queue": ti.queue} |
There was a problem hiding this comment.
When metrics.dag_tags_in_metrics is on, dr.stats_tags costs two lazy loads per task start here: dag_tags_for_stats touches self.dag_model and then dag_model.tags, and the dr select above only eager-loads consumed_asset_events. The earlier query does join DagModel, but it selects DagModel.owners as a column rather than the entity, so nothing lands in the identity map and the lazy load still fires.
get_running_dag_runs_to_examine eager-loads exactly this to keep it out of the scheduler loop, and the comment on dag_tags_for_stats calls the remaining paths low frequency. ti_run is once per task start on the execution API, so it belongs with the first group. Adding the same eager load to the dr select, conditional on the config so the join is not paid when the feature is off, would cover it.
Non-blocking, and the config is off by default.
There was a problem hiding this comment.
Agreed — the earlier select takes DagModel.owners as a column, not the entity. Left out
to keep the diff on the blocking question; happy to add it here or in a follow-up.
| # dag.<dag_id>.<task_id>.queued_duration is emitted by stats.timing automatically. | ||
| tags = {**dr.stats_tags, "task_id": ti.task_id, "queue": ti.queue} | ||
| if team_name: | ||
| tags["team_name"] = team_name |
There was a problem hiding this comment.
On the _team_name point from your last reply: it is the convention the scheduler already uses for exactly this, ti.dag_run._team_name = team plus two more sites, and stats_tags reads it through getattr. Setting dr._team_name = team_name alongside the dr.team_name = team_name above would make these three lines and the comment redundant, and keeps the tag set derived in one place if stats_tags grows again. Optional, since what you have produces the same tags today.
On the test: yes please, pin the multi-team leg with conf_vars. Your read of the gap is right, and get_team_name_for_ti returns None unless core.multi_team is on, so today that key agrees for the wrong reason on both sides.
There was a problem hiding this comment.
Test added with conf_vars, and it fails if the team tag is dropped.
_team_name: agreed — left out since it touches the same lines as the axis decision.
With multi-team off the tag is absent from both the metric and any expectation derived from stats_tags, so the assertion held for the wrong reason and would not have caught the tag going missing.
The comment claimed to mirror a legacy emit that fired on every transition to RUNNING; that emit skipped retries and fired on deferral resumes, so the claim was inverted on both paths. It also credited the trigger with setting next_method, which the DEFERRED transition does, and asserted one sample per try without accounting for operators deferred before they are ever queued.
These operators are deferred straight from SCHEDULED without being queued, so the resume the next_method guard skips is the only queue wait they ever have.
Summary
task.queued_duration(and its registry-derived legacy namedag.<dag_id>.<task_id>.queued_duration) stopped firing entirely after the Airflow 3 worker switched to the Task SDK / supervisor / Execution API.The metric was only emitted by
TaskInstance.emit_state_change_metric, which is only reachable from_check_and_change_state_before_execution— the legacy LocalTaskJob path. Airflow 3 workers flip TI state toRUNNINGthrough theti_runExecution API endpoint instead, which bypasses the emit site.This is the same regression pattern as #62019 (missing
ti.start/ti.finish).Fix
Emit
task.queued_durationfromti_runat the moment it transitions the TI from QUEUED to RUNNING. Skip the emit on transitions that are not the genuine first run of a try:end_date;next_method(the trigger sets it andend_datestaysNoneacross a deferral, soend_datealone cannot catch it);queued_dttm(rare race / test setups).The legacy dotted name is emitted automatically by
stats.timingvia themetrics_template.yamlregistry — no manual second call needed.Test plan
test_ti_run_emits_queued_duration_metricconfirmed to fail before the fix and pass after (verified by stashing the production change and re-running the test).test_ti_run_skips_queued_duration_metriccovers all three skip conditions (end_dateset /deferral_resume/queued_dttmmissing).TestTIRunStatetests still pass.ruff format/ruff check/mypy-airflow-core/prek run --from-ref upstream/main --stage pre-commitall green.closes: #63503
closes: #66067
Was generative AI tooling used to co-author this PR?
Generated-by: Claude Code (Opus 4.7) following the guidelines
Important
🛠️ Maintainer triage note for @myps6415 · by
@potiuk· 2026-07-02 17:46 UTCSome review feedback from
@ashb,@henry3260is waiting on you:@ashb,@henry3260need a reply or a fix.The ball is in your court — you've been assigned to this PR. Reply or push a fix in each thread, then mark them resolved.
Automated triage — may be imperfect; a maintainer takes the next look.