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