Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line numberDiff line numberDiff line change
Expand Up@@ -19,62 +19,60 @@

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;
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;

/**
* An implementation of {@code MetricKey} based on a MonitoringInfo's URN and label to represent the
* 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;
private Map<String, String> labels = new HashMap<String, String>();
private final String urn;
private final Map<String, String> labels;

private MonitoringInfoMetricName(String urn, Map<String, String> labels) {
checkArgument(!Strings.isNullOrEmpty(urn), "MonitoringInfoMetricName urn must be non-empty");
checkArgument(labels != null, "MonitoringInfoMetricName labels must be non-null");
// 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<String, String> entry : labels.entrySet()) {
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.getOrDefault(MonitoringInfoConstants.Labels.NAMESPACE, null);
} else if (labels.containsKey(MonitoringInfoConstants.Labels.PCOLLECTION)) {
// System-generated metric
return labels.getOrDefault(MonitoringInfoConstants.Labels.PCOLLECTION, null);
} else if (labels.containsKey(MonitoringInfoConstants.Labels.PTRANSFORM)) {
// System-generated metric
return labels.getOrDefault(MonitoringInfoConstants.Labels.PTRANSFORM, null);
} 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
public String getName() {
if (labels.containsKey(MonitoringInfoConstants.Labels.NAME)) {
return labels.getOrDefault(MonitoringInfoConstants.Labels.NAME, null);
} 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. */
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -466,7 +466,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,
String.valueOf(googleCloudStorageOptions.getProjectId()));
baseLabels.put(MonitoringInfoConstants.Labels.GCS_BUCKET, path.getBucket());

ServiceCallMetric serviceCallMetric =
Expand DownExpand Up@@ -580,7 +581,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,
String.valueOf(googleCloudStorageOptions.getProjectId()));
baseLabels.put(MonitoringInfoConstants.Labels.GCS_BUCKET, path.getBucket());

ServiceCallMetric serviceCallMetric =
Expand Down