Skip to content

Commit 02168d3

Browse files
Revert "[Python] Honor disableCounterMetrics, disableStringSetMetrics, and disableBoundedTrieMetrics experiments" (#38901)
* Revert "[Python] Honor disableCounterMetrics, disableStringSetMetrics, and di…" This reverts commit c01ceea. * Apply suggestions from code review Co-authored-by: gemini-code-assist[bot] <176961590+gemini-code-assist[bot]@users.noreply.github.com> --------- Co-authored-by: gemini-code-assist[bot] <176961590+gemini-code-assist[bot]@users.noreply.github.com>
1 parent 6a2dbb6 commit 02168d3

5 files changed

Lines changed: 3 additions & 160 deletions

File tree

CHANGES.md

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -69,7 +69,6 @@
6969

7070
## New Features / Improvements
7171

72-
* X feature added (Java/Python) ([#X](https://github.com/apache/beam/issues/X)).
7372
* (Java) Enabled state tag encoding v2 by default for new Dataflow Streaming Engine jobs. It can be disabled by passing `--experiments=disable_streaming_engine_state_tag_encoding_v2` or `--updateCompatibilityVersion=2.74.0` pipeline option. Note that the tag encoding version cannot change during a job update. Jobs using tag encoding v2 (enabled by default for new jobs on 2.75.0+) cannot be downgraded to Beam versions prior to 2.73.0, as only versions 2.73.0 and later support tag encoding v2. ([#38705](https://github.com/apache/beam/issues/38705)).
7473

7574
## Breaking Changes

sdks/python/apache_beam/metrics/metric.py

Lines changed: 3 additions & 51 deletions
Original file line numberDiff line numberDiff line change
@@ -46,51 +46,18 @@
4646
from apache_beam.metrics.metricbase import Histogram
4747
from apache_beam.metrics.metricbase import MetricName
4848
from apache_beam.metrics.metricbase import StringSet
49-
from apache_beam.options.pipeline_options import DebugOptions
5049

5150
if TYPE_CHECKING:
5251
from apache_beam.internal.metrics.metric import MetricLogger
5352
from apache_beam.metrics.execution import MetricKey
5453
from apache_beam.metrics.metricbase import Metric
55-
from apache_beam.options.pipeline_options import PipelineOptions
5654
from apache_beam.utils.histogram import BucketType
5755

5856
__all__ = ['Metrics', 'MetricsFilter', 'Lineage']
5957

6058
_LOGGER = logging.getLogger(__name__)
6159

6260

63-
class MetricsFlag(object):
64-
"""Process-wide flags controlling which user metric kinds are emitted."""
65-
counter_disabled = False
66-
string_set_disabled = False
67-
bounded_trie_disabled = False
68-
_initialized = False
69-
70-
@classmethod
71-
def set_default_pipeline_options(cls, options: 'PipelineOptions') -> None:
72-
if cls._initialized:
73-
return
74-
debug_options = options.view_as(DebugOptions)
75-
if debug_options.lookup_experiment('disableCounterMetrics'):
76-
cls.counter_disabled = True
77-
_LOGGER.info('Counter metrics are disabled.')
78-
if debug_options.lookup_experiment('disableStringSetMetrics'):
79-
cls.string_set_disabled = True
80-
_LOGGER.info('StringSet metrics are disabled.')
81-
if debug_options.lookup_experiment('disableBoundedTrieMetrics'):
82-
cls.bounded_trie_disabled = True
83-
_LOGGER.info('BoundedTrie metrics are disabled.')
84-
cls._initialized = True
85-
86-
@classmethod
87-
def reset(cls) -> None:
88-
cls.counter_disabled = False
89-
cls.string_set_disabled = False
90-
cls.bounded_trie_disabled = False
91-
cls._initialized = False
92-
93-
9461
class Metrics(object):
9562
"""Lets users create/access metric objects during pipeline execution."""
9663
@staticmethod
@@ -237,17 +204,12 @@ class DelegatingCounter(Counter):
237204
def __init__(
238205
self, metric_name: MetricName, process_wide: bool = False) -> None:
239206
super().__init__(metric_name)
240-
self._updater = MetricUpdater(
207+
self.inc = MetricUpdater( # type: ignore[method-assign]
241208
cells.CounterCell,
242209
metric_name,
243210
default_value=1,
244211
process_wide=process_wide)
245212

246-
def inc(self, n: int = 1) -> None:
247-
if MetricsFlag.counter_disabled:
248-
return
249-
self._updater(n)
250-
251213
class DelegatingDistribution(Distribution):
252214
"""Metrics Distribution Delegates functionality to MetricsEnvironment."""
253215
def __init__(
@@ -269,23 +231,13 @@ class DelegatingStringSet(StringSet):
269231
"""Metrics StringSet that Delegates functionality to MetricsEnvironment."""
270232
def __init__(self, metric_name: MetricName) -> None:
271233
super().__init__(metric_name)
272-
self._updater = MetricUpdater(cells.StringSetCell, metric_name)
273-
274-
def add(self, value: str) -> None:
275-
if MetricsFlag.string_set_disabled:
276-
return
277-
self._updater(value)
234+
self.add = MetricUpdater(cells.StringSetCell, metric_name) # type: ignore[method-assign]
278235

279236
class DelegatingBoundedTrie(BoundedTrie):
280237
"""Metrics BoundedTrie that Delegates functionality to MetricsEnvironment."""
281238
def __init__(self, metric_name: MetricName) -> None:
282239
super().__init__(metric_name)
283-
self._updater = MetricUpdater(cells.BoundedTrieCell, metric_name)
284-
285-
def add(self, value) -> None:
286-
if MetricsFlag.bounded_trie_disabled:
287-
return
288-
self._updater(value)
240+
self.add = MetricUpdater(cells.BoundedTrieCell, metric_name) # type: ignore[method-assign]
289241

290242

291243
class MetricResults(object):

sdks/python/apache_beam/metrics/metric_test.py

Lines changed: 0 additions & 104 deletions
Original file line numberDiff line numberDiff line change
@@ -32,9 +32,7 @@
3232
from apache_beam.metrics.metric import MetricResults
3333
from apache_beam.metrics.metric import Metrics
3434
from apache_beam.metrics.metric import MetricsFilter
35-
from apache_beam.metrics.metric import MetricsFlag
3635
from apache_beam.metrics.metricbase import MetricName
37-
from apache_beam.options.pipeline_options import PipelineOptions
3836
from apache_beam.runners.direct.direct_runner import BundleBasedDirectRunner
3937
from apache_beam.runners.worker import statesampler
4038
from apache_beam.testing.metric_result_matchers import DistributionMatcher
@@ -123,108 +121,6 @@ def test_get_namespace_error(self):
123121
with self.assertRaises(ValueError):
124122
Metrics.get_namespace(object())
125123

126-
def test_metrics_flag(self):
127-
MetricsFlag.reset()
128-
try:
129-
self.assertFalse(MetricsFlag.counter_disabled)
130-
self.assertFalse(MetricsFlag.string_set_disabled)
131-
self.assertFalse(MetricsFlag.bounded_trie_disabled)
132-
133-
options = PipelineOptions(['--experiments=disableCounterMetrics'])
134-
MetricsFlag.set_default_pipeline_options(options)
135-
self.assertTrue(MetricsFlag.counter_disabled)
136-
self.assertFalse(MetricsFlag.string_set_disabled)
137-
self.assertFalse(MetricsFlag.bounded_trie_disabled)
138-
139-
MetricsFlag.reset()
140-
options = PipelineOptions(['--experiments=disableStringSetMetrics'])
141-
MetricsFlag.set_default_pipeline_options(options)
142-
self.assertFalse(MetricsFlag.counter_disabled)
143-
self.assertTrue(MetricsFlag.string_set_disabled)
144-
self.assertFalse(MetricsFlag.bounded_trie_disabled)
145-
146-
MetricsFlag.reset()
147-
options = PipelineOptions(['--experiments=disableBoundedTrieMetrics'])
148-
MetricsFlag.set_default_pipeline_options(options)
149-
self.assertFalse(MetricsFlag.counter_disabled)
150-
self.assertFalse(MetricsFlag.string_set_disabled)
151-
self.assertTrue(MetricsFlag.bounded_trie_disabled)
152-
153-
MetricsFlag.reset()
154-
options = PipelineOptions([
155-
'--experiments=disableCounterMetrics',
156-
'--experiments=disableStringSetMetrics',
157-
'--experiments=disableBoundedTrieMetrics',
158-
])
159-
MetricsFlag.set_default_pipeline_options(options)
160-
self.assertTrue(MetricsFlag.counter_disabled)
161-
self.assertTrue(MetricsFlag.string_set_disabled)
162-
self.assertTrue(MetricsFlag.bounded_trie_disabled)
163-
finally:
164-
MetricsFlag.reset()
165-
166-
def test_disabled_counter_is_noop(self):
167-
sampler = statesampler.StateSampler('', counters.CounterFactory())
168-
statesampler.set_current_tracker(sampler)
169-
state = sampler.scoped_state(
170-
'mystep', 'myState', metrics_container=MetricsContainer('mystep'))
171-
MetricsFlag.reset()
172-
try:
173-
sampler.start()
174-
with state:
175-
container = MetricsEnvironment.current_container()
176-
Metrics.counter('ns', 'baseline').inc()
177-
self.assertEqual(len(container.metrics), 1)
178-
options = PipelineOptions(['--experiments=disableCounterMetrics'])
179-
MetricsFlag.set_default_pipeline_options(options)
180-
Metrics.counter('ns', 'after_disable').inc()
181-
Metrics.counter('ns', 'after_disable').inc(5)
182-
Metrics.counter('ns', 'after_disable').dec()
183-
self.assertEqual(len(container.metrics), 1)
184-
finally:
185-
sampler.stop()
186-
MetricsFlag.reset()
187-
188-
def test_disabled_string_set_is_noop(self):
189-
sampler = statesampler.StateSampler('', counters.CounterFactory())
190-
statesampler.set_current_tracker(sampler)
191-
state = sampler.scoped_state(
192-
'mystep', 'myState', metrics_container=MetricsContainer('mystep'))
193-
MetricsFlag.reset()
194-
try:
195-
sampler.start()
196-
with state:
197-
container = MetricsEnvironment.current_container()
198-
Metrics.string_set('ns', 'baseline').add('seed')
199-
self.assertEqual(len(container.metrics), 1)
200-
options = PipelineOptions(['--experiments=disableStringSetMetrics'])
201-
MetricsFlag.set_default_pipeline_options(options)
202-
Metrics.string_set('ns', 'after_disable').add('value')
203-
self.assertEqual(len(container.metrics), 1)
204-
finally:
205-
sampler.stop()
206-
MetricsFlag.reset()
207-
208-
def test_disabled_bounded_trie_is_noop(self):
209-
sampler = statesampler.StateSampler('', counters.CounterFactory())
210-
statesampler.set_current_tracker(sampler)
211-
state = sampler.scoped_state(
212-
'mystep', 'myState', metrics_container=MetricsContainer('mystep'))
213-
MetricsFlag.reset()
214-
try:
215-
sampler.start()
216-
with state:
217-
container = MetricsEnvironment.current_container()
218-
Metrics.bounded_trie('ns', 'baseline').add(['a'])
219-
self.assertEqual(len(container.metrics), 1)
220-
options = PipelineOptions(['--experiments=disableBoundedTrieMetrics'])
221-
MetricsFlag.set_default_pipeline_options(options)
222-
Metrics.bounded_trie('ns', 'after_disable').add(['a', 'b'])
223-
self.assertEqual(len(container.metrics), 1)
224-
finally:
225-
sampler.stop()
226-
MetricsFlag.reset()
227-
228124
def test_counter_empty_name(self):
229125
with self.assertRaises(ValueError):
230126
Metrics.counter("namespace", "")

sdks/python/apache_beam/pipeline.py

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -73,7 +73,6 @@
7373
from apache_beam.coders import typecoders
7474
from apache_beam.internal import pickler
7575
from apache_beam.io.filesystems import FileSystems
76-
from apache_beam.metrics.metric import MetricsFlag
7776
from apache_beam.options.pipeline_options import CrossLanguageOptions
7877
from apache_beam.options.pipeline_options import DebugOptions
7978
from apache_beam.options.pipeline_options import PipelineOptions
@@ -193,7 +192,6 @@ def __init__(
193192
self._options = PipelineOptions([])
194193

195194
FileSystems.set_options(self._options)
196-
MetricsFlag.set_default_pipeline_options(self._options)
197195

198196
if runner is None:
199197
runner = self._options.view_as(StandardOptions).runner

sdks/python/apache_beam/runners/worker/sdk_worker_main.py

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -33,7 +33,6 @@
3333

3434
from apache_beam.internal import pickler
3535
from apache_beam.io import filesystems
36-
from apache_beam.metrics import metric
3736
from apache_beam.options.pipeline_options import DebugOptions
3837
from apache_beam.options.pipeline_options import GoogleCloudOptions
3938
from apache_beam.options.pipeline_options import PipelineOptions
@@ -124,7 +123,6 @@ def create_harness(environment, dry_run=False):
124123
RuntimeValueProvider.set_runtime_options(pipeline_options_dict)
125124
sdk_pipeline_options = PipelineOptions.from_dictionary(pipeline_options_dict)
126125
filesystems.FileSystems.set_options(sdk_pipeline_options)
127-
metric.MetricsFlag.set_default_pipeline_options(sdk_pipeline_options)
128126
pickle_library = sdk_pipeline_options.view_as(SetupOptions).pickle_library
129127
pickler.set_library(pickle_library)
130128

0 commit comments

Comments
 (0)