This is an automated email from the ASF dual-hosted git repository.
Abacn pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git
The following commit(s) were added to refs/heads/master by this push:
new 9f5cd57a009 [Python] Honor disableCounterMetrics,
disableStringSetMetrics and disableBoundedTrieMetrics experiments (#40165)
9f5cd57a009 is described below
commit 9f5cd57a0096cb0018ec02f31c010e366192fca9
Author: Anish Mehta <[email protected]>
AuthorDate: Wed Sep 30 09:08:41 2026 +0530
[Python] Honor disableCounterMetrics, disableStringSetMetrics and
disableBoundedTrieMetrics experiments (#40165)
* [Python] Honor disableCounterMetrics, disableStringSetMetrics and
disableBoundedTrieMetrics experiments
Re-lands the feature from #38749 (reverted in #38901) without changing the
public surface of the metric objects. The Java SDK lets high throughput jobs
turn off metric kinds that pressure the metrics backend via these
experiments;
the Python SDK now does the same.
Instead of replacing DelegatingCounter.inc / DelegatingStringSet.add /
DelegatingBoundedTrie.add with methods (which broke callers passing the
value
as a keyword and code relying on them being MetricUpdater instances), the
gate
lives in MetricUpdater.__call__ and is keyed by cell type. With no
experiment
set it is a truthiness check on an empty set, so the hot path is unchanged.
MetricsFlag.set_default_pipeline_options mirrors the Java MetricsFlag and is
applied when a Pipeline is constructed and at SDK worker harness start-up,
first call wins as in Java.
Fixes #38746
Co-Authored-By: Claude Opus 5 (1M context) <[email protected]>
* Check only this test's key in the shared process-wide container
The process-wide metrics container is shared with every other test in the
pytest-xdist worker process, so asserting it is empty fails whenever
another test (the GCP suites in the cloud tox env) registered a
process-wide counter first.
* Move MetricsFlag into metrics_flag.py
Keeps the disable*Metrics logic in a module with minimal imports to
reduce circular import risk. execution.py reads the disabled cell types
from metrics_flag, metric.py no longer imports execution or options,
and DebugOptions is imported inside set_default_pipeline_options.
* Store disabled cell types by name so metrics_flag doesn't import cells
---------
Co-authored-by: Claude Opus 5 (1M context) <[email protected]>
---
CHANGES.md | 1 +
sdks/python/apache_beam/metrics/execution.pxd | 1 +
sdks/python/apache_beam/metrics/execution.py | 11 +
sdks/python/apache_beam/metrics/metrics_flag.py | 100 +++++++++
.../apache_beam/metrics/metrics_flag_test.py | 225 +++++++++++++++++++++
sdks/python/apache_beam/pipeline.py | 2 +
.../apache_beam/runners/worker/sdk_worker_main.py | 2 +
7 files changed, 342 insertions(+)
diff --git a/CHANGES.md b/CHANGES.md
index 408f97fad90..336773b93e5 100644
--- a/CHANGES.md
+++ b/CHANGES.md
@@ -69,6 +69,7 @@
## New Features / Improvements
* (Python) Expanded the SDK worker heap dump
(`--experiments=enable_heap_dump`) with process RSS, CPython allocator/GC
stats, and glibc `mallinfo2` native-heap/fragmentation stats to help
distinguish native-heap from Python-object memory growth
([#39244](https://github.com/apache/beam/issues/39244)).
+* The `disableCounterMetrics`, `disableStringSetMetrics` and
`disableBoundedTrieMetrics` experiments are now honored by the Python SDK, as
they already were in Java (Python)
([#38746](https://github.com/apache/beam/issues/38746)).
## Breaking Changes
diff --git a/sdks/python/apache_beam/metrics/execution.pxd
b/sdks/python/apache_beam/metrics/execution.pxd
index 6311158ea1a..06b54282922 100644
--- a/sdks/python/apache_beam/metrics/execution.pxd
+++ b/sdks/python/apache_beam/metrics/execution.pxd
@@ -31,6 +31,7 @@ cdef class _TypedMetricName(object):
cdef object _DEFAULT
+cdef set _DISABLED_CELL_TYPES
cdef class MetricUpdater(object):
diff --git a/sdks/python/apache_beam/metrics/execution.py
b/sdks/python/apache_beam/metrics/execution.py
index e304658e09a..12403189806 100644
--- a/sdks/python/apache_beam/metrics/execution.py
+++ b/sdks/python/apache_beam/metrics/execution.py
@@ -39,10 +39,12 @@ from typing import Any
from typing import Dict
from typing import FrozenSet
from typing import Optional
+from typing import Set
from typing import Type
from typing import Union
from typing import cast
+from apache_beam.metrics import metrics_flag
from apache_beam.metrics import monitoring_infos
from apache_beam.metrics.cells import BoundedTrieCell
from apache_beam.metrics.cells import CounterCell
@@ -201,6 +203,11 @@ class _TypedMetricName(object):
_DEFAULT = None # type: Any
+# Names of the metric cell types whose updates are dropped process-wide, owned
+# by apache_beam.metrics.metrics_flag.MetricsFlag from the disable*Metrics
+# experiments. Empty (the default) means every update is delivered.
+_DISABLED_CELL_TYPES = metrics_flag.DISABLED_CELL_TYPES # type: Set[str]
+
class MetricUpdater(object):
"""A callable that updates the metric as quickly as possible."""
@@ -216,6 +223,10 @@ class MetricUpdater(object):
def __call__(self, value=_DEFAULT):
# type: (Any) -> None
+ if _DISABLED_CELL_TYPES and (getattr(self.typed_metric_name.cell_type,
+ '__name__',
+ None) in _DISABLED_CELL_TYPES):
+ return
if value is _DEFAULT:
if self.default_value is _DEFAULT:
raise ValueError(
diff --git a/sdks/python/apache_beam/metrics/metrics_flag.py
b/sdks/python/apache_beam/metrics/metrics_flag.py
new file mode 100644
index 00000000000..af00dd2b114
--- /dev/null
+++ b/sdks/python/apache_beam/metrics/metrics_flag.py
@@ -0,0 +1,100 @@
+#
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements. See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
+#
+
+"""Process-wide switches that stop kinds of user metrics from being reported.
+
+This module is imported by ``apache_beam.metrics.execution`` and by the
+pipeline and worker start-up code, so it deliberately keeps its imports to a
+minimum to avoid circular imports.
+"""
+
+# pytype: skip-file
+
+import logging
+from typing import TYPE_CHECKING
+from typing import Set
+
+if TYPE_CHECKING:
+ from apache_beam.options.pipeline_options import PipelineOptions
+
+__all__ = ['MetricsFlag']
+
+_LOGGER = logging.getLogger(__name__)
+
+# Names of the metric cell types whose updates are dropped process-wide. Empty
+# (the default) means every update is delivered. apache_beam.metrics.execution
+# holds a reference to this set, so it is only ever mutated in place.
+DISABLED_CELL_TYPES = set() # type: Set[str]
+
+
+class MetricsFlag(object):
+ """Process-wide switches that stop kinds of user metrics from being reported.
+
+ High throughput jobs may want to turn off metrics that put pressure on the
+ metrics backend. Mirroring the Java SDK, the ``disableCounterMetrics``,
+ ``disableStringSetMetrics`` and ``disableBoundedTrieMetrics`` experiments
make
+ the corresponding ``Metrics.counter``, ``Metrics.string_set`` and
+ ``Metrics.bounded_trie`` updates no-ops. The metric objects themselves are
+ unchanged, so code that holds on to them keeps working.
+ """
+ _EXPERIMENTS = (
+ ('disableCounterMetrics', 'CounterCell', 'Counter'),
+ ('disableStringSetMetrics', 'StringSetCell', 'StringSet'),
+ ('disableBoundedTrieMetrics', 'BoundedTrieCell', 'BoundedTrie'),
+ )
+ _initialized = False
+
+ @classmethod
+ def set_default_pipeline_options(cls, options: 'PipelineOptions') -> None:
+ """Initializes the flags from ``options`` if not already done so.
+
+ Called when a ``Pipeline`` is constructed and at SDK worker harness
+ start-up. As in the Java SDK, the first call wins so that user code running
+ on a worker cannot change the flags the harness was started with.
+ """
+ if cls._initialized:
+ return
+ # Imported here rather than at module level, as a precaution against
+ # circular imports between apache_beam.metrics and apache_beam.options.
+ from apache_beam.options.pipeline_options import DebugOptions
+ debug_options = options.view_as(DebugOptions)
+ disabled = set()
+ for experiment, cell_type_name, kind in cls._EXPERIMENTS:
+ if debug_options.lookup_experiment(experiment):
+ disabled.add(cell_type_name)
+ _LOGGER.info('%s metrics are disabled.', kind)
+ DISABLED_CELL_TYPES.clear()
+ DISABLED_CELL_TYPES.update(disabled)
+ cls._initialized = True
+
+ @classmethod
+ def counter_disabled(cls) -> bool:
+ return 'CounterCell' in DISABLED_CELL_TYPES
+
+ @classmethod
+ def string_set_disabled(cls) -> bool:
+ return 'StringSetCell' in DISABLED_CELL_TYPES
+
+ @classmethod
+ def bounded_trie_disabled(cls) -> bool:
+ return 'BoundedTrieCell' in DISABLED_CELL_TYPES
+
+ @classmethod
+ def reset(cls) -> None:
+ """Clears the flags so the next ``set_default_pipeline_options``
applies."""
+ DISABLED_CELL_TYPES.clear()
+ cls._initialized = False
diff --git a/sdks/python/apache_beam/metrics/metrics_flag_test.py
b/sdks/python/apache_beam/metrics/metrics_flag_test.py
new file mode 100644
index 00000000000..bbee3baaf85
--- /dev/null
+++ b/sdks/python/apache_beam/metrics/metrics_flag_test.py
@@ -0,0 +1,225 @@
+#
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements. See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
+#
+
+# pytype: skip-file
+
+import pickle
+import unittest
+
+import apache_beam as beam
+from apache_beam.metrics.cells import DistributionData
+from apache_beam.metrics.execution import MetricKey
+from apache_beam.metrics.execution import MetricsContainer
+from apache_beam.metrics.execution import MetricsEnvironment
+from apache_beam.metrics.execution import MetricUpdater
+from apache_beam.metrics.metric import Metrics
+from apache_beam.metrics.metric import MetricsFilter
+from apache_beam.metrics.metricbase import MetricName
+from apache_beam.metrics.metrics_flag import MetricsFlag
+from apache_beam.options.pipeline_options import PipelineOptions
+from apache_beam.runners.worker import statesampler
+from apache_beam.testing.test_pipeline import TestPipeline
+from apache_beam.testing.util import assert_that
+from apache_beam.testing.util import equal_to
+from apache_beam.utils import counters
+
+
+class MetricsFlagTest(unittest.TestCase):
+ """Covers the disable*Metrics experiments.
+
+ See https://github.com/apache/beam/issues/38746.
+ """
+ def setUp(self):
+ MetricsFlag.reset()
+ self.sampler = statesampler.StateSampler('', counters.CounterFactory())
+ statesampler.set_current_tracker(self.sampler)
+ self.state = self.sampler.scoped_state(
+ 'mystep', 'myState', metrics_container=MetricsContainer('mystep'))
+ self.sampler.start()
+
+ def tearDown(self):
+ self.sampler.stop()
+ MetricsFlag.reset()
+
+ @staticmethod
+ def _set_experiments(*experiments):
+ MetricsFlag.set_default_pipeline_options(
+ PipelineOptions(['--experiments=%s' % exp for exp in experiments]))
+
+ def test_flags_follow_experiments(self):
+ self.assertFalse(MetricsFlag.counter_disabled())
+ self.assertFalse(MetricsFlag.string_set_disabled())
+ self.assertFalse(MetricsFlag.bounded_trie_disabled())
+
+ for experiment, expected in [
+ ('disableCounterMetrics', (True, False, False)),
+ ('disableStringSetMetrics', (False, True, False)),
+ ('disableBoundedTrieMetrics', (False, False, True)),
+ ]:
+ MetricsFlag.reset()
+ self._set_experiments(experiment)
+ self.assertEqual((
+ MetricsFlag.counter_disabled(),
+ MetricsFlag.string_set_disabled(),
+ MetricsFlag.bounded_trie_disabled()),
+ expected,
+ experiment)
+
+ MetricsFlag.reset()
+ self._set_experiments(
+ 'disableCounterMetrics',
+ 'disableStringSetMetrics',
+ 'disableBoundedTrieMetrics')
+ self.assertTrue(MetricsFlag.counter_disabled())
+ self.assertTrue(MetricsFlag.string_set_disabled())
+ self.assertTrue(MetricsFlag.bounded_trie_disabled())
+
+ def test_first_call_wins(self):
+ self._set_experiments('disableCounterMetrics')
+ # Later options, e.g. from user code constructing a Pipeline on a worker,
+ # do not change the flags the harness was started with.
+ self._set_experiments('disableStringSetMetrics')
+ self.assertTrue(MetricsFlag.counter_disabled())
+ self.assertFalse(MetricsFlag.string_set_disabled())
+
+ def test_update_call_shapes_keep_working(self):
+ # The first attempt at this feature (#38749) was reverted because it
+ # changed the signature of DelegatingCounter.inc; the metric objects must
+ # stay MetricUpdater callables that accept the value as a keyword too.
+ with self.state:
+ counter = Metrics.counter('ns', 'counter')
+ self.assertIsInstance(counter.inc, MetricUpdater)
+ counter.inc()
+ counter.inc(4)
+ counter.inc(value=5)
+ counter.dec()
+ counter.dec(2)
+ string_set = Metrics.string_set('ns', 'set')
+ string_set.add('a')
+ string_set.add(value='b')
+ container = MetricsEnvironment.current_container()
+ self.assertEqual(
+ container.get_counter(MetricName('ns', 'counter')).get_cumulative(),
+ 7)
+ self.assertEqual(
+ container.get_string_set(MetricName(
+ 'ns', 'set')).get_cumulative().string_set, {'a', 'b'})
+
+ def test_disabled_counter_is_noop(self):
+ with self.state:
+ container = MetricsEnvironment.current_container()
+ Metrics.counter('ns', 'before').inc()
+ self.assertEqual(len(container.metrics), 1)
+
+ self._set_experiments('disableCounterMetrics')
+ created_before = Metrics.counter('ns', 'before')
+ created_before.inc()
+ created_before.inc(value=5)
+ created_before.dec()
+ Metrics.counter('ns', 'after').inc(3)
+ self.assertEqual(len(container.metrics), 1)
+ self.assertEqual(
+ container.get_counter(MetricName('ns', 'before')).get_cumulative(),
1)
+
+ # Other kinds keep reporting.
+ Metrics.distribution('ns', 'dist').update(3)
+ Metrics.gauge('ns', 'gauge').set(2)
+ Metrics.string_set('ns', 'set').add('x')
+ Metrics.bounded_trie('ns', 'trie').add(('x', ))
+ self.assertEqual(len(container.metrics), 5)
+
+ def test_disabled_string_set_is_noop(self):
+ with self.state:
+ container = MetricsEnvironment.current_container()
+ Metrics.string_set('ns', 'before').add('seed')
+ self.assertEqual(len(container.metrics), 1)
+
+ self._set_experiments('disableStringSetMetrics')
+ Metrics.string_set('ns', 'before').add('more')
+ Metrics.string_set('ns', 'after').add('value')
+ self.assertEqual(len(container.metrics), 1)
+ self.assertEqual(
+ container.get_string_set(MetricName(
+ 'ns', 'before')).get_cumulative().string_set, {'seed'})
+ Metrics.counter('ns', 'counter').inc()
+ self.assertEqual(len(container.metrics), 2)
+
+ def test_disabled_bounded_trie_is_noop(self):
+ with self.state:
+ container = MetricsEnvironment.current_container()
+ Metrics.bounded_trie('ns', 'before').add(('a', ))
+ self.assertEqual(len(container.metrics), 1)
+
+ self._set_experiments('disableBoundedTrieMetrics')
+ Metrics.bounded_trie('ns', 'before').add(('a', 'b'))
+ Metrics.bounded_trie('ns', 'after').add(('c', ))
+ self.assertEqual(len(container.metrics), 1)
+ self.assertEqual(
+ list(
+ container.get_bounded_trie(MetricName(
+ 'ns', 'before')).get_cumulative().flattened()),
+ [('a', False)])
+
+ def test_disabled_process_wide_counter_is_noop(self):
+ self._set_experiments('disableCounterMetrics')
+ name = MetricName('ns', 'process_wide')
+ counter = Metrics.DelegatingCounter(name, process_wide=True)
+ counter.inc()
+ # The process-wide container is shared with every other test in this
+ # process, so only check that this counter never reached it.
+ self.assertNotIn(
+ MetricKey(None, name),
+ MetricsEnvironment.process_wide_container().get_cumulative().counters)
+
+ def test_disabled_flag_applies_to_unpickled_metrics(self):
+ # DoFns holding metric objects are pickled at submission time and
+ # unpickled on the worker, where the harness sets the flags.
+ counter = Metrics.counter('ns', 'pickled')
+ counter = pickle.loads(pickle.dumps(counter))
+ self.assertIsInstance(counter.inc, MetricUpdater)
+ self._set_experiments('disableCounterMetrics')
+ with self.state:
+ counter.inc()
+ self.assertEqual(len(MetricsEnvironment.current_container().metrics), 0)
+
+ def test_disabled_counters_in_pipeline(self):
+ class SomeDoFn(beam.DoFn):
+ def process(self, element):
+ Metrics.counter(self.__class__, 'elements').inc()
+ Metrics.distribution(self.__class__, 'element_dist').update(element)
+ yield element
+
+ MetricsFlag.reset()
+ pipeline = TestPipeline(
+ options=PipelineOptions(['--experiments=disableCounterMetrics']))
+ results = pipeline | beam.Create([1, 2, 3]) | beam.ParDo(SomeDoFn())
+ assert_that(results, equal_to([1, 2, 3]))
+ res = pipeline.run()
+ res.wait_until_finish()
+
+ self.assertEqual(
+ res.metrics().query(MetricsFilter().with_name('elements'))['counters'],
+ [])
+ distributions = res.metrics().query(
+ MetricsFilter().with_name('element_dist'))['distributions']
+ self.assertEqual(len(distributions), 1)
+ self.assertEqual(
+ distributions[0].committed.data, DistributionData(6, 3, 1, 3))
+
+
+if __name__ == '__main__':
+ unittest.main()
diff --git a/sdks/python/apache_beam/pipeline.py
b/sdks/python/apache_beam/pipeline.py
index 750868f7443..fea42bbc0cb 100644
--- a/sdks/python/apache_beam/pipeline.py
+++ b/sdks/python/apache_beam/pipeline.py
@@ -73,6 +73,7 @@ from apache_beam import pvalue
from apache_beam.coders import typecoders
from apache_beam.internal import pickler
from apache_beam.io.filesystems import FileSystems
+from apache_beam.metrics.metrics_flag import MetricsFlag
from apache_beam.options.pipeline_options import CrossLanguageOptions
from apache_beam.options.pipeline_options import DebugOptions
from apache_beam.options.pipeline_options import PipelineOptions
@@ -192,6 +193,7 @@ class Pipeline(HasDisplayData):
self._options = PipelineOptions([])
FileSystems.set_options(self._options)
+ MetricsFlag.set_default_pipeline_options(self._options)
if runner is None:
runner = self._options.view_as(StandardOptions).runner
diff --git a/sdks/python/apache_beam/runners/worker/sdk_worker_main.py
b/sdks/python/apache_beam/runners/worker/sdk_worker_main.py
index 8bee86f010f..cf18030deec 100644
--- a/sdks/python/apache_beam/runners/worker/sdk_worker_main.py
+++ b/sdks/python/apache_beam/runners/worker/sdk_worker_main.py
@@ -33,6 +33,7 @@ from google.protobuf import text_format
from apache_beam.internal import pickler
from apache_beam.io import filesystems
+from apache_beam.metrics.metrics_flag import MetricsFlag
from apache_beam.options.pipeline_options import DebugOptions
from apache_beam.options.pipeline_options import GoogleCloudOptions
from apache_beam.options.pipeline_options import PipelineOptions
@@ -129,6 +130,7 @@ def create_harness(environment, dry_run=False):
RuntimeValueProvider.set_runtime_options(pipeline_options_dict)
sdk_pipeline_options = PipelineOptions.from_dictionary(pipeline_options_dict)
filesystems.FileSystems.set_options(sdk_pipeline_options)
+ MetricsFlag.set_default_pipeline_options(sdk_pipeline_options)
pickle_library = sdk_pipeline_options.view_as(SetupOptions).pickle_library
pickler.set_library(pickle_library)