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: []