Skip to content

Commit 20546cd

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 fbed9e5 commit 20546cd

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
@@ -2704,26 +2704,25 @@ def _create_dag_runs(
27042704
active_non_backfill_runs=active_runs_of_dags[dag_model.dag_id],
27052705
)
27062706

2707-
if data_interval is not None:
2708-
skipped_summary = self._collect_skipped_intervals(
2709-
serdag=serdag,
2710-
new_info=next_info,
2711-
session=session,
2712-
listener_has_impls=listener_has_impls,
2713-
)
2714-
if skipped_summary is not None:
2715-
if listener_has_impls:
2716-
skipped_intervals_listener_events.append((serdag.dag_id, skipped_summary))
2717-
if serdag.has_on_skipped_intervals_callback:
2718-
skip_callback_requests.append(
2719-
DagSkippedIntervalsCallbackRequest.from_summary(
2720-
filepath=dag_model.relative_fileloc or "",
2721-
bundle_name=dag_model.bundle_name,
2722-
bundle_version=created_run.bundle_version,
2723-
dag_id=serdag.dag_id,
2724-
summary=skipped_summary,
2725-
)
2707+
skipped_summary = self._collect_skipped_intervals(
2708+
serdag=serdag,
2709+
new_info=next_info,
2710+
session=session,
2711+
listener_has_impls=listener_has_impls,
2712+
)
2713+
if skipped_summary is not None:
2714+
if listener_has_impls:
2715+
skipped_intervals_listener_events.append((serdag.dag_id, skipped_summary))
2716+
if serdag.has_on_skipped_intervals_callback:
2717+
skip_callback_requests.append(
2718+
DagSkippedIntervalsCallbackRequest.from_summary(
2719+
filepath=dag_model.relative_fileloc or "",
2720+
bundle_name=dag_model.bundle_name,
2721+
bundle_version=created_run.bundle_version,
2722+
dag_id=serdag.dag_id,
2723+
summary=skipped_summary,
27262724
)
2725+
)
27272726

27282727
# Exceptions like ValueError, ParamValidationError, etc. are raised by
27292728
# 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)