Skip to content

Commit 508ee0f

Browse files
committed
Align skipped-interval callback metrics and simplify guard logic
Keep Dag callback exception tags consistent across callback paths and centralize skipped-interval precondition handling in one place to reduce duplicate checks.
1 parent 8493d5a commit 508ee0f

3 files changed

Lines changed: 32 additions & 21 deletions

File tree

airflow-core/src/airflow/dag_processing/processor.py

Lines changed: 13 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -452,7 +452,19 @@ def _execute_dag_skipped_intervals_callback(
452452
callback(context)
453453
except Exception:
454454
log.exception("Callback failed", dag_id=request.dag_id)
455-
stats.incr("dag.callback_exceptions", tags={"dag_id": request.dag_id})
455+
stats.incr(
456+
"dag.callback_exceptions",
457+
tags=prune_dict(
458+
{
459+
"dag_id": request.dag_id,
460+
"team_name": (
461+
DagModel.get_team_name(request.dag_id)
462+
if conf.getboolean("core", "multi_team")
463+
else None
464+
),
465+
}
466+
),
467+
)
456468

457469

458470
def _execute_task_callbacks(dagbag: DagBag, request: TaskCallbackRequest, log: FilteringBoundLogger) -> None:

airflow-core/src/airflow/jobs/scheduler_job_runner.py

Lines changed: 18 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -2692,26 +2692,25 @@ def _create_dag_runs(
26922692
active_non_backfill_runs=active_runs_of_dags[dag_model.dag_id],
26932693
)
26942694

2695-
if data_interval is not None:
2696-
skipped_summary = self._collect_skipped_intervals(
2697-
serdag=serdag,
2698-
new_info=next_info,
2699-
session=session,
2700-
listener_has_impls=listener_has_impls,
2701-
)
2702-
if skipped_summary is not None:
2703-
if listener_has_impls:
2704-
skipped_intervals_listener_events.append((serdag.dag_id, skipped_summary))
2705-
if serdag.has_on_skipped_intervals_callback:
2706-
skip_callback_requests.append(
2707-
DagSkippedIntervalsCallbackRequest.from_summary(
2708-
filepath=dag_model.relative_fileloc or "",
2709-
bundle_name=dag_model.bundle_name,
2710-
bundle_version=created_run.bundle_version,
2711-
dag_id=serdag.dag_id,
2712-
summary=skipped_summary,
2713-
)
2695+
skipped_summary = self._collect_skipped_intervals(
2696+
serdag=serdag,
2697+
new_info=next_info,
2698+
session=session,
2699+
listener_has_impls=listener_has_impls,
2700+
)
2701+
if skipped_summary is not None:
2702+
if listener_has_impls:
2703+
skipped_intervals_listener_events.append((serdag.dag_id, skipped_summary))
2704+
if serdag.has_on_skipped_intervals_callback:
2705+
skip_callback_requests.append(
2706+
DagSkippedIntervalsCallbackRequest.from_summary(
2707+
filepath=dag_model.relative_fileloc or "",
2708+
bundle_name=dag_model.bundle_name,
2709+
bundle_version=created_run.bundle_version,
2710+
dag_id=serdag.dag_id,
2711+
summary=skipped_summary,
27142712
)
2713+
)
27152714

27162715
# Exceptions like ValueError, ParamValidationError, etc. are raised by
27172716
# DagModel.create_dagrun() when dag is misconfigured. The scheduler should not

task-sdk/src/airflow/sdk/execution_time/schema/versions/v2026_10_30.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -44,7 +44,7 @@ class AddArgBindingsToSupervisorTIRunContext(VersionChange):
4444

4545

4646
class AddDagSkippedIntervalsCallbackRequest(VersionChange):
47-
"""Introduce ``DagSkippedIntervalsCallbackRequest`` in ``CallbackRequest`` union."""
47+
"""Introduce ``DagSkippedIntervalsCallbackRequest`` in the ``CallbackRequest`` union."""
4848

4949
description = __doc__
5050

0 commit comments

Comments
 (0)