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

Reply via email to