This is an automated email from the ASF dual-hosted git repository. yhu pushed a commit to branch revert-33691-delta-sset in repository https://gitbox.apache.org/repos/asf/beam.git
commit cd36df5f044330449d0945fc71ad090b88781d9b Author: Yi Hu <[email protected]> AuthorDate: Fri Jan 31 10:53:47 2025 -0500 Revert "Report delta StringSet counters for streaming (#33691)" This reverts commit bbe63947aa6d7c865a6664526f8d1e8117cca619. --- .../apache/beam/runners/core/metrics/StringSetCell.java | 8 -------- .../dataflow/worker/BatchModeExecutionContext.java | 2 +- .../worker/MetricsToCounterUpdateConverter.java | 5 ++--- .../dataflow/worker/StreamingStepMetricsContainer.java | 17 +++++------------ .../worker/StreamingStepMetricsContainerTest.java | 8 ++++---- 5 files changed, 12 insertions(+), 28 deletions(-) diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/StringSetCell.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/StringSetCell.java index 6bf34306e27..fc8dcb49894 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/StringSetCell.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/StringSetCell.java @@ -72,14 +72,6 @@ public class StringSetCell implements StringSet, MetricCell<StringSetData> { return setValue.get(); } - // Used by Streaming metric container to extract deltas since streaming metrics are - // reported as deltas rather than cumulative as in batch. - // For delta we take the current value then reset the cell to empty so the next call only see - // delta/updates from last call. - public StringSetData getAndReset() { - return setValue.getAndUpdate(unused -> StringSetData.empty()); - } - @Override public MetricName getName() { return name; diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/BatchModeExecutionContext.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/BatchModeExecutionContext.java index cf07e5f7240..aeef7784c2c 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/BatchModeExecutionContext.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/BatchModeExecutionContext.java @@ -519,7 +519,7 @@ public class BatchModeExecutionContext .transform( update -> MetricsToCounterUpdateConverter.fromStringSet( - update.getKey(), true, update.getUpdate()))); + update.getKey(), update.getUpdate()))); }); } diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/MetricsToCounterUpdateConverter.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/MetricsToCounterUpdateConverter.java index 4866d201122..dbedc51528a 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/MetricsToCounterUpdateConverter.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/MetricsToCounterUpdateConverter.java @@ -98,8 +98,7 @@ public class MetricsToCounterUpdateConverter { .setIntegerGauge(integerGaugeProto); } - public static CounterUpdate fromStringSet( - MetricKey key, boolean isCumulative, StringSetData stringSetData) { + public static CounterUpdate fromStringSet(MetricKey key, StringSetData stringSetData) { CounterStructuredNameAndMetadata name = structuredNameAndMetadata(key, Kind.SET); StringList stringList = new StringList(); @@ -107,7 +106,7 @@ public class MetricsToCounterUpdateConverter { return new CounterUpdate() .setStructuredNameAndMetadata(name) - .setCumulative(isCumulative) + .setCumulative(false) .setStringList(stringList); } diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingStepMetricsContainer.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingStepMetricsContainer.java index dd7db5ae9d1..03ba501941a 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingStepMetricsContainer.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingStepMetricsContainer.java @@ -34,7 +34,6 @@ import org.apache.beam.runners.core.metrics.DistributionData; import org.apache.beam.runners.core.metrics.GaugeCell; import org.apache.beam.runners.core.metrics.MetricsMap; import org.apache.beam.runners.core.metrics.StringSetCell; -import org.apache.beam.runners.core.metrics.StringSetData; import org.apache.beam.sdk.metrics.Counter; import org.apache.beam.sdk.metrics.Distribution; import org.apache.beam.sdk.metrics.Gauge; @@ -74,7 +73,7 @@ public class StreamingStepMetricsContainer implements MetricsContainer { private final ConcurrentHashMap<MetricName, GaugeCell> perWorkerGauges = new ConcurrentHashMap<>(); - private MetricsMap<MetricName, StringSetCell> stringSets = new MetricsMap<>(StringSetCell::new); + private MetricsMap<MetricName, StringSetCell> stringSet = new MetricsMap<>(StringSetCell::new); private MetricsMap<MetricName, DeltaDistributionCell> distributions = new MetricsMap<>(DeltaDistributionCell::new); @@ -183,7 +182,7 @@ public class StreamingStepMetricsContainer implements MetricsContainer { @Override public StringSet getStringSet(MetricName metricName) { - return stringSets.get(metricName); + return stringSet.get(metricName); } @Override @@ -203,11 +202,9 @@ public class StreamingStepMetricsContainer implements MetricsContainer { } public Iterable<CounterUpdate> extractUpdates() { - // Streaming metrics are updated as delta and not cumulative. return counterUpdates() .append(distributionUpdates()) - .append(gaugeUpdates()) - .append(stringSetUpdates()); + .append(gaugeUpdates().append(stringSetUpdates())); } private FluentIterable<CounterUpdate> counterUpdates() { @@ -250,18 +247,14 @@ public class StreamingStepMetricsContainer implements MetricsContainer { } private FluentIterable<CounterUpdate> stringSetUpdates() { - return FluentIterable.from(stringSets.entries()) + return FluentIterable.from(stringSet.entries()) .transform( new Function<Entry<MetricName, StringSetCell>, CounterUpdate>() { @Override public @Nullable CounterUpdate apply( @Nonnull Map.Entry<MetricName, StringSetCell> entry) { - StringSetData value = entry.getValue().getAndReset(); - if (value.stringSet().isEmpty()) { - return null; - } return MetricsToCounterUpdateConverter.fromStringSet( - MetricKey.create(stepName, entry.getKey()), false, value); + MetricKey.create(stepName, entry.getKey()), entry.getValue().getCumulative()); } }) .filter(Predicates.notNull()); diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingStepMetricsContainerTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingStepMetricsContainerTest.java index 4b758aa6cd4..37c5ad26128 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingStepMetricsContainerTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingStepMetricsContainerTest.java @@ -315,14 +315,14 @@ public class StreamingStepMetricsContainerTest { .setStringList(new StringList().setElements(Arrays.asList("ij", "kl", "mn"))); updates = StreamingStepMetricsContainer.extractMetricUpdates(registry); - assertThat(updates, containsInAnyOrder(name2Update)); + assertThat(updates, containsInAnyOrder(name1Update, name2Update)); - // test deltas c1.getStringSet(name1).add("op"); - name1Update.setStringList(new StringList().setElements(Arrays.asList("op"))); + name1Update.setStringList( + new StringList().setElements(Arrays.asList("ab", "cd", "ef", "gh", "op"))); updates = StreamingStepMetricsContainer.extractMetricUpdates(registry); - assertThat(updates, containsInAnyOrder(name1Update)); + assertThat(updates, containsInAnyOrder(name1Update, name2Update)); } @Test
