From 5111ba953be94dc8f2014426726a840442258388 Mon Sep 17 00:00:00 2001 From: dpcollins-google <40498610+dpcollins-google@users.noreply.github.com> Date: Tue, 15 Mar 2022 14:31:05 -0400 Subject: [PATCH 1/7] Catch MonitoringInfoMetricName null keys or values in the constructor --- .../beam/runners/core/metrics/MonitoringInfoMetricName.java | 2 ++ 1 file changed, 2 insertions(+) diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MonitoringInfoMetricName.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MonitoringInfoMetricName.java index bd3ccc871a50..ba659bc4b1d4 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MonitoringInfoMetricName.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MonitoringInfoMetricName.java @@ -48,6 +48,8 @@ private MonitoringInfoMetricName(String urn, Map labels) { // and ensure all necessary labels are set for the specific URN. this.urn = urn; for (Entry entry : labels.entrySet()) { + checkArgument(entry.getKey() != null, "MonitoringInfoMetricName keys must be non-null"); + checkArgument(entry.getValue() != null, "MonitoringInfoMetricName values must be non-null"); this.labels.put(entry.getKey(), entry.getValue()); } } From 3c0e27e681b63e176bde67159a39085fc3eb9724 Mon Sep 17 00:00:00 2001 From: Daniel Collins Date: Wed, 16 Mar 2022 09:39:57 -0400 Subject: [PATCH 2/7] remove nullness suppression and fix bugs --- .../runners/core/metrics/MonitoringInfoMetricName.java | 8 ++++---- .../org/apache/beam/sdk/extensions/gcp/util/GcsUtil.java | 7 +++++-- 2 files changed, 9 insertions(+), 6 deletions(-) diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MonitoringInfoMetricName.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MonitoringInfoMetricName.java index ba659bc4b1d4..46dddcb4d73d 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MonitoringInfoMetricName.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MonitoringInfoMetricName.java @@ -33,9 +33,6 @@ * key instead of only a name+namespace. This is useful when defining system defined metrics with a * specific urn via a {@code CounterContainer}. */ -@SuppressWarnings({ - "nullness" // TODO(https://issues.apache.org/jira/browse/BEAM-10402) -}) public class MonitoringInfoMetricName extends MetricName { private String urn; @@ -49,7 +46,10 @@ private MonitoringInfoMetricName(String urn, Map labels) { this.urn = urn; for (Entry entry : labels.entrySet()) { checkArgument(entry.getKey() != null, "MonitoringInfoMetricName keys must be non-null"); - checkArgument(entry.getValue() != null, "MonitoringInfoMetricName values must be non-null"); + checkArgument( + entry.getValue() != null, + "MonitoringInfoMetricName values must be non-null, but was null for key `%s`", + entry.getKey()); this.labels.put(entry.getKey(), entry.getValue()); } } diff --git a/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtil.java b/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtil.java index 668d94bbfd89..1163a4256fca 100644 --- a/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtil.java +++ b/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtil.java @@ -85,6 +85,7 @@ import org.apache.beam.sdk.util.FluentBackoff; import org.apache.beam.sdk.util.MoreFutures; import org.apache.beam.vendor.guava.v26_0_jre.com.google.common.annotations.VisibleForTesting; +import org.apache.beam.vendor.guava.v26_0_jre.com.google.common.base.MoreObjects; import org.apache.beam.vendor.guava.v26_0_jre.com.google.common.collect.ImmutableList; import org.apache.beam.vendor.guava.v26_0_jre.com.google.common.collect.Lists; import org.apache.beam.vendor.guava.v26_0_jre.com.google.common.collect.Sets; @@ -466,7 +467,8 @@ SeekableByteChannel open(GcsPath path, GoogleCloudStorageReadOptions readOptions MonitoringInfoConstants.Labels.RESOURCE, GcpResourceIdentifiers.cloudStorageBucket(path.getBucket())); baseLabels.put( - MonitoringInfoConstants.Labels.GCS_PROJECT_ID, googleCloudStorageOptions.getProjectId()); + MonitoringInfoConstants.Labels.GCS_PROJECT_ID, + MoreObjects.firstNonNull(googleCloudStorageOptions.getProjectId(), "")); baseLabels.put(MonitoringInfoConstants.Labels.GCS_BUCKET, path.getBucket()); ServiceCallMetric serviceCallMetric = @@ -580,7 +582,8 @@ public WritableByteChannel create(GcsPath path, CreateOptions options) throws IO MonitoringInfoConstants.Labels.RESOURCE, GcpResourceIdentifiers.cloudStorageBucket(path.getBucket())); baseLabels.put( - MonitoringInfoConstants.Labels.GCS_PROJECT_ID, googleCloudStorageOptions.getProjectId()); + MonitoringInfoConstants.Labels.GCS_PROJECT_ID, + MoreObjects.firstNonNull(googleCloudStorageOptions.getProjectId(), "")); baseLabels.put(MonitoringInfoConstants.Labels.GCS_BUCKET, path.getBucket()); ServiceCallMetric serviceCallMetric = From ce4c48efcbaa63752f448d95fe99243af5bf164d Mon Sep 17 00:00:00 2001 From: Daniel Collins Date: Wed, 16 Mar 2022 10:24:32 -0400 Subject: [PATCH 3/7] fixes --- .../runners/core/metrics/MonitoringInfoMetricName.java | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MonitoringInfoMetricName.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MonitoringInfoMetricName.java index 46dddcb4d73d..6038f6a430f5 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MonitoringInfoMetricName.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MonitoringInfoMetricName.java @@ -58,13 +58,13 @@ private MonitoringInfoMetricName(String urn, Map labels) { public String getNamespace() { if (labels.containsKey(MonitoringInfoConstants.Labels.NAMESPACE)) { // User-generated metric - return labels.getOrDefault(MonitoringInfoConstants.Labels.NAMESPACE, null); + return labels.get(MonitoringInfoConstants.Labels.NAMESPACE); } else if (labels.containsKey(MonitoringInfoConstants.Labels.PCOLLECTION)) { // System-generated metric - return labels.getOrDefault(MonitoringInfoConstants.Labels.PCOLLECTION, null); + return labels.get(MonitoringInfoConstants.Labels.PCOLLECTION); } else if (labels.containsKey(MonitoringInfoConstants.Labels.PTRANSFORM)) { // System-generated metric - return labels.getOrDefault(MonitoringInfoConstants.Labels.PTRANSFORM, null); + return labels.get(MonitoringInfoConstants.Labels.PTRANSFORM); } else { return urn.split(":", 2)[0]; } @@ -73,7 +73,7 @@ public String getNamespace() { @Override public String getName() { if (labels.containsKey(MonitoringInfoConstants.Labels.NAME)) { - return labels.getOrDefault(MonitoringInfoConstants.Labels.NAME, null); + return labels.get(MonitoringInfoConstants.Labels.NAME); } else { return urn.split(":", 2)[1]; } From 5417d565ee5046d69ed81e736f324a87e6414229 Mon Sep 17 00:00:00 2001 From: Daniel Collins Date: Fri, 18 Mar 2022 18:32:17 -0400 Subject: [PATCH 4/7] fixes --- .../metrics/MonitoringInfoMetricName.java | 40 +++++++++---------- .../beam/sdk/extensions/gcp/util/GcsUtil.java | 4 +- 2 files changed, 21 insertions(+), 23 deletions(-) diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MonitoringInfoMetricName.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MonitoringInfoMetricName.java index 6038f6a430f5..2a1345df64ba 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MonitoringInfoMetricName.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MonitoringInfoMetricName.java @@ -26,6 +26,7 @@ import org.apache.beam.model.pipeline.v1.MetricsApi; import org.apache.beam.sdk.metrics.MetricName; import org.apache.beam.vendor.guava.v26_0_jre.com.google.common.base.Strings; +import org.apache.beam.vendor.guava.v26_0_jre.com.google.common.collect.ImmutableMap; import org.checkerframework.checker.nullness.qual.Nullable; /** @@ -35,8 +36,8 @@ */ public class MonitoringInfoMetricName extends MetricName { - private String urn; - private Map labels = new HashMap(); + private final String urn; + private final Map labels; private MonitoringInfoMetricName(String urn, Map labels) { checkArgument(!Strings.isNullOrEmpty(urn), "MonitoringInfoMetricName urn must be non-empty"); @@ -44,30 +45,27 @@ private MonitoringInfoMetricName(String urn, Map labels) { // TODO(ajamato): Move SimpleMonitoringInfoBuilder to :runners:core-construction-java // and ensure all necessary labels are set for the specific URN. this.urn = urn; - for (Entry entry : labels.entrySet()) { - checkArgument(entry.getKey() != null, "MonitoringInfoMetricName keys must be non-null"); - checkArgument( - entry.getValue() != null, - "MonitoringInfoMetricName values must be non-null, but was null for key `%s`", - entry.getKey()); - this.labels.put(entry.getKey(), entry.getValue()); - } + this.labels = ImmutableMap.copyOf(labels); } @Override public String getNamespace() { - if (labels.containsKey(MonitoringInfoConstants.Labels.NAMESPACE)) { - // User-generated metric - return labels.get(MonitoringInfoConstants.Labels.NAMESPACE); - } else if (labels.containsKey(MonitoringInfoConstants.Labels.PCOLLECTION)) { - // System-generated metric - return labels.get(MonitoringInfoConstants.Labels.PCOLLECTION); - } else if (labels.containsKey(MonitoringInfoConstants.Labels.PTRANSFORM)) { - // System-generated metric - return labels.get(MonitoringInfoConstants.Labels.PTRANSFORM); - } else { - return urn.split(":", 2)[0]; + // User-generated metric + String ret = labels.get(MonitoringInfoConstants.Labels.NAMESPACE); + if (ret != null) { + return ret; + } + // System-generated metric + ret = labels.get(MonitoringInfoConstants.Labels.PCOLLECTION); + if (ret != null) { + return ret; + } + // System-generated metric + ret = labels.get(MonitoringInfoConstants.Labels.PTRANSFORM); + if (ret != null) { + return ret; } + return urn.split(":", 2)[0]; } @Override diff --git a/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtil.java b/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtil.java index 1163a4256fca..4761d2831ecd 100644 --- a/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtil.java +++ b/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtil.java @@ -468,7 +468,7 @@ SeekableByteChannel open(GcsPath path, GoogleCloudStorageReadOptions readOptions GcpResourceIdentifiers.cloudStorageBucket(path.getBucket())); baseLabels.put( MonitoringInfoConstants.Labels.GCS_PROJECT_ID, - MoreObjects.firstNonNull(googleCloudStorageOptions.getProjectId(), "")); + String.valueOf(googleCloudStorageOptions.getProjectId())); baseLabels.put(MonitoringInfoConstants.Labels.GCS_BUCKET, path.getBucket()); ServiceCallMetric serviceCallMetric = @@ -583,7 +583,7 @@ public WritableByteChannel create(GcsPath path, CreateOptions options) throws IO GcpResourceIdentifiers.cloudStorageBucket(path.getBucket())); baseLabels.put( MonitoringInfoConstants.Labels.GCS_PROJECT_ID, - MoreObjects.firstNonNull(googleCloudStorageOptions.getProjectId(), "")); + String.valueOf(googleCloudStorageOptions.getProjectId())); baseLabels.put(MonitoringInfoConstants.Labels.GCS_BUCKET, path.getBucket()); ServiceCallMetric serviceCallMetric = From 2a33d94444bee11af65ef6f76c60fddc62138dcf Mon Sep 17 00:00:00 2001 From: Daniel Collins Date: Fri, 18 Mar 2022 19:08:31 -0400 Subject: [PATCH 5/7] fixes --- .../java/org/apache/beam/sdk/extensions/gcp/util/GcsUtil.java | 1 - 1 file changed, 1 deletion(-) diff --git a/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtil.java b/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtil.java index 4761d2831ecd..9e4e51680e1a 100644 --- a/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtil.java +++ b/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtil.java @@ -85,7 +85,6 @@ import org.apache.beam.sdk.util.FluentBackoff; import org.apache.beam.sdk.util.MoreFutures; import org.apache.beam.vendor.guava.v26_0_jre.com.google.common.annotations.VisibleForTesting; -import org.apache.beam.vendor.guava.v26_0_jre.com.google.common.base.MoreObjects; import org.apache.beam.vendor.guava.v26_0_jre.com.google.common.collect.ImmutableList; import org.apache.beam.vendor.guava.v26_0_jre.com.google.common.collect.Lists; import org.apache.beam.vendor.guava.v26_0_jre.com.google.common.collect.Sets; From db7e179a40b1d2662c8cd192565c4cfd1b240ef6 Mon Sep 17 00:00:00 2001 From: Daniel Collins Date: Sat, 19 Mar 2022 10:45:05 -0400 Subject: [PATCH 6/7] fixes --- .../beam/runners/core/metrics/MonitoringInfoMetricName.java | 2 -- 1 file changed, 2 deletions(-) diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MonitoringInfoMetricName.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MonitoringInfoMetricName.java index 2a1345df64ba..8dc200f52a0d 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MonitoringInfoMetricName.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MonitoringInfoMetricName.java @@ -19,9 +19,7 @@ import static org.apache.beam.vendor.guava.v26_0_jre.com.google.common.base.Preconditions.checkArgument; -import java.util.HashMap; import java.util.Map; -import java.util.Map.Entry; import java.util.Objects; import org.apache.beam.model.pipeline.v1.MetricsApi; import org.apache.beam.sdk.metrics.MetricName; From 847fe17f18991e6841fa3f5d74537b2ccee33d76 Mon Sep 17 00:00:00 2001 From: Daniel Collins Date: Tue, 22 Mar 2022 15:29:43 -0400 Subject: [PATCH 7/7] fixes --- .../runners/core/metrics/MonitoringInfoMetricName.java | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MonitoringInfoMetricName.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MonitoringInfoMetricName.java index 8dc200f52a0d..c6dda30354d5 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MonitoringInfoMetricName.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MonitoringInfoMetricName.java @@ -68,11 +68,11 @@ public String getNamespace() { @Override public String getName() { - if (labels.containsKey(MonitoringInfoConstants.Labels.NAME)) { - return labels.get(MonitoringInfoConstants.Labels.NAME); - } else { - return urn.split(":", 2)[1]; + String ret = labels.get(MonitoringInfoConstants.Labels.NAME); + if (ret != null) { + return ret; } + return urn.split(":", 2)[1]; } /** @return the urn of this MonitoringInfo metric. */