-
Notifications
You must be signed in to change notification settings - Fork 4.7k
[Python] Honor disableCounterMetrics, disableStringSetMetrics and disableBoundedTrieMetrics experiments #40165
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: master
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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 | ||
|
|
@@ -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): | ||
| """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) | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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) | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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 | ||
|
|
||
There was a problem hiding this comment.
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?