uranusjr commented on code in PR #68359:
URL: https://github.com/apache/airflow/pull/68359#discussion_r3800042960


##########
airflow-core/src/airflow/dag_processing/processor.py:
##########
@@ -424,6 +429,32 @@ def _execute_dag_callbacks(dagbag: DagBag, request: 
DagCallbackRequest, log: Fil
             )
 
 
+def _execute_dag_skipped_intervals_callback(
+    dagbag: DagBag, request: DagSkippedIntervalsCallbackRequest, log: 
FilteringBoundLogger
+) -> None:
+    dag, _ = _get_dag_with_task(dagbag, request.dag_id)
+    callbacks = dag.on_skipped_intervals_callback
+    if not callbacks:
+        log.warning("Skipped intervals callback requested, but dag didn't have 
any", dag_id=request.dag_id)
+        return
+
+    callbacks = callbacks if isinstance(callbacks, list) else [callbacks]
+    summary = request.to_summary()
+    context: SkippedIntervalsCallbackContext = {
+        "dag": dag,
+        "reason": "skipped_intervals",
+        "skipped_range": summary.skipped_range,
+    }
+
+    for callback in callbacks:
+        log.info("Executing on_skipped_intervals_callback", 
dag_id=request.dag_id)
+        try:
+            callback(context)
+        except Exception:
+            log.exception("Callback failed", dag_id=request.dag_id)
+            stats.incr("dag.callback_exceptions", tags={"dag_id": 
request.dag_id})

Review Comment:
   This drops the `team_name` tag that `_execute_dag_callbacks` attaches 
(L417-429), so dag.callback_exceptions will be inconsistently tagged depending 
on which callback failed. Worth aligning.



-- 
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