This is an automated email from the ASF dual-hosted git repository.

yhu pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git


The following commit(s) were added to refs/heads/master by this push:
     new f64501427bf Extract ThrottleTimeCounter for PerWorkerMetrics (#30245)
f64501427bf is described below

commit f64501427bf2e0e5482859e66a3d8e1f7d6e6c94
Author: JayajP <[email protected]>
AuthorDate: Wed Feb 14 12:46:21 2024 -0800

    Extract ThrottleTimeCounter for PerWorkerMetrics (#30245)
    
    * Extract ThrottleTimeCounter for PerWorkerMetrics
    
    * spotless
    
    * Update unit tests
---
 .../dataflow/worker/streaming/StageInfo.java       | 20 +++++++
 .../dataflow/worker/streaming/StageInfoTest.java   | 70 ++++++++++++++++++++++
 .../sdk/io/gcp/bigquery/BigQuerySinkMetrics.java   |  4 +-
 3 files changed, 92 insertions(+), 2 deletions(-)

diff --git 
a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/StageInfo.java
 
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/StageInfo.java
index cb6cbec7d4b..8f14ea26a46 100644
--- 
a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/StageInfo.java
+++ 
b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/StageInfo.java
@@ -21,6 +21,7 @@ import static 
org.apache.beam.runners.dataflow.worker.DataflowSystemMetrics.THRO
 
 import com.google.api.services.dataflow.model.CounterStructuredName;
 import com.google.api.services.dataflow.model.CounterUpdate;
+import com.google.api.services.dataflow.model.MetricValue;
 import com.google.api.services.dataflow.model.PerStepNamespaceMetrics;
 import com.google.auto.value.AutoValue;
 import java.util.ArrayList;
@@ -34,6 +35,7 @@ import 
org.apache.beam.runners.dataflow.worker.counters.Counter;
 import org.apache.beam.runners.dataflow.worker.counters.CounterSet;
 import 
org.apache.beam.runners.dataflow.worker.counters.DataflowCounterUpdateExtractor;
 import org.apache.beam.runners.dataflow.worker.counters.NameContext;
+import org.apache.beam.sdk.io.gcp.bigquery.BigQuerySinkMetrics;
 import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Iterables;
 
 /** Contains a few of the stage specific fields. E.g. metrics container 
registry, counters etc. */
@@ -117,6 +119,24 @@ public abstract class StageInfo {
     Iterables.addAll(
         metrics,
         
StreamingStepMetricsContainer.extractPerWorkerMetricUpdates(metricsContainerRegistry()));
+    translateKnownPerWorkerCounters(metrics);
     return metrics;
   }
+
+  private void translateKnownPerWorkerCounters(List<PerStepNamespaceMetrics> 
metrics) {
+    for (PerStepNamespaceMetrics perStepnamespaceMetrics : metrics) {
+      if (!BigQuerySinkMetrics.METRICS_NAMESPACE.equals(
+          perStepnamespaceMetrics.getMetricsNamespace())) {
+        continue;
+      }
+      for (MetricValue metric : perStepnamespaceMetrics.getMetricValues()) {
+        if (BigQuerySinkMetrics.THROTTLED_TIME.equals(metric.getMetric())) {
+          Long msecs = metric.getValueInt64();
+          if (msecs != null && msecs > 0) {
+            throttledMsecs().addValue(msecs);
+          }
+        }
+      }
+    }
+  }
 }
diff --git 
a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/StageInfoTest.java
 
b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/StageInfoTest.java
new file mode 100644
index 00000000000..9b46466c5bf
--- /dev/null
+++ 
b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/StageInfoTest.java
@@ -0,0 +1,70 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.beam.runners.dataflow.worker.streaming;
+
+import static org.hamcrest.MatcherAssert.assertThat;
+import static org.hamcrest.Matchers.equalTo;
+
+import org.apache.beam.runners.dataflow.worker.StreamingStepMetricsContainer;
+import org.apache.beam.sdk.io.gcp.bigquery.BigQuerySinkMetrics;
+import org.apache.beam.sdk.metrics.MetricName;
+import org.apache.beam.sdk.metrics.MetricsEnvironment;
+import org.apache.beam.sdk.util.HistogramData;
+import org.junit.Test;
+import org.junit.runner.RunWith;
+import org.junit.runners.JUnit4;
+
+@RunWith(JUnit4.class)
+public class StageInfoTest {
+  @Test
+  public void testTranslateKnownPerWorkerCounters_validCounters() throws 
Exception {
+    StageInfo stageInfo = StageInfo.create("user_name", "system_name");
+    StreamingStepMetricsContainer metricsContainer =
+        stageInfo.metricsContainerRegistry().getContainer("s1");
+    MetricsEnvironment.setCurrentContainer(metricsContainer);
+    StreamingStepMetricsContainer.setEnablePerWorkerMetrics(true);
+
+    
BigQuerySinkMetrics.throttledTimeCounter(BigQuerySinkMetrics.RpcMethod.APPEND_ROWS).inc(100);
+    stageInfo.extractPerWorkerMetricValues();
+
+    assertThat(stageInfo.throttledMsecs().getAggregate(), equalTo(100L));
+  }
+
+  /**
+   * Test the scenario where there is a metric with a known metric name - 
{@code ThrottledTime} -
+   * that is not a counter. {@code translateKnownPerWorkerCounters} should 
handle this gracefully.
+   */
+  @Test
+  public void testTranslateKnownPerWorkerCounters_malformedCounters() throws 
Exception {
+    StageInfo stageInfo = StageInfo.create("user_name", "system_name");
+    StreamingStepMetricsContainer metricsContainer =
+        stageInfo.metricsContainerRegistry().getContainer("s1");
+    MetricsEnvironment.setCurrentContainer(metricsContainer);
+    StreamingStepMetricsContainer.setEnablePerWorkerMetrics(true);
+
+    MetricName name =
+        
BigQuerySinkMetrics.throttledTimeCounter(BigQuerySinkMetrics.RpcMethod.APPEND_ROWS)
+            .getName();
+
+    HistogramData.BucketType linearBuckets = HistogramData.LinearBuckets.of(0, 
10.0, 10);
+    metricsContainer.getPerWorkerHistogram(name, linearBuckets).update(10.0);
+
+    stageInfo.extractPerWorkerMetricValues();
+    assertThat(stageInfo.throttledMsecs().getAggregate(), equalTo(0L));
+  }
+}
diff --git 
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/BigQuerySinkMetrics.java
 
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/BigQuerySinkMetrics.java
index 34e3b704b4f..0375cf9ab33 100644
--- 
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/BigQuerySinkMetrics.java
+++ 
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/BigQuerySinkMetrics.java
@@ -57,10 +57,10 @@ public class BigQuerySinkMetrics {
   private static final String RPC_REQUESTS = "RpcRequestsCount";
   private static final String RPC_LATENCY = "RpcLatency";
   private static final String APPEND_ROWS_ROW_STATUS = "RowsAppendedCount";
-  private static final String THROTTLED_TIME = "ThrottledTime";
+  public static final String THROTTLED_TIME = "ThrottledTime";
 
   // StorageWriteAPI Method names
-  enum RpcMethod {
+  public enum RpcMethod {
     APPEND_ROWS,
     FLUSH_ROWS,
     FINALIZE_STREAM

Reply via email to