Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions CHANGES.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
1 change: 1 addition & 0 deletions sdks/python/apache_beam/metrics/execution.pxd
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@ cdef class _TypedMetricName(object):


cdef object _DEFAULT
cdef set _DISABLED_CELL_TYPES


cdef class MetricUpdater(object):
Expand Down
23 changes: 23 additions & 0 deletions sdks/python/apache_beam/metrics/execution.py
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,9 @@
from typing import Any
from typing import Dict
from typing import FrozenSet
from typing import Iterable
from typing import Optional
from typing import Set
from typing import Type
from typing import Union
from typing import cast
Expand Down Expand Up @@ -201,6 +203,24 @@ def __reduce__(self):

_DEFAULT = None # type: Any

# Metric cell types whose updates are dropped process-wide. Populated from the
# disable*Metrics experiments by apache_beam.metrics.metric.MetricsFlag; empty
# (the default) means every update is delivered.
_DISABLED_CELL_TYPES = set() # type: Set[Any]


def set_disabled_cell_types(cell_types):
# type: (Iterable[Any]) -> None

"""Replaces the set of metric cell types whose updates are dropped."""
_DISABLED_CELL_TYPES.clear()
_DISABLED_CELL_TYPES.update(cell_types)


def is_cell_type_disabled(cell_type):
# type: (Any) -> bool
return cell_type in _DISABLED_CELL_TYPES


class MetricUpdater(object):
"""A callable that updates the metric as quickly as possible."""
Expand All @@ -216,6 +236,9 @@ def __init__(

def __call__(self, value=_DEFAULT):
# type: (Any) -> None
if _DISABLED_CELL_TYPES and (self.typed_metric_name.cell_type
in _DISABLED_CELL_TYPES):
return
if value is _DEFAULT:
if self.default_value is _DEFAULT:
raise ValueError(
Expand Down
59 changes: 59 additions & 0 deletions sdks/python/apache_beam/metrics/metric.py
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@
from typing import Union

from apache_beam.metrics import cells
from apache_beam.metrics import execution
from apache_beam.metrics.cells import HistogramCellFactory
from apache_beam.metrics.execution import MetricResult
from apache_beam.metrics.execution import MetricUpdater
Expand All @@ -46,18 +47,76 @@
from apache_beam.metrics.metricbase import Histogram
from apache_beam.metrics.metricbase import MetricName
from apache_beam.metrics.metricbase import StringSet
from apache_beam.options.pipeline_options import DebugOptions

if TYPE_CHECKING:
from apache_beam.internal.metrics.metric import MetricLogger
from apache_beam.metrics.execution import MetricKey
from apache_beam.metrics.metricbase import Metric
from apache_beam.options.pipeline_options import PipelineOptions
from apache_beam.utils.histogram import BucketType

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

_LOGGER = logging.getLogger(__name__)


class MetricsFlag(object):

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Besides the method signature change, I remember another breakage introduced in the first attempt was a circular import. For precaution can we move the new logic into a new metrics_flag.py with minimum required import, to avoid potential issues in callers?

"""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', cells.CounterCell, 'Counter'),
('disableStringSetMetrics', cells.StringSetCell, 'StringSet'),
('disableBoundedTrieMetrics', cells.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
debug_options = options.view_as(DebugOptions)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

consider import debug options in place. Currently no non-test file in apache_beam/metrics/ imported apache_beam.options at runtime, as a pre-caution for circular import risk

disabled = set()
for experiment, cell_type, kind in cls._EXPERIMENTS:
if debug_options.lookup_experiment(experiment):
disabled.add(cell_type)
_LOGGER.info('%s metrics are disabled.', kind)
execution.set_disabled_cell_types(disabled)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

if we put metrics_flag in a separate module, we can query disabled cells from execution.py and no need to import execution here

cls._initialized = True

@classmethod
def counter_disabled(cls) -> bool:
return execution.is_cell_type_disabled(cells.CounterCell)

@classmethod
def string_set_disabled(cls) -> bool:
return execution.is_cell_type_disabled(cells.StringSetCell)

@classmethod
def bounded_trie_disabled(cls) -> bool:
return execution.is_cell_type_disabled(cells.BoundedTrieCell)

@classmethod
def reset(cls) -> None:
"""Clears the flags so the next ``set_default_pipeline_options`` applies."""
execution.set_disabled_cell_types(())
cls._initialized = False


class Metrics(object):
"""Lets users create/access metric objects during pipeline execution."""
@staticmethod
Expand Down
187 changes: 187 additions & 0 deletions sdks/python/apache_beam/metrics/metric_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@
#

# pytype: skip-file
import pickle
import re
import unittest

Expand All @@ -28,11 +29,14 @@
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 Lineage
from apache_beam.metrics.metric import MetricResults
from apache_beam.metrics.metric import Metrics
from apache_beam.metrics.metric import MetricsFilter
from apache_beam.metrics.metric import MetricsFlag
from apache_beam.metrics.metricbase import MetricName
from apache_beam.options.pipeline_options import PipelineOptions
from apache_beam.runners.direct.direct_runner import BundleBasedDirectRunner
from apache_beam.runners.worker import statesampler
from apache_beam.testing.metric_result_matchers import DistributionMatcher
Expand Down Expand Up @@ -250,6 +254,189 @@ def test_create_counter_distribution(self):
sampler.stop()


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


class LineageTest(unittest.TestCase):
def test_fq_name(self):
test_cases = {
Expand Down
2 changes: 2 additions & 0 deletions sdks/python/apache_beam/pipeline.py
Original file line number Diff line number Diff line change
Expand Up @@ -73,6 +73,7 @@
from apache_beam.coders import typecoders
from apache_beam.internal import pickler
from apache_beam.io.filesystems import FileSystems
from apache_beam.metrics.metric 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
Expand Down Expand Up @@ -192,6 +193,7 @@ def __init__(
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
Expand Down
2 changes: 2 additions & 0 deletions sdks/python/apache_beam/runners/worker/sdk_worker_main.py
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@

from apache_beam.internal import pickler
from apache_beam.io import filesystems
from apache_beam.metrics.metric 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
Expand Down Expand Up @@ -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)

Expand Down
Loading