This is an automated email from the ASF dual-hosted git repository.
shunping 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 720326829a7 Fill gaps between GcsUtilV2 and GcsUtilV1 (#40244)
720326829a7 is described below
commit 720326829a71b16b3553a246305fbf9031a70d3d
Author: Shunping Huang <[email protected]>
AuthorDate: Thu Sep 24 22:03:02 2026 -0400
Fill gaps between GcsUtilV2 and GcsUtilV1 (#40244)
* Fix a leaking channel in gcsutilv2.
* Configure GcsUtilV2's storage client from pipeline options
GcsUtilV2 left credentials to the application default chain, silently
ignoring --gcpCredentialFactoryClass, impersonation and explicit service
account keys that GcsUtilV1 honors. A custom --gcsEndpoint was ignored
outright.
Pass the pipeline's credentials through, mapping a null credential (an
explicit opt-out) to NoCredentials, and apply gcsEndpoint as the client
host.
* Emit GCS metrics from GcsUtilV2 read and write paths
GcsUtilV2 returned raw channels, so enabling use_gcsutil_v2 silently
dropped every GCS metric that GcsUtilV1 reports: the per-bucket byte
counters from GcsCountersOptions, the gcs_http_*_wire_bytes counters
gated on gcsPerformanceMetrics, and the api_request_count
ServiceCallMetric.
Wrap the channels returned by open() and create() in the existing
Counting{Seekable,Writable}ByteChannel helpers and report a
ServiceCallMetric for GcsGet and GcsInsert, matching V1's labels.
* Implement GCS Http request metrics for GcsUtilV2
* Add tests for GCS HTTP request metrics for GcsUtilV2
* Log which GcsUtil delegate is in use.
* Fix spotbug finding
* Fix another spotbug finding
* Address reviewer comments to add metric and configuration tests.
---
.../org/apache/beam/gradle/BeamModulePlugin.groovy | 1 +
.../google-cloud-platform-core/build.gradle | 1 +
.../beam/sdk/extensions/gcp/util/GcsUtil.java | 10 +-
.../beam/sdk/extensions/gcp/util/GcsUtilV2.java | 268 ++++++++++++-
.../gcp/util/GcsUtilParameterizedIT.java | 155 ++++++++
.../sdk/extensions/gcp/util/GcsUtilV2Test.java | 418 +++++++++++++++++++++
6 files changed, 840 insertions(+), 13 deletions(-)
diff --git
a/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy
b/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy
index 88c998ed910..0c34e1529e8 100644
--- a/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy
+++ b/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy
@@ -768,6 +768,7 @@ class BeamModulePlugin implements Plugin<Project> {
google_cloud_bigtable_emulator :
"com.google.cloud:google-cloud-bigtable-emulator", //
google_cloud_platform_libraries_bom sets version
google_cloud_core :
"com.google.cloud:google-cloud-core", // google_cloud_platform_libraries_bom
sets version
google_cloud_core_grpc :
"com.google.cloud:google-cloud-core-grpc", //
google_cloud_platform_libraries_bom sets version
+ google_cloud_core_http :
"com.google.cloud:google-cloud-core-http", //
google_cloud_platform_libraries_bom sets version
google_cloud_datacatalog_v1beta1 :
"com.google.cloud:google-cloud-datacatalog", //
google_cloud_platform_libraries_bom sets version
google_cloud_dataflow_java_proto_library_all:
"com.google.cloud.dataflow:google-cloud-dataflow-java-proto-library-all:0.5.160304",
google_cloud_datastore_v1_proto_client :
"com.google.cloud.datastore:datastore-v1-proto-client:3.4.0", //
[bomupgrader] sets version
diff --git a/sdks/java/extensions/google-cloud-platform-core/build.gradle
b/sdks/java/extensions/google-cloud-platform-core/build.gradle
index 75e0c508727..9f2b475a713 100644
--- a/sdks/java/extensions/google-cloud-platform-core/build.gradle
+++ b/sdks/java/extensions/google-cloud-platform-core/build.gradle
@@ -44,6 +44,7 @@ dependencies {
implementation library.java.google_auth_library_oauth2_http
implementation library.java.google_api_client
implementation library.java.google_cloud_core
+ implementation library.java.google_cloud_core_http
implementation library.java.google_cloud_storage
implementation library.java.bigdataoss_gcsio
implementation library.java.bigdataoss_util
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 070cf74d7c1..9e2d82ae1eb 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
@@ -47,8 +47,12 @@ import org.apache.beam.sdk.options.PipelineOptions;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting;
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Sets;
import org.checkerframework.checker.nullness.qual.Nullable;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
public class GcsUtil {
+ private static final Logger LOG = LoggerFactory.getLogger(GcsUtil.class);
+
/**
* Namespace for every GCS metric. The namespace is dropped when Dataflow
exports counters to
* Cloud Monitoring, so the layer is carried by the metric name instead:
{@code gcs_http_*} for
@@ -119,8 +123,12 @@ public class GcsUtil {
this.delegate = new GcsUtilV1.GcsUtilFactory().create(options);
if (ExperimentalOptions.hasExperiment(options, "use_gcsutil_v2")) {
this.delegateV2 = new GcsUtilV2.GcsUtilFactory().create(options);
+ // INFO only for V2, which is opt-in. V1 is still the default for every
pipeline,
+ // so logging it at INFO would be noise.
+ LOG.info("Using GcsUtilV2 (java-storage) for GCS operations.");
} else {
this.delegateV2 = null;
+ LOG.debug("Using GcsUtilV1 (gcsio) for GCS operations.");
}
}
@@ -297,7 +305,7 @@ public class GcsUtil {
public WritableByteChannel create(GcsPath path, CreateOptions options)
throws IOException {
if (delegateV2 != null) {
- delegateV2.create(path, options.delegate);
+ return delegateV2.create(path, options.delegate);
}
return delegate.create(path, options.delegate);
}
diff --git
a/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV2.java
b/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV2.java
index bbdab0982b5..9d223f81050 100644
---
a/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV2.java
+++
b/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV2.java
@@ -21,10 +21,16 @@ import static
org.apache.beam.sdk.io.FileSystemUtils.wildcardToRegexp;
import static
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkArgument;
import static
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkNotNull;
+import com.google.api.client.http.HttpRequestInitializer;
import com.google.api.gax.paging.Page;
+import com.google.auth.Credentials;
import com.google.auto.value.AutoValue;
+import com.google.cloud.NoCredentials;
import com.google.cloud.ReadChannel;
+import com.google.cloud.ServiceOptions;
+import com.google.cloud.TransportOptions;
import com.google.cloud.WriteChannel;
+import com.google.cloud.http.HttpTransportOptions;
import com.google.cloud.storage.Blob;
import com.google.cloud.storage.BlobId;
import com.google.cloud.storage.BlobInfo;
@@ -47,6 +53,8 @@ import com.google.cloud.storage.StorageException;
import com.google.cloud.storage.StorageOptions;
import java.io.FileNotFoundException;
import java.io.IOException;
+import java.net.MalformedURLException;
+import java.net.URL;
import java.nio.ByteBuffer;
import java.nio.channels.SeekableByteChannel;
import java.nio.channels.WritableByteChannel;
@@ -54,11 +62,25 @@ import java.nio.file.AccessDeniedException;
import java.nio.file.FileAlreadyExistsException;
import java.util.ArrayList;
import java.util.Arrays;
+import java.util.HashMap;
import java.util.List;
+import java.util.Objects;
+import java.util.Optional;
+import java.util.function.Consumer;
import java.util.regex.Pattern;
+import org.apache.beam.runners.core.metrics.GcpResourceIdentifiers;
+import org.apache.beam.runners.core.metrics.MonitoringInfoConstants;
+import org.apache.beam.runners.core.metrics.ServiceCallMetric;
import org.apache.beam.sdk.extensions.gcp.options.GcpOptions;
import org.apache.beam.sdk.extensions.gcp.options.GcsOptions;
+import
org.apache.beam.sdk.extensions.gcp.util.channels.CountingSeekableByteChannel;
+import
org.apache.beam.sdk.extensions.gcp.util.channels.CountingWritableByteChannel;
import org.apache.beam.sdk.extensions.gcp.util.gcsfs.GcsPath;
+import org.apache.beam.sdk.metrics.Counter;
+import org.apache.beam.sdk.metrics.MetricName;
+import org.apache.beam.sdk.metrics.Metrics;
+import org.apache.beam.sdk.metrics.MetricsContainer;
+import org.apache.beam.sdk.metrics.MetricsEnvironment;
import org.apache.beam.sdk.options.DefaultValueFactory;
import org.apache.beam.sdk.options.PipelineOptions;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting;
@@ -82,6 +104,12 @@ class GcsUtilV2 {
private final @Nullable Integer uploadBufferSizeBytes;
+ private final GcsUtilV1.GcsCountersOptions gcsCountersOptions;
+
+ private final boolean gcsPerformanceMetrics;
+
+ private final @Nullable String projectId;
+
/** Maximum number of items to retrieve per Objects.List request. */
private static final long MAX_LIST_BLOBS_PER_CALL = 1024;
@@ -95,9 +123,204 @@ class GcsUtilV2 {
private static final long MEGABYTES_COPIED_PER_CHUNK = 2048L;
GcsUtilV2(PipelineOptions options) {
- String projectId = options.as(GcpOptions.class).getProject();
- storage =
StorageOptions.newBuilder().setProjectId(projectId).build().getService();
- uploadBufferSizeBytes =
options.as(GcsOptions.class).getGcsUploadBufferSizeBytes();
+ GcsOptions gcsOptions = options.as(GcsOptions.class);
+ this.projectId = options.as(GcpOptions.class).getProject();
+ StorageOptions.Builder storageOptionsBuilder =
+ StorageOptions.newBuilder().setProjectId(this.projectId);
+
+ // Use the pipeline's configured credentials rather than falling back to
application default
+ // credentials, so that --gcpCredentialFactoryClass, impersonation and
explicit service account
+ // keys are honored. A null credential means the pipeline opted out of
authentication
+ // (e.g. NoopCredentialFactory), which maps to NoCredentials for this
client.
+ Credentials credentials = gcsOptions.getGcpCredential();
+ storageOptionsBuilder.setCredentials(
+ credentials != null ? credentials : NoCredentials.getInstance());
+
+ // GcsOptions#getGcsEndpoint may carry a service path (as the JSON client
in Transport expects),
+ // but this client derives its own path, so only the root is applicable
here.
+ String endpoint = gcsOptions.getGcsEndpoint();
+ if (endpoint != null) {
+ storageOptionsBuilder.setHost(rootUrlOf(endpoint));
+ }
+
+ storage = storageOptionsBuilder.build().getService();
+ uploadBufferSizeBytes = gcsOptions.getGcsUploadBufferSizeBytes();
+ this.gcsCountersOptions =
+ GcsUtilV1.GcsCountersOptions.create(
+ gcsOptions.getEnableBucketReadMetricCounter()
+ ? gcsOptions.getGcsReadCounterPrefix()
+ : null,
+ gcsOptions.getEnableBucketWriteMetricCounter()
+ ? gcsOptions.getGcsWriteCounterPrefix()
+ : null);
+ this.gcsPerformanceMetrics =
Boolean.TRUE.equals(gcsOptions.getGcsPerformanceMetrics());
+ }
+
+ /**
+ * Creates an integer consumer that updates the counter identified by a
prefix and a bucket name.
+ */
+ private static Consumer<Integer> createCounterConsumer(String
counterNamePrefix, String bucket) {
+ return Metrics.counter(GcsUtil.class, String.format("%s_%s",
counterNamePrefix, bucket))::inc;
+ }
+
+ /** Returns the {@link MetricsContainer} to attribute wire-byte counters to,
if enabled. */
+ @VisibleForTesting
+ @Nullable MetricsContainer performanceMetricsContainer() {
+ return gcsPerformanceMetrics ? MetricsEnvironment.getCurrentContainer() :
null;
+ }
+
+ @VisibleForTesting
+ WritableByteChannel wrapInCounting(
+ WritableByteChannel writableByteChannel,
+ String bucket,
+ @Nullable MetricsContainer container) {
+ Consumer<Integer> writeConsumer =
+ Optional.ofNullable(gcsCountersOptions.getWriteCounterPrefix())
+ .map(prefix -> createCounterConsumer(prefix, bucket))
+ .orElse(null);
+
+ if (this.gcsPerformanceMetrics && container != null) {
+ Counter perfWriteCounter =
+ container.getCounter(
+ MetricName.named(GcsUtil.METRIC_NAMESPACE,
"gcs_http_write_wire_bytes_sent"));
+ Consumer<Integer> perfConsumer = perfWriteCounter::inc;
+ writeConsumer = writeConsumer == null ? perfConsumer :
writeConsumer.andThen(perfConsumer);
+ }
+
+ if (writeConsumer == null) {
+ return writableByteChannel;
+ }
+ return new CountingWritableByteChannel(writableByteChannel, writeConsumer);
+ }
+
+ @VisibleForTesting
+ SeekableByteChannel wrapInCounting(
+ SeekableByteChannel seekableByteChannel,
+ String bucket,
+ @Nullable MetricsContainer container) {
+ Consumer<Integer> readConsumer =
+ Optional.ofNullable(gcsCountersOptions.getReadCounterPrefix())
+ .map(prefix -> createCounterConsumer(prefix, bucket))
+ .orElse(null);
+
+ if (this.gcsPerformanceMetrics && container != null) {
+ Counter perfReadCounter =
+ container.getCounter(
+ MetricName.named(GcsUtil.METRIC_NAMESPACE,
"gcs_http_read_wire_bytes_received"));
+ Consumer<Integer> perfConsumer = perfReadCounter::inc;
+ readConsumer = readConsumer == null ? perfConsumer :
readConsumer.andThen(perfConsumer);
+ }
+
+ if (readConsumer == null) {
+ return seekableByteChannel;
+ }
+ return CountingSeekableByteChannel.createWithBytesReadConsumer(
+ seekableByteChannel, readConsumer);
+ }
+
+ /** Builds the API request metric for {@code method} (e.g. {@code GcsGet})
on {@code bucket}. */
+ private ServiceCallMetric serviceCallMetric(String method, String bucket) {
+ HashMap<String, String> baseLabels = new HashMap<>();
+ baseLabels.put(MonitoringInfoConstants.Labels.PTRANSFORM, "");
+ baseLabels.put(MonitoringInfoConstants.Labels.SERVICE, "Storage");
+ baseLabels.put(MonitoringInfoConstants.Labels.METHOD, method);
+ baseLabels.put(
+ MonitoringInfoConstants.Labels.RESOURCE,
GcpResourceIdentifiers.cloudStorageBucket(bucket));
+ baseLabels.put(MonitoringInfoConstants.Labels.GCS_PROJECT_ID,
String.valueOf(projectId));
+ baseLabels.put(MonitoringInfoConstants.Labels.GCS_BUCKET, bucket);
+ return new
ServiceCallMetric(MonitoringInfoConstants.Urns.API_REQUEST_COUNT, baseLabels);
+ }
+
+ /**
+ * {@link HttpTransportOptions} that wraps the request initializer so that
HTTP-level counters
+ * (request counts, request shape, and status classes) are incremented
against a pre-bound {@link
+ * MetricsContainer}.
+ *
+ * <p>The container is bound eagerly rather than resolved per request
because requests may execute
+ * on background threads, where {@link
MetricsEnvironment#getCurrentContainer} would not resolve
+ * to the step that initiated the operation.
+ */
+ private static class MetricsHttpTransportOptions extends
HttpTransportOptions {
+ private static final long serialVersionUID = 1L;
+
+ // Not serializable, and only meaningful in the process that created it.
+ private final transient @Nullable MetricsContainer container;
+ private final boolean isWrite;
+
+ MetricsHttpTransportOptions(
+ HttpTransportOptions base, @Nullable MetricsContainer container,
boolean isWrite) {
+ super(base.toBuilder());
+ this.container = container;
+ this.isWrite = isWrite;
+ }
+
+ @Override
+ public HttpRequestInitializer getHttpRequestInitializer(ServiceOptions<?,
?> serviceOptions) {
+ // withMetricsContainer returns the delegate unchanged when the
container is null, which is
+ // also the case after deserialization.
+ return Transport.withMetricsContainer(
+ super.getHttpRequestInitializer(serviceOptions), container, isWrite);
+ }
+
+ @Override
+ public boolean equals(@Nullable Object obj) {
+ if (this == obj) {
+ return true;
+ }
+ if (!(obj instanceof MetricsHttpTransportOptions)) {
+ return false;
+ }
+ if (!super.equals(obj)) {
+ return false;
+ }
+ MetricsHttpTransportOptions other = (MetricsHttpTransportOptions) obj;
+ return isWrite == other.isWrite && Objects.equals(container,
other.container);
+ }
+
+ @Override
+ public int hashCode() {
+ return Objects.hash(super.hashCode(), container, isWrite);
+ }
+ }
+
+ /**
+ * Returns a client whose HTTP requests are counted against {@code
container}, or the shared
+ * client when HTTP metrics are not being collected.
+ *
+ * <p>A distinct client is required because the interceptors are installed
on the transport, which
+ * is fixed when the client is built. This mirrors {@code GcsUtilV1}, which
likewise builds a
+ * scoped client per operation while performance metrics are enabled.
+ */
+ @VisibleForTesting
+ Storage storageWithHttpMetrics(@Nullable MetricsContainer container, boolean
isWrite) {
+ if (container == null) {
+ return storage;
+ }
+ StorageOptions options = storage.getOptions();
+ TransportOptions transportOptions = options.getTransportOptions();
+ if (!(transportOptions instanceof HttpTransportOptions)) {
+ // A non-HTTP transport (e.g. gRPC) has no HttpRequestInitializer to
wrap.
+ return storage;
+ }
+ return options.toBuilder()
+ .setTransportOptions(
+ new MetricsHttpTransportOptions(
+ (HttpTransportOptions) transportOptions, container, isWrite))
+ .build()
+ .getService();
+ }
+
+ /** Returns the {@code scheme://host[:port]} prefix of {@code endpoint},
discarding any path. */
+ private static String rootUrlOf(String endpoint) {
+ try {
+ URL url = new URL(endpoint);
+ return url.getProtocol()
+ + "://"
+ + url.getHost()
+ + (url.getPort() > 0 ? ":" + url.getPort() : "");
+ } catch (MalformedURLException e) {
+ throw new IllegalArgumentException("Invalid gcsEndpoint URL: " +
endpoint, e);
+ }
}
// AccessDeniedException/FileAlreadyExistsException permit a null "other"
argument, and these
@@ -123,8 +346,14 @@ class GcsUtilV2 {
}
public Blob getBlob(GcsPath gcsPath, BlobGetOption... options) throws
IOException {
+ return getBlob(storage, gcsPath, options);
+ }
+
+ /** As {@link #getBlob(GcsPath, BlobGetOption...)}, but issued through a
specific client. */
+ private Blob getBlob(Storage client, GcsPath gcsPath, BlobGetOption...
options)
+ throws IOException {
try {
- Blob blob = storage.get(gcsPath.getBucket(), gcsPath.getObject(),
options);
+ Blob blob = client.get(gcsPath.getBucket(), gcsPath.getObject(),
options);
if (blob == null) {
throw new FileNotFoundException(
String.format("The specified file does not exist: %s",
gcsPath.toString()));
@@ -561,11 +790,21 @@ class GcsUtilV2 {
public SeekableByteChannel open(GcsPath path, BlobSourceOption...
sourceOptions)
throws IOException {
- Blob blob = getBlob(path, BlobGetOption.fields(BlobField.SIZE));
- ReadChannel reader = blob.getStorage().reader(blob.getBlobId(),
sourceOptions);
- // disable internal buffering, and make the channel non-blocking
- reader.setChunkSize(0);
- return new GcsSeekableByteChannel(reader, blob.getSize());
+ ServiceCallMetric serviceCallMetric = serviceCallMetric("GcsGet",
path.getBucket());
+ MetricsContainer container = performanceMetricsContainer();
+ try {
+ Storage client = storageWithHttpMetrics(container, false);
+ Blob blob = getBlob(client, path, BlobGetOption.fields(BlobField.SIZE));
+ ReadChannel reader = client.reader(blob.getBlobId(), sourceOptions);
+ // disable internal buffering, and make the channel non-blocking
+ reader.setChunkSize(0);
+ serviceCallMetric.call("ok");
+ return wrapInCounting(
+ new GcsSeekableByteChannel(reader, blob.getSize()),
path.getBucket(), container);
+ } catch (StorageException e) {
+ serviceCallMetric.call(e.getCode());
+ throw translateStorageException(path, e);
+ }
}
/** A bridge that allows a GCS WriteChannel to behave as a
WritableByteChannel. */
@@ -601,6 +840,8 @@ class GcsUtilV2 {
public WritableByteChannel create(
GcsPath path, GcsUtilV1.CreateOptions options, BlobWriteOption...
writeOptions)
throws IOException {
+ ServiceCallMetric serviceCallMetric = serviceCallMetric("GcsInsert",
path.getBucket());
+ MetricsContainer container = performanceMetricsContainer();
try {
// Define the metadata for the new object
BlobInfo.Builder builder = BlobInfo.newBuilder(path.getBucket(),
path.getObject());
@@ -611,13 +852,14 @@ class GcsUtilV2 {
BlobInfo blobInfo = builder.build();
+ Storage client = storageWithHttpMetrics(container, true);
List<BlobWriteOption> writeOptionList = new
ArrayList<>(Arrays.asList(writeOptions));
if (options.getExpectFileToNotExist()) {
writeOptionList.add(BlobWriteOption.doesNotExist());
} else {
// We do not merge this check with the getExpectFileToNotExist()
branch above
// because we don't want to always make the storage.get() RPC call.
- Blob blob = storage.get(path.getBucket(), path.getObject());
+ Blob blob = client.get(path.getBucket(), path.getObject());
if (blob == null) {
writeOptionList.add(BlobWriteOption.doesNotExist());
} else {
@@ -626,7 +868,7 @@ class GcsUtilV2 {
}
// Open a WriteChannel from the storage service
WriteChannel writer =
- storage.writer(blobInfo, writeOptionList.toArray(new
BlobWriteOption[0]));
+ client.writer(blobInfo, writeOptionList.toArray(new
BlobWriteOption[0]));
Integer uploadBufferSizeBytes =
options.getUploadBufferSizeBytes() != null
? options.getUploadBufferSizeBytes()
@@ -635,10 +877,12 @@ class GcsUtilV2 {
writer.setChunkSize(uploadBufferSizeBytes);
}
+ serviceCallMetric.call("ok");
// Return the bridge wrapper
- return new GcsWritableByteChannel(writer, path);
+ return wrapInCounting(new GcsWritableByteChannel(writer, path),
path.getBucket(), container);
} catch (StorageException e) {
+ serviceCallMetric.call(e.getCode());
throw translateStorageException(path, e);
}
}
diff --git
a/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilParameterizedIT.java
b/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilParameterizedIT.java
index 5759bb10a65..17752337ce9 100644
---
a/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilParameterizedIT.java
+++
b/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilParameterizedIT.java
@@ -20,6 +20,8 @@ package org.apache.beam.sdk.extensions.gcp.util;
import static org.junit.Assert.assertArrayEquals;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertNotNull;
+import static org.junit.Assert.assertNull;
import static org.junit.Assert.assertThrows;
import static org.junit.Assert.assertTrue;
@@ -45,13 +47,22 @@ import java.security.NoSuchAlgorithmException;
import java.util.Arrays;
import java.util.Collections;
import java.util.List;
+import java.util.Map;
import java.util.stream.Collectors;
+import org.apache.beam.runners.core.metrics.CounterCell;
+import org.apache.beam.runners.core.metrics.GcpResourceIdentifiers;
+import org.apache.beam.runners.core.metrics.MetricUpdates.MetricUpdate;
+import org.apache.beam.runners.core.metrics.MetricsContainerImpl;
+import org.apache.beam.runners.core.metrics.MonitoringInfoConstants;
+import org.apache.beam.runners.core.metrics.MonitoringInfoMetricName;
import org.apache.beam.sdk.extensions.gcp.options.GcsOptions;
import org.apache.beam.sdk.extensions.gcp.util.GcsUtil.CreateOptions;
import org.apache.beam.sdk.extensions.gcp.util.GcsUtilV2.MissingStrategy;
import org.apache.beam.sdk.extensions.gcp.util.GcsUtilV2.OverwriteStrategy;
import org.apache.beam.sdk.extensions.gcp.util.gcsfs.GcsPath;
import org.apache.beam.sdk.io.fs.MoveOptions;
+import org.apache.beam.sdk.metrics.MetricName;
+import org.apache.beam.sdk.metrics.MetricsEnvironment;
import org.apache.beam.sdk.options.ExperimentalOptions;
import org.apache.beam.sdk.testing.TestPipeline;
import org.apache.beam.sdk.testing.TestPipelineOptions;
@@ -75,6 +86,9 @@ import org.junit.runners.Parameterized.Parameters;
@Category(UsesKms.class)
public class GcsUtilParameterizedIT {
+ private static final String READ_COUNTER_PREFIX = "it_read_bytes";
+ private static final String WRITE_COUNTER_PREFIX = "it_write_bytes";
+
@Parameters(name = "{0}")
public static Iterable<String> data() {
return Arrays.asList("use_gcsutil_v1", "use_gcsutil_v2");
@@ -683,4 +697,145 @@ public class GcsUtilParameterizedIT {
tearDownTestBucketHelper(bucketName);
}
}
+
+ //
---------------------------------------------------------------------------------------------
+ // Metrics parity: the same assertions must hold for both GcsUtilV1 and
GcsUtilV2.
+ //
---------------------------------------------------------------------------------------------
+
+ /** Returns a {@link GcsUtil} with every GCS metric flag turned on. */
+ private GcsUtil gcsUtilWithAllMetrics() {
+ GcsOptions gcsOptions = options.as(GcsOptions.class);
+ gcsOptions.setGcsPerformanceMetrics(true);
+ gcsOptions.setEnableBucketReadMetricCounter(true);
+ gcsOptions.setEnableBucketWriteMetricCounter(true);
+ gcsOptions.setGcsReadCounterPrefix(READ_COUNTER_PREFIX);
+ gcsOptions.setGcsWriteCounterPrefix(WRITE_COUNTER_PREFIX);
+ // Built directly, as getGcsUtil() returns the instance cached in setUp().
+ return new GcsUtil(gcsOptions);
+ }
+
+ private static long counter(MetricsContainerImpl container, MetricName name)
{
+ CounterCell cell = container.tryGetCounter(name);
+ assertNotNull("counter " + name + " was not reported", cell);
+ return cell.getCumulative();
+ }
+
+ private static long gcsCounter(MetricsContainerImpl container, String name) {
+ return counter(container, MetricName.named(GcsUtil.METRIC_NAMESPACE,
name));
+ }
+
+ private static long bucketCounter(MetricsContainerImpl container, String
prefix, String bucket) {
+ return counter(container, MetricName.named(GcsUtil.class, prefix + "_" +
bucket));
+ }
+
+ /**
+ * Sums the API request counter for {@code method} and {@code status} on
{@code bucket}.
+ *
+ * <p>The {@code GCS_PROJECT_ID} label is deliberately not matched: V1
reports the project of its
+ * gcsio options (which is never set, so it reports {@code "null"}), while
V2 reports the
+ * pipeline's project.
+ */
+ private static long apiRequestCount(
+ MetricsContainerImpl container, String method, String status, String
bucket) {
+ long total = 0;
+ for (MetricUpdate<Long> update :
container.getCumulative().counterUpdates()) {
+ MetricName name = update.getKey().metricName();
+ if (!(name instanceof MonitoringInfoMetricName)) {
+ continue;
+ }
+ MonitoringInfoMetricName miName = (MonitoringInfoMetricName) name;
+ Map<String, String> labels = miName.getLabels();
+ if
(MonitoringInfoConstants.Urns.API_REQUEST_COUNT.equals(miName.getUrn())
+ &&
"Storage".equals(labels.get(MonitoringInfoConstants.Labels.SERVICE))
+ && method.equals(labels.get(MonitoringInfoConstants.Labels.METHOD))
+ && status.equals(labels.get(MonitoringInfoConstants.Labels.STATUS))
+ &&
bucket.equals(labels.get(MonitoringInfoConstants.Labels.GCS_BUCKET))
+ && GcpResourceIdentifiers.cloudStorageBucket(bucket)
+ .equals(labels.get(MonitoringInfoConstants.Labels.RESOURCE))) {
+ total += update.getUpdate();
+ }
+ }
+ return total;
+ }
+
+ @Test
+ public void testReadMetrics() throws IOException {
+ final String bucket = "apache-beam-samples";
+ final GcsPath gcsPath = GcsPath.fromComponents(bucket,
"shakespeare/kinglear.txt");
+ final long expectedSize = 157283L;
+ GcsUtil metricsGcsUtil = gcsUtilWithAllMetrics();
+
+ MetricsContainerImpl container = new MetricsContainerImpl("step");
+ MetricsContainerImpl processWide = new MetricsContainerImpl(null);
+ MetricsEnvironment.setCurrentContainer(container);
+ MetricsEnvironment.setProcessWideContainer(processWide);
+ try {
+ try (SeekableByteChannel channel = metricsGcsUtil.open(gcsPath)) {
+ ByteBuffer buffer = ByteBuffer.allocate((int) expectedSize + 1024);
+ assertEquals(expectedSize,
StorageChannelUtils.blockingFillFrom(buffer, channel));
+ }
+ } finally {
+ MetricsEnvironment.setCurrentContainer(null);
+ MetricsEnvironment.setProcessWideContainer(null);
+ }
+
+ // --enableBucketReadMetricCounter
+ assertEquals(expectedSize, bucketCounter(container, READ_COUNTER_PREFIX,
bucket));
+ // --gcsPerformanceMetrics: wire bytes and HTTP counters
+ assertEquals(expectedSize, gcsCounter(container,
"gcs_http_read_wire_bytes_received"));
+ assertTrue(gcsCounter(container, "gcs_http_read_request_count") >= 1);
+ assertTrue(gcsCounter(container, "gcs_http_read_status_2xx") >= 1);
+ assertEquals(
+ gcsCounter(container, "gcs_http_read_request_count"),
+ gcsCounter(container, "gcs_http_read_request_count_ranged")
+ + gcsCounter(container, "gcs_http_read_request_count_unbounded")
+ + gcsCounter(container, "gcs_http_read_request_count_other"));
+ // A read must not produce write-side counters.
+ assertNull(
+ container.tryGetCounter(
+ MetricName.named(GcsUtil.METRIC_NAMESPACE,
"gcs_http_write_request_count")));
+ // API request metric
+ assertEquals(1, apiRequestCount(processWide, "GcsGet", "ok", bucket));
+ }
+
+ @Test
+ public void testWriteMetrics() throws IOException {
+ final String bucket =
+ "apache-beam-temp-metrics-" +
java.util.UUID.randomUUID().toString().substring(0, 8);
+ final GcsPath targetPath = GcsPath.fromComponents(bucket,
"test-object.txt");
+ final byte[] content = "Hello, GCS
metrics!".getBytes(StandardCharsets.UTF_8);
+ GcsUtil metricsGcsUtil = gcsUtilWithAllMetrics();
+
+ MetricsContainerImpl container = new MetricsContainerImpl("step");
+ MetricsContainerImpl processWide = new MetricsContainerImpl(null);
+ try {
+ createTestBucketHelper(bucket, false);
+
+ MetricsEnvironment.setCurrentContainer(container);
+ MetricsEnvironment.setProcessWideContainer(processWide);
+ try (WritableByteChannel writer =
+ metricsGcsUtil.create(
+ targetPath,
CreateOptions.builder().setExpectFileToNotExist(true).build())) {
+ writer.write(ByteBuffer.wrap(content));
+ } finally {
+ MetricsEnvironment.setCurrentContainer(null);
+ MetricsEnvironment.setProcessWideContainer(null);
+ }
+ } finally {
+ tearDownTestBucketHelper(bucket);
+ }
+
+ // --enableBucketWriteMetricCounter
+ assertEquals(content.length, bucketCounter(container,
WRITE_COUNTER_PREFIX, bucket));
+ // --gcsPerformanceMetrics: wire bytes and HTTP counters
+ assertEquals(content.length, gcsCounter(container,
"gcs_http_write_wire_bytes_sent"));
+ assertTrue(gcsCounter(container, "gcs_http_write_request_count") >= 1);
+ assertTrue(gcsCounter(container, "gcs_http_write_status_2xx") >= 1);
+ // A write must not produce read-side counters.
+ assertNull(
+ container.tryGetCounter(
+ MetricName.named(GcsUtil.METRIC_NAMESPACE,
"gcs_http_read_request_count")));
+ // API request metric
+ assertEquals(1, apiRequestCount(processWide, "GcsInsert", "ok", bucket));
+ }
}
diff --git
a/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV2Test.java
b/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV2Test.java
new file mode 100644
index 00000000000..2b9e55b98b8
--- /dev/null
+++
b/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV2Test.java
@@ -0,0 +1,418 @@
+/*
+ * 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.sdk.extensions.gcp.util;
+
+import static java.nio.charset.StandardCharsets.UTF_8;
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertNotSame;
+import static org.junit.Assert.assertNull;
+import static org.junit.Assert.assertSame;
+import static org.junit.Assert.assertThrows;
+
+import com.google.api.client.http.GenericUrl;
+import com.google.api.client.http.HttpRequestInitializer;
+import com.google.api.client.testing.http.MockHttpTransport;
+import com.google.api.client.testing.http.MockLowLevelHttpResponse;
+import com.google.auth.Credentials;
+import com.google.cloud.NoCredentials;
+import com.google.cloud.http.HttpTransportOptions;
+import com.google.cloud.storage.Storage;
+import com.google.cloud.storage.StorageOptions;
+import java.io.ByteArrayOutputStream;
+import java.io.IOException;
+import java.nio.ByteBuffer;
+import java.nio.channels.Channels;
+import java.nio.channels.SeekableByteChannel;
+import java.nio.channels.WritableByteChannel;
+import
org.apache.beam.repackaged.core.org.apache.commons.compress.utils.SeekableInMemoryByteChannel;
+import org.apache.beam.runners.core.metrics.CounterCell;
+import org.apache.beam.runners.core.metrics.MetricsContainerImpl;
+import org.apache.beam.sdk.extensions.gcp.auth.NoopCredentialFactory;
+import org.apache.beam.sdk.extensions.gcp.auth.TestCredential;
+import org.apache.beam.sdk.extensions.gcp.options.GcsOptions;
+import org.apache.beam.sdk.metrics.MetricName;
+import org.apache.beam.sdk.metrics.MetricsEnvironment;
+import org.apache.beam.sdk.options.PipelineOptionsFactory;
+import org.checkerframework.checker.nullness.qual.Nullable;
+import org.junit.After;
+import org.junit.Test;
+import org.junit.runner.RunWith;
+import org.junit.runners.JUnit4;
+
+/**
+ * Unit tests for {@link GcsUtilV2}.
+ *
+ * <p>Covers two things without making any request to GCS:
+ *
+ * <ul>
+ * <li>The storage client is configured from the pipeline options (project,
credentials,
+ * endpoint), matching what {@link GcsUtilV1} honors.
+ * <li>HTTP metrics and byte counters are installed only when their flags
ask for them, and only
+ * when there is a container to report into. The counting logic itself
is covered by {@link
+ * TransportTest}.
+ * </ul>
+ *
+ * <p>End-to-end parity with {@link GcsUtilV1} against real GCS is covered by
{@link
+ * GcsUtilParameterizedIT}.
+ */
+@RunWith(JUnit4.class)
+public class GcsUtilV2Test {
+
+ private static final String BUCKET = "test-bucket";
+ private static final String READ_PREFIX = "test_read_prefix";
+ private static final String WRITE_PREFIX = "test_write_prefix";
+ private static final byte[] PAYLOAD = "some_bytes".getBytes(UTF_8);
+
+ @After
+ public void tearDown() {
+ MetricsEnvironment.setCurrentContainer(null);
+ }
+
+ private static GcsOptions gcsOptions() {
+ GcsOptions options = PipelineOptionsFactory.as(GcsOptions.class);
+ // Avoid resolving application default credentials; no request leaves the
process.
+ options.setGcpCredential(new TestCredential());
+ options.setProject("test-project");
+ return options;
+ }
+
+ private static GcsUtilV2 gcsUtilV2(boolean performanceMetrics) {
+ GcsOptions options = gcsOptions();
+ options.setGcsPerformanceMetrics(performanceMetrics);
+ return new GcsUtilV2(options);
+ }
+
+ private static GcsUtilV2 gcsUtilV2(boolean performanceMetrics, boolean
bucketCounters) {
+ GcsOptions options = gcsOptions();
+ options.setGcsPerformanceMetrics(performanceMetrics);
+ options.setEnableBucketReadMetricCounter(bucketCounters);
+ options.setEnableBucketWriteMetricCounter(bucketCounters);
+ options.setGcsReadCounterPrefix(READ_PREFIX);
+ options.setGcsWriteCounterPrefix(WRITE_PREFIX);
+ return new GcsUtilV2(options);
+ }
+
+ /** The shared (unscoped) storage client of {@code gcsUtil}. */
+ private static StorageOptions storageOptionsOf(GcsUtilV2 gcsUtil) {
+ return gcsUtil.storageWithHttpMetrics(null, false).getOptions();
+ }
+
+ //
---------------------------------------------------------------------------------------------
+ // Configuration
+ //
---------------------------------------------------------------------------------------------
+
+ @Test
+ public void testStorageClientUsesPipelineProject() {
+ assertEquals("test-project",
storageOptionsOf(gcsUtilV2(false)).getProjectId());
+ }
+
+ @Test
+ public void testStorageClientUsesPipelineCredentials() {
+ GcsOptions options = gcsOptions();
+ Credentials credentials = options.getGcpCredential();
+
+ assertSame(credentials, storageOptionsOf(new
GcsUtilV2(options)).getCredentials());
+ }
+
+ @Test
+ public void testNullCredentialMapsToNoCredentials() {
+ GcsOptions options = gcsOptions();
+ options.setCredentialFactoryClass(NoopCredentialFactory.class);
+ options.setGcpCredential(null);
+
+ assertSame(
+ NoCredentials.getInstance(), storageOptionsOf(new
GcsUtilV2(options)).getCredentials());
+ }
+
+ @Test
+ public void testDefaultHostWithoutGcsEndpoint() {
+ String defaultHost =
+ StorageOptions.newBuilder()
+ .setProjectId("test-project")
+ .setCredentials(NoCredentials.getInstance())
+ .build()
+ .getHost();
+
+ assertEquals(defaultHost, storageOptionsOf(gcsUtilV2(false)).getHost());
+ }
+
+ /** Mirrors {@code GcsUtilTest#testGcsEndpoint}: only the root of the
endpoint applies to V2. */
+ @Test
+ public void testGcsEndpointRootIsUsedAsHost() {
+ GcsOptions options = gcsOptions();
+ options.setGcsEndpoint("http://localhost:4443/storage/v1/");
+
+ assertEquals("http://localhost:4443", storageOptionsOf(new
GcsUtilV2(options)).getHost());
+ }
+
+ @Test
+ public void testGcsEndpointWithoutPort() {
+ GcsOptions options = gcsOptions();
+ options.setGcsEndpoint("https://storage.example.com/storage/v1/");
+
+ assertEquals("https://storage.example.com", storageOptionsOf(new
GcsUtilV2(options)).getHost());
+ }
+
+ @Test
+ public void testInvalidGcsEndpointIsRejected() {
+ GcsOptions options = gcsOptions();
+ options.setGcsEndpoint("not a url");
+
+ assertThrows(IllegalArgumentException.class, () -> new GcsUtilV2(options));
+ }
+
+ //
---------------------------------------------------------------------------------------------
+ // HTTP metrics binding
+ //
---------------------------------------------------------------------------------------------
+
+ @Test
+ public void testSharedClientIsReusedWithoutAContainer() {
+ GcsUtilV2 gcsUtil = gcsUtilV2(true);
+
+ Storage read = gcsUtil.storageWithHttpMetrics(null, false);
+ Storage write = gcsUtil.storageWithHttpMetrics(null, true);
+
+ // With nothing to count, no per-operation client should be built.
+ assertSame(read, write);
+ }
+
+ @Test
+ public void testScopedClientIsBuiltForAContainer() {
+ GcsUtilV2 gcsUtil = gcsUtilV2(true);
+ MetricsContainerImpl container = new MetricsContainerImpl(null);
+
+ Storage shared = gcsUtil.storageWithHttpMetrics(null, false);
+ Storage scoped = gcsUtil.storageWithHttpMetrics(container, false);
+
+ assertNotSame(shared, scoped);
+ assertNotSame(
+ shared.getOptions().getTransportOptions(),
scoped.getOptions().getTransportOptions());
+ // Everything other than the transport must carry over from the shared
client.
+ assertEquals(shared.getOptions().getProjectId(),
scoped.getOptions().getProjectId());
+ assertEquals(shared.getOptions().getHost(), scoped.getOptions().getHost());
+ assertSame(shared.getOptions().getCredentials(),
scoped.getOptions().getCredentials());
+ }
+
+ private static long counter(MetricsContainerImpl container, String name) {
+ return container.getCounter(MetricName.named(GcsUtil.METRIC_NAMESPACE,
name)).getCumulative();
+ }
+
+ /** Executes one request through {@code initializer} against a mock
transport. */
+ private static void executeOneRequest(HttpRequestInitializer initializer)
throws IOException {
+ MockHttpTransport transport =
+ new MockHttpTransport.Builder()
+ .setLowLevelHttpResponse(new
MockLowLevelHttpResponse().setStatusCode(200))
+ .build();
+ transport
+ .createRequestFactory(initializer)
+ .buildGetRequest(new GenericUrl("https://storage.googleapis.com/test"))
+ .execute();
+ }
+
+ private static HttpRequestInitializer initializerOf(Storage client) {
+ return ((HttpTransportOptions) client.getOptions().getTransportOptions())
+ .getHttpRequestInitializer(client.getOptions());
+ }
+
+ /**
+ * The scoped client must hand out a request initializer that counts, since
that is the only way
+ * the HTTP counters reach the container.
+ */
+ @Test
+ public void testScopedClientInitializerCountsRequests() throws IOException {
+ GcsUtilV2 gcsUtil = gcsUtilV2(true);
+ MetricsContainerImpl container = new MetricsContainerImpl(null);
+
+ executeOneRequest(initializerOf(gcsUtil.storageWithHttpMetrics(container,
false)));
+
+ assertEquals(1, counter(container, "gcs_http_read_request_count"));
+ assertEquals(1, counter(container, "gcs_http_read_status_2xx"));
+ }
+
+ @Test
+ public void testWriteDirectionIsCountedSeparately() throws IOException {
+ GcsUtilV2 gcsUtil = gcsUtilV2(true);
+ MetricsContainerImpl container = new MetricsContainerImpl(null);
+
+ executeOneRequest(initializerOf(gcsUtil.storageWithHttpMetrics(container,
true)));
+
+ assertEquals(1, counter(container, "gcs_http_write_request_count"));
+ assertEquals(0, counter(container, "gcs_http_read_request_count"));
+ }
+
+ /**
+ * The flag is the actual gate: with performance metrics disabled no
container is resolved, so
+ * open() and create() pass null and no scoped client or HTTP counter is
ever produced.
+ */
+ @Test
+ public void testPerformanceMetricsFlagGatesTheContainer() {
+ MetricsContainerImpl current = new MetricsContainerImpl(null);
+ MetricsEnvironment.setCurrentContainer(current);
+ try {
+ assertSame(current, gcsUtilV2(true).performanceMetricsContainer());
+ assertNull(gcsUtilV2(false).performanceMetricsContainer());
+ } finally {
+ MetricsEnvironment.setCurrentContainer(null);
+ }
+ }
+
+ //
---------------------------------------------------------------------------------------------
+ // Byte counters (mirrors GcsUtilTest#testReadMetrics / #testWriteMetrics
for V1)
+ //
---------------------------------------------------------------------------------------------
+
+ private static @Nullable Long counterOrNull(MetricsContainerImpl container,
MetricName name) {
+ CounterCell cell = container.tryGetCounter(name);
+ return cell == null ? null : cell.getCumulative();
+ }
+
+ private static @Nullable Long bucketCounter(MetricsContainerImpl container,
String prefix) {
+ return counterOrNull(container, MetricName.named(GcsUtil.class, prefix +
"_" + BUCKET));
+ }
+
+ private static @Nullable Long gcsCounter(MetricsContainerImpl container,
String name) {
+ return counterOrNull(container, MetricName.named(GcsUtil.METRIC_NAMESPACE,
name));
+ }
+
+ /** Reads {@link #PAYLOAD} through V2's read wrapper, bound to {@code
bound}. */
+ private static void readPayload(GcsUtilV2 gcsUtil, @Nullable
MetricsContainerImpl bound)
+ throws IOException {
+ try (SeekableByteChannel channel =
+ gcsUtil.wrapInCounting(new SeekableInMemoryByteChannel(PAYLOAD),
BUCKET, bound)) {
+ assertEquals(PAYLOAD.length,
channel.read(ByteBuffer.allocate(PAYLOAD.length)));
+ }
+ }
+
+ /** Writes {@link #PAYLOAD} through V2's write wrapper, bound to {@code
bound}. */
+ private static void writePayload(GcsUtilV2 gcsUtil, @Nullable
MetricsContainerImpl bound)
+ throws IOException {
+ try (WritableByteChannel channel =
+ gcsUtil.wrapInCounting(Channels.newChannel(new
ByteArrayOutputStream()), BUCKET, bound)) {
+ assertEquals(PAYLOAD.length, channel.write(ByteBuffer.wrap(PAYLOAD)));
+ }
+ }
+
+ @Test
+ public void testReadCountersWhenAllEnabled() throws IOException {
+ MetricsContainerImpl container = new MetricsContainerImpl(null);
+ MetricsEnvironment.setCurrentContainer(container);
+
+ readPayload(gcsUtilV2(true, true), container);
+
+ assertEquals(Long.valueOf(PAYLOAD.length), bucketCounter(container,
READ_PREFIX));
+ assertEquals(
+ Long.valueOf(PAYLOAD.length), gcsCounter(container,
"gcs_http_read_wire_bytes_received"));
+ assertNull(bucketCounter(container, WRITE_PREFIX));
+ assertNull(gcsCounter(container, "gcs_http_write_wire_bytes_sent"));
+ }
+
+ @Test
+ public void testWriteCountersWhenAllEnabled() throws IOException {
+ MetricsContainerImpl container = new MetricsContainerImpl(null);
+ MetricsEnvironment.setCurrentContainer(container);
+
+ writePayload(gcsUtilV2(true, true), container);
+
+ assertEquals(Long.valueOf(PAYLOAD.length), bucketCounter(container,
WRITE_PREFIX));
+ assertEquals(
+ Long.valueOf(PAYLOAD.length), gcsCounter(container,
"gcs_http_write_wire_bytes_sent"));
+ assertNull(bucketCounter(container, READ_PREFIX));
+ assertNull(gcsCounter(container, "gcs_http_read_wire_bytes_received"));
+ }
+
+ @Test
+ public void testChannelsAreNotWrappedWhenAllDisabled() throws IOException {
+ MetricsContainerImpl container = new MetricsContainerImpl(null);
+ MetricsEnvironment.setCurrentContainer(container);
+ GcsUtilV2 gcsUtil = gcsUtilV2(false, false);
+ SeekableByteChannel readChannel = new SeekableInMemoryByteChannel(PAYLOAD);
+ WritableByteChannel writeChannel = Channels.newChannel(new
ByteArrayOutputStream());
+
+ assertSame(readChannel, gcsUtil.wrapInCounting(readChannel, BUCKET,
container));
+ assertSame(writeChannel, gcsUtil.wrapInCounting(writeChannel, BUCKET,
container));
+ }
+
+ @Test
+ public void testOnlyBucketCountersWhenPerformanceMetricsDisabled() throws
IOException {
+ MetricsContainerImpl container = new MetricsContainerImpl(null);
+ MetricsEnvironment.setCurrentContainer(container);
+ GcsUtilV2 gcsUtil = gcsUtilV2(false, true);
+
+ readPayload(gcsUtil, container);
+ writePayload(gcsUtil, container);
+
+ assertEquals(Long.valueOf(PAYLOAD.length), bucketCounter(container,
READ_PREFIX));
+ assertEquals(Long.valueOf(PAYLOAD.length), bucketCounter(container,
WRITE_PREFIX));
+ assertNull(gcsCounter(container, "gcs_http_read_wire_bytes_received"));
+ assertNull(gcsCounter(container, "gcs_http_write_wire_bytes_sent"));
+ }
+
+ @Test
+ public void testOnlyWireBytesWhenBucketCountersDisabled() throws IOException
{
+ MetricsContainerImpl container = new MetricsContainerImpl(null);
+ MetricsEnvironment.setCurrentContainer(container);
+ GcsUtilV2 gcsUtil = gcsUtilV2(true, false);
+
+ readPayload(gcsUtil, container);
+ writePayload(gcsUtil, container);
+
+ assertNull(bucketCounter(container, READ_PREFIX));
+ assertNull(bucketCounter(container, WRITE_PREFIX));
+ assertEquals(
+ Long.valueOf(PAYLOAD.length), gcsCounter(container,
"gcs_http_read_wire_bytes_received"));
+ assertEquals(
+ Long.valueOf(PAYLOAD.length), gcsCounter(container,
"gcs_http_write_wire_bytes_sent"));
+ }
+
+ /** Without a bound container there is nowhere to report wire bytes, so
nothing is wrapped. */
+ @Test
+ public void testWireBytesRequireABoundContainer() {
+ GcsUtilV2 gcsUtil = gcsUtilV2(true, false);
+ SeekableByteChannel readChannel = new SeekableInMemoryByteChannel(PAYLOAD);
+ WritableByteChannel writeChannel = Channels.newChannel(new
ByteArrayOutputStream());
+
+ assertSame(readChannel, gcsUtil.wrapInCounting(readChannel, BUCKET, null));
+ assertSame(writeChannel, gcsUtil.wrapInCounting(writeChannel, BUCKET,
null));
+ }
+
+ /**
+ * Wire bytes go to the container bound when the channel was created, even
if a different
+ * container is current while the bytes are read or written.
+ */
+ @Test
+ public void testWireBytesAreAttributedToTheBoundContainer() throws
IOException {
+ MetricsContainerImpl bound = new MetricsContainerImpl("bound");
+ MetricsContainerImpl other = new MetricsContainerImpl("other");
+ GcsUtilV2 gcsUtil = gcsUtilV2(true, false);
+ SeekableByteChannel readChannel =
+ gcsUtil.wrapInCounting(new SeekableInMemoryByteChannel(PAYLOAD),
BUCKET, bound);
+ WritableByteChannel writeChannel =
+ gcsUtil.wrapInCounting(Channels.newChannel(new
ByteArrayOutputStream()), BUCKET, bound);
+
+ MetricsEnvironment.setCurrentContainer(other);
+ readChannel.read(ByteBuffer.allocate(PAYLOAD.length));
+ writeChannel.write(ByteBuffer.wrap(PAYLOAD));
+ readChannel.close();
+ writeChannel.close();
+
+ assertEquals(
+ Long.valueOf(PAYLOAD.length), gcsCounter(bound,
"gcs_http_read_wire_bytes_received"));
+ assertEquals(Long.valueOf(PAYLOAD.length), gcsCounter(bound,
"gcs_http_write_wire_bytes_sent"));
+ assertNull(gcsCounter(other, "gcs_http_read_wire_bytes_received"));
+ assertNull(gcsCounter(other, "gcs_http_write_wire_bytes_sent"));
+ }
+}