This is an automated email from the ASF dual-hosted git repository.

msumit pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git


The following commit(s) were added to refs/heads/main by this push:
     new ac88accf4ed Add bundle_name tag to missing dag processor metrics 
(#73717)
ac88accf4ed is described below

commit ac88accf4edc6ed70557d5e9554e36adc04ed131
Author: Sameer Raj <[email protected]>
AuthorDate: Tue Sep 29 10:23:03 2026 +0530

    Add bundle_name tag to missing dag processor metrics (#73717)
    
    * Attribute Dag processor timeouts to bundles
    
    Dag bundles can contain the same relative file path. Without bundle 
attribution, timeout counts for those files collapse into one metric series.
    
    * Test normalized bundle tags for processor timeouts
    
    Bundle names can contain characters that require normalization before they 
are used as metric tags. The regression case ensures the timeout metric keeps 
the same tag format as other Dag processing metrics.
    
    * Attribute Dag processor process metrics to bundles
    
    Process counts for files with the same relative path can collapse across 
bundles. Consistent bundle tags at every process lifecycle emission keep those 
series distinct.
    
    * Avoid ordering assumptions in processor metric test
    
    The metric contract requires distinct bundle tags for identical relative 
paths, regardless of the order in which the packets are emitted.
---
 airflow-core/src/airflow/dag_processing/manager.py | 19 ++++-
 .../tests/unit/dag_processing/test_manager.py      | 92 ++++++++++++++++------
 .../observability/metrics/metrics_template.yaml    |  5 +-
 3 files changed, 87 insertions(+), 29 deletions(-)

diff --git a/airflow-core/src/airflow/dag_processing/manager.py 
b/airflow-core/src/airflow/dag_processing/manager.py
index 16edef45603..c6ce3c2dcbe 100644
--- a/airflow-core/src/airflow/dag_processing/manager.py
+++ b/airflow-core/src/airflow/dag_processing/manager.py
@@ -1246,6 +1246,7 @@ class DagFileProcessorManager(LoggingMixin):
                     tags=prune_dict(
                         {
                             "file_path": file.normalized_file_path_for_stats,
+                            "bundle_name": 
normalize_name_for_stats(file.bundle_name),
                             "action": "stop",
                             "team_name": bundle_to_team.get(file.bundle_name),
                         }
@@ -1477,6 +1478,7 @@ class DagFileProcessorManager(LoggingMixin):
                 tags=prune_dict(
                     {
                         "file_path": file.normalized_file_path_for_stats,
+                        "bundle_name": 
normalize_name_for_stats(file.bundle_name),
                         "action": "start",
                         "team_name": bundle_to_team.get(file.bundle_name),
                     }
@@ -1660,16 +1662,28 @@ class DagFileProcessorManager(LoggingMixin):
                     self.processor_timeout,
                 )
                 file_path_tag = file.normalized_file_path_for_stats
+                bundle_name_tag = normalize_name_for_stats(file.bundle_name)
                 team_name = bundle_to_team.get(file.bundle_name)
                 stats.decr(
                     "dag_processing.processes",
                     tags=prune_dict(
-                        {"file_path": file_path_tag, "action": "timeout", 
"team_name": team_name}
+                        {
+                            "file_path": file_path_tag,
+                            "bundle_name": bundle_name_tag,
+                            "action": "timeout",
+                            "team_name": team_name,
+                        }
                     ),
                 )
                 stats.incr(
                     "dag_processing.processor_timeouts",
-                    tags=prune_dict({"file_path": file_path_tag, "team_name": 
team_name}),
+                    tags=prune_dict(
+                        {
+                            "file_path": file_path_tag,
+                            "bundle_name": bundle_name_tag,
+                            "team_name": team_name,
+                        }
+                    ),
                 )
                 processor.kill(signal.SIGKILL)
 
@@ -1739,6 +1753,7 @@ class DagFileProcessorManager(LoggingMixin):
                 tags=prune_dict(
                     {
                         "file_path": file.normalized_file_path_for_stats,
+                        "bundle_name": 
normalize_name_for_stats(file.bundle_name),
                         "action": "terminate",
                         "team_name": bundle_to_team.get(file.bundle_name),
                     }
diff --git a/airflow-core/tests/unit/dag_processing/test_manager.py 
b/airflow-core/tests/unit/dag_processing/test_manager.py
index fbdc9bc9b75..b581a3d742d 100644
--- a/airflow-core/tests/unit/dag_processing/test_manager.py
+++ b/airflow-core/tests/unit/dag_processing/test_manager.py
@@ -596,7 +596,7 @@ class TestDagFileProcessorManager:
 
         stats_incr_mock.assert_called_once_with(
             "dag_processing.processes",
-            tags={"file_path": "folder_file_2.py", "action": "start"},
+            tags={"file_path": "folder_file_2.py", "bundle_name": "testing", 
"action": "start"},
         )
 
         # Because of the config: '[dag_processor] parsing_processes = 2'
@@ -701,7 +701,7 @@ class TestDagFileProcessorManager:
         assert manager._processors == {}
         stats_decr_mock.assert_called_once_with(
             "dag_processing.processes",
-            tags={"file_path": "callbacks_with_spaces.py", "action": "stop"},
+            tags={"file_path": "callbacks_with_spaces.py", "bundle_name": 
"testing", "action": "stop"},
         )
         processor.kill.assert_called_once_with(signal.SIGKILL)
 
@@ -1614,33 +1614,55 @@ class TestDagFileProcessorManager:
 
         assert call_order == ["kill", "close"]
 
-    def test_kill_timed_out_processors_kill(self):
+    def 
test_kill_timed_out_processors_tags_same_file_in_different_bundles(self):
         manager = DagFileProcessorManager(max_runs=1, processor_timeout=5)
         # Set start_time to ensure timeout occurs: start_time = current_time - 
(timeout + 1) = always (timeout + 1) seconds
         start_time = time.monotonic() - manager.processor_timeout - 1
-        processor, _ = self.mock_processor(start_time=start_time)
+        processor_a, _ = self.mock_processor(start_time=start_time)
+        processor_b, _ = self.mock_processor(start_time=start_time)
+        rel_path = Path("folder/abc txt.py")
         manager._processors = {
-            DagFileInfo(
-                bundle_name="testing", rel_path=Path("folder/abc txt.py"), 
bundle_path=TEST_DAGS_FOLDER
-            ): processor
+            DagFileInfo(bundle_name="bundle_a", rel_path=rel_path, 
bundle_path=TEST_DAGS_FOLDER): processor_a,
+            DagFileInfo(bundle_name="bundle b", rel_path=rel_path, 
bundle_path=TEST_DAGS_FOLDER): processor_b,
         }
         with (
-            mock.patch.object(type(processor), "kill") as mock_kill,
+            mock.patch.object(type(processor_a), "kill") as mock_kill,
             mock.patch("airflow.dag_processing.manager.stats.decr") as 
stats_decr_mock,
             mock.patch("airflow.dag_processing.manager.stats.incr") as 
stats_incr_mock,
         ):
             manager._kill_timed_out_processors()
-        mock_kill.assert_called_once_with(signal.SIGKILL)
-        stats_decr_mock.assert_called_once_with(
-            "dag_processing.processes",
-            tags={"file_path": "folder_abc_txt.py", "action": "timeout"},
+        assert mock_kill.call_args_list == [mock.call(signal.SIGKILL)] * 2
+        assert stats_decr_mock.call_count == 2
+        stats_decr_mock.assert_has_calls(
+            [
+                mock.call(
+                    "dag_processing.processes",
+                    tags={"file_path": "folder_abc_txt.py", "bundle_name": 
"bundle_a", "action": "timeout"},
+                ),
+                mock.call(
+                    "dag_processing.processes",
+                    tags={"file_path": "folder_abc_txt.py", "bundle_name": 
"bundle_b", "action": "timeout"},
+                ),
+            ],
+            any_order=True,
         )
-        stats_incr_mock.assert_called_once_with(
-            "dag_processing.processor_timeouts",
-            tags={"file_path": "folder_abc_txt.py"},
+        assert stats_incr_mock.call_count == 2
+        stats_incr_mock.assert_has_calls(
+            [
+                mock.call(
+                    "dag_processing.processor_timeouts",
+                    tags={"file_path": "folder_abc_txt.py", "bundle_name": 
"bundle_a"},
+                ),
+                mock.call(
+                    "dag_processing.processor_timeouts",
+                    tags={"file_path": "folder_abc_txt.py", "bundle_name": 
"bundle_b"},
+                ),
+            ],
+            any_order=True,
         )
-        assert len(manager._processors) == 0
-        processor.logger_filehandle.close.assert_called()
+        assert not manager._processors
+        processor_a.logger_filehandle.close.assert_called_once_with()
+        processor_b.logger_filehandle.close.assert_called_once_with()
 
     def 
test_kill_timed_out_processors_tolerates_stale_file_handle_on_close(self):
         """A stale NFS file handle on close (e.g. OpenShift) must not crash 
the manager."""
@@ -1692,7 +1714,7 @@ class TestDagFileProcessorManager:
 
         stats_decr_mock.assert_called_once_with(
             "dag_processing.processes",
-            tags={"file_path": "folder_abc_txt.py", "action": "terminate"},
+            tags={"file_path": "folder_abc_txt.py", "bundle_name": "testing", 
"action": "terminate"},
         )
         processor.kill.assert_called_once_with(signal.SIGTERM, 
escalation_delay=5.0)
 
@@ -4155,7 +4177,12 @@ class TestMultiTeamMetrics:
 
         mock_incr.assert_any_call(
             "dag_processing.processes",
-            tags={"file_path": "dag_file.py", "action": "start", "team_name": 
"team_alpha"},
+            tags={
+                "file_path": "dag_file.py",
+                "bundle_name": "testing",
+                "action": "start",
+                "team_name": "team_alpha",
+            },
         )
 
     @conf_vars({("core", "multi_team"): "true"})
@@ -4178,11 +4205,16 @@ class TestMultiTeamMetrics:
 
         mock_decr.assert_called_once_with(
             "dag_processing.processes",
-            tags={"file_path": "dag_file.py", "action": "timeout", 
"team_name": "team_alpha"},
+            tags={
+                "file_path": "dag_file.py",
+                "bundle_name": "testing",
+                "action": "timeout",
+                "team_name": "team_alpha",
+            },
         )
         mock_incr.assert_any_call(
             "dag_processing.processor_timeouts",
-            tags={"file_path": "dag_file.py", "team_name": "team_alpha"},
+            tags={"file_path": "dag_file.py", "bundle_name": "testing", 
"team_name": "team_alpha"},
         )
 
     @conf_vars({("core", "multi_team"): "true"})
@@ -4238,13 +4270,18 @@ class TestMultiTeamMetrics:
             pytest.param(
                 True,
                 "team_alpha",
-                {"file_path": "dag_file.py", "action": "stop", "team_name": 
"team_alpha"},
+                {
+                    "file_path": "dag_file.py",
+                    "bundle_name": "testing",
+                    "action": "stop",
+                    "team_name": "team_alpha",
+                },
                 id="with_team",
             ),
             pytest.param(
                 False,
                 None,
-                {"file_path": "dag_file.py", "action": "stop"},
+                {"file_path": "dag_file.py", "bundle_name": "testing", 
"action": "stop"},
                 id="without_team",
             ),
         ],
@@ -4279,13 +4316,18 @@ class TestMultiTeamMetrics:
             pytest.param(
                 True,
                 "team_alpha",
-                {"file_path": "dag_file.py", "action": "terminate", 
"team_name": "team_alpha"},
+                {
+                    "file_path": "dag_file.py",
+                    "bundle_name": "testing",
+                    "action": "terminate",
+                    "team_name": "team_alpha",
+                },
                 id="with_team",
             ),
             pytest.param(
                 False,
                 None,
-                {"file_path": "dag_file.py", "action": "terminate"},
+                {"file_path": "dag_file.py", "bundle_name": "testing", 
"action": "terminate"},
                 id="without_team",
             ),
         ],
diff --git 
a/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml
 
b/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml
index 6cbcf3dc304..fde9ce33205 100644
--- 
a/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml
+++ 
b/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml
@@ -124,14 +124,15 @@ metrics:
 
   - name: "dag_processing.processes"
     description: "Relative number of currently running Dag parsing processes 
(ie this delta is negative when,
-    since the last metric was sent, processes have completed). Metric with 
file_path and action tagging."
+    since the last metric was sent, processes have completed). Metric with 
file_path, bundle_name, and action
+    tagging, and team_name tagging in multi-team mode."
     type: "counter"
     legacy_name: "-"
     name_variables: []
 
   - name: "dag_processing.processor_timeouts"
     description: "Number of file processors that have been killed due to 
taking too long.
-    Metric with file_path tagging."
+    Metric with file_path and bundle_name tagging, and team_name tagging in 
multi-team mode."
     type: "counter"
     legacy_name: "-"
     name_variables: []

Reply via email to