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 28279d2e9f5 Extend GCS performance metrics and unify the GCS metric 
namespace (#40142)
28279d2e9f5 is described below

commit 28279d2e9f59268696ab62267c26ea4403d768f9
Author: Shunping Huang <[email protected]>
AuthorDate: Tue Sep 22 12:19:45 2026 -0400

    Extend GCS performance metrics and unify the GCS metric namespace (#40142)
    
    * Add more gcs performance metrics in gcsutilv1
    
    * Extend GCS operation metrics and unify GCS metric name space.
    
    * Move the GCS performance metrics flag into GcsCountersOptions
    
    * Fix the inaccurate javadocs
    
    * Cache the storage instance for different metric containers when gcs 
metrics are enabled
    
    * Use weak-keyed Cache with removal listener for per-container 
GoogleCloudStorage
---
 .../sdk/extensions/gcp/storage/GcsFileSystem.java  | 157 ++++++++++-----
 .../beam/sdk/extensions/gcp/util/GcsUtil.java      |  14 ++
 .../beam/sdk/extensions/gcp/util/GcsUtilV1.java    | 219 +++++++++++++++------
 .../beam/sdk/extensions/gcp/util/Transport.java    | 165 +++++++++++++++-
 .../extensions/gcp/storage/GcsFileSystemTest.java  |  61 ++++++
 .../beam/sdk/extensions/gcp/util/GcsUtilTest.java  |   8 +-
 .../sdk/extensions/gcp/util/TransportTest.java     | 106 ++++++++++
 .../java/org/apache/beam/sdk/io/text/TextIOIT.java |   8 +-
 8 files changed, 620 insertions(+), 118 deletions(-)

diff --git 
a/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/storage/GcsFileSystem.java
 
b/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/storage/GcsFileSystem.java
index 1bee44eb38c..5eca4e9e2c2 100644
--- 
a/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/storage/GcsFileSystem.java
+++ 
b/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/storage/GcsFileSystem.java
@@ -77,25 +77,58 @@ class GcsFileSystem extends FileSystem<GcsResourceId> {
 
   private final GcsOptions options;
 
-  /** Number of copy operations performed. */
-  private Counter numCopies;
+  /** The {@code _count} and {@code _msec} counter pair for a single 
operation. */
+  private static class OpMetrics {
+    private final Counter count;
+    private final Counter msec;
+
+    OpMetrics(String operation) {
+      this.count = Metrics.counter(GcsUtil.METRIC_NAMESPACE, "gcs_op_" + 
operation + "_count");
+      this.msec = Metrics.counter(GcsUtil.METRIC_NAMESPACE, "gcs_op_" + 
operation + "_msec");
+    }
+  }
 
-  /** Number of renames operations performed. */
-  private Counter numRenames;
+  /**
+   * Per-operation metrics, or null when {@link 
GcsOptions#getGcsPerformanceMetrics()} is off, which
+   * is the default. Every recording site tolerates null, so nothing is 
emitted unless asked for.
+   *
+   * <p>For {@link #open} and {@link #create} the elapsed time covers only 
channel setup, not the
+   * transfer; the bytes moved are counted by the {@code gcs_http_*} wire-byte 
counters.
+   */
+  private @Nullable OpMetrics copyMetrics;
 
-  /** Time spent performing copies. */
-  private Counter copyTimeMsec;
+  private @Nullable OpMetrics renameMetrics;
+  private @Nullable OpMetrics deleteMetrics;
+  private @Nullable OpMetrics matchGlobMetrics;
+  private @Nullable OpMetrics matchNonGlobMetrics;
+  private @Nullable OpMetrics openMetrics;
+  private @Nullable OpMetrics createMetrics;
 
-  /** Time spent performing renames. */
-  private Counter renameTimeMsec;
+  /** Object listing pages fetched while expanding globs. */
+  private @Nullable Counter matchGlobPages;
 
   GcsFileSystem(GcsOptions options) {
     this.options = checkNotNull(options, "options");
     if (options.getGcsPerformanceMetrics()) {
-      numCopies = Metrics.counter(GcsFileSystem.class, "num_copies");
-      copyTimeMsec = Metrics.counter(GcsFileSystem.class, "copy_time_msec");
-      numRenames = Metrics.counter(GcsFileSystem.class, "num_renames");
-      renameTimeMsec = Metrics.counter(GcsFileSystem.class, 
"rename_time_msec");
+      copyMetrics = new OpMetrics("copy");
+      renameMetrics = new OpMetrics("rename");
+      deleteMetrics = new OpMetrics("delete");
+      matchGlobMetrics = new OpMetrics("match_glob");
+      matchNonGlobMetrics = new OpMetrics("match_nonglob");
+      openMetrics = new OpMetrics("open");
+      createMetrics = new OpMetrics("create");
+      matchGlobPages = Metrics.counter(GcsUtil.METRIC_NAMESPACE, 
"gcs_op_match_glob_pages");
+    }
+  }
+
+  /**
+   * Records a finished operation. Called from a finally block so that failed 
operations, which are
+   * often the slow ones, are counted too.
+   */
+  private static void record(@Nullable OpMetrics metrics, long count, 
Stopwatch stopwatch) {
+    if (metrics != null) {
+      metrics.count.inc(count);
+      metrics.msec.inc(stopwatch.elapsed(TimeUnit.MILLISECONDS));
     }
   }
 
@@ -152,12 +185,22 @@ class GcsFileSystem extends FileSystem<GcsResourceId> {
           builder.setUploadBufferSizeBytes(
               ((GcsCreateOptions) createOptions).gcsUploadBufferSizeBytes());
     }
-    return options.getGcsUtil().create(resourceId.getGcsPath(), 
builder.build());
+    Stopwatch stopwatch = Stopwatch.createStarted();
+    try {
+      return options.getGcsUtil().create(resourceId.getGcsPath(), 
builder.build());
+    } finally {
+      record(createMetrics, 1, stopwatch);
+    }
   }
 
   @Override
   protected ReadableByteChannel open(GcsResourceId resourceId) throws 
IOException {
-    return options.getGcsUtil().open(resourceId.getGcsPath());
+    Stopwatch stopwatch = Stopwatch.createStarted();
+    try {
+      return options.getGcsUtil().open(resourceId.getGcsPath());
+    } finally {
+      record(openMetrics, 1, stopwatch);
+    }
   }
 
   @Override
@@ -167,19 +210,23 @@ class GcsFileSystem extends FileSystem<GcsResourceId> {
       MoveOptions... moveOptions)
       throws IOException {
     Stopwatch stopwatch = Stopwatch.createStarted();
-    options
-        .getGcsUtil()
-        .rename(toFilenames(srcResourceIds), toFilenames(destResourceIds), 
moveOptions);
-    stopwatch.stop();
-    if (options.getGcsPerformanceMetrics()) {
-      numRenames.inc(srcResourceIds.size());
-      renameTimeMsec.inc(stopwatch.elapsed(TimeUnit.MILLISECONDS));
+    try {
+      options
+          .getGcsUtil()
+          .rename(toFilenames(srcResourceIds), toFilenames(destResourceIds), 
moveOptions);
+    } finally {
+      record(renameMetrics, srcResourceIds.size(), stopwatch);
     }
   }
 
   @Override
   protected void delete(Collection<GcsResourceId> resourceIds) throws 
IOException {
-    options.getGcsUtil().remove(toFilenames(resourceIds));
+    Stopwatch stopwatch = Stopwatch.createStarted();
+    try {
+      options.getGcsUtil().remove(toFilenames(resourceIds));
+    } finally {
+      record(deleteMetrics, resourceIds.size(), stopwatch);
+    }
   }
 
   @Override
@@ -202,11 +249,10 @@ class GcsFileSystem extends FileSystem<GcsResourceId> {
   protected void copy(List<GcsResourceId> srcResourceIds, List<GcsResourceId> 
destResourceIds)
       throws IOException {
     Stopwatch stopwatch = Stopwatch.createStarted();
-    options.getGcsUtil().copy(toFilenames(srcResourceIds), 
toFilenames(destResourceIds));
-    stopwatch.stop();
-    if (options.getGcsPerformanceMetrics()) {
-      numCopies.inc(srcResourceIds.size());
-      copyTimeMsec.inc(stopwatch.elapsed(TimeUnit.MILLISECONDS));
+    try {
+      options.getGcsUtil().copy(toFilenames(srcResourceIds), 
toFilenames(destResourceIds));
+    } finally {
+      record(copyMetrics, srcResourceIds.size(), stopwatch);
     }
   }
 
@@ -265,26 +311,35 @@ class GcsFileSystem extends FileSystem<GcsResourceId> {
         prefix,
         p.toString());
 
-    String pageToken = null;
-    List<Metadata> results = new ArrayList<>();
-    do {
-      Objects objects = 
options.getGcsUtil().listObjects(gcsPattern.getBucket(), prefix, pageToken);
-      if (objects.getItems() == null) {
-        break;
-      }
+    Stopwatch stopwatch = Stopwatch.createStarted();
+    try {
+      String pageToken = null;
+      List<Metadata> results = new ArrayList<>();
+      do {
+        Objects objects =
+            options.getGcsUtil().listObjects(gcsPattern.getBucket(), prefix, 
pageToken);
+        if (matchGlobPages != null) {
+          matchGlobPages.inc();
+        }
+        if (objects.getItems() == null) {
+          break;
+        }
 
-      // Filter objects based on the regex.
-      for (StorageObject o : objects.getItems()) {
-        String name = o.getName();
-        // Skip directories, which end with a slash.
-        if (p.matcher(name).matches() && !name.endsWith("/")) {
-          LOG.debug("Matched object: {}", name);
-          results.add(toMetadata(o));
+        // Filter objects based on the regex.
+        for (StorageObject o : objects.getItems()) {
+          String name = o.getName();
+          // Skip directories, which end with a slash.
+          if (p.matcher(name).matches() && !name.endsWith("/")) {
+            LOG.debug("Matched object: {}", name);
+            results.add(toMetadata(o));
+          }
         }
-      }
-      pageToken = objects.getNextPageToken();
-    } while (pageToken != null);
-    return MatchResult.create(Status.OK, results);
+        pageToken = objects.getNextPageToken();
+      } while (pageToken != null);
+      return MatchResult.create(Status.OK, results);
+    } finally {
+      record(matchGlobMetrics, 1, stopwatch);
+    }
   }
 
   /**
@@ -295,7 +350,17 @@ class GcsFileSystem extends FileSystem<GcsResourceId> {
    */
   @VisibleForTesting
   List<MatchResult> matchNonGlobs(List<GcsPath> gcsPaths) throws IOException {
-    List<StorageObjectOrIOException> results = 
options.getGcsUtil().getObjects(gcsPaths);
+    if (gcsPaths.isEmpty()) {
+      // match() always calls this, so recording here would bury the real 
calls in empty ones.
+      return ImmutableList.of();
+    }
+    Stopwatch stopwatch = Stopwatch.createStarted();
+    List<StorageObjectOrIOException> results;
+    try {
+      results = options.getGcsUtil().getObjects(gcsPaths);
+    } finally {
+      record(matchNonGlobMetrics, gcsPaths.size(), stopwatch);
+    }
 
     ImmutableList.Builder<MatchResult> ret = ImmutableList.builder();
     for (StorageObjectOrIOException result : results) {
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 5ed97d935c6..070cf74d7c1 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
@@ -49,9 +49,23 @@ import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Sets;
 import org.checkerframework.checker.nullness.qual.Nullable;
 
 public class GcsUtil {
+  /**
+   * 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
+   * transport-level counters and {@code gcs_op_*} for operation-level ones.
+   */
+  public static final String METRIC_NAMESPACE = "Gcs";
+
   @VisibleForTesting GcsUtilV1 delegate;
   @VisibleForTesting @Nullable GcsUtilV2 delegateV2;
 
+  /**
+   * @deprecated no {@link GcsUtil} API accepts this type, so an instance 
cannot be used for
+   *     anything. GCS counters are configured from {@link
+   *     org.apache.beam.sdk.extensions.gcp.options.GcsOptions} when the 
{@link GcsUtil} is
+   *     constructed. Scheduled for removal.
+   */
+  @Deprecated
   public static class GcsCountersOptions {
     final GcsUtilV1.GcsCountersOptions delegate;
 
diff --git 
a/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV1.java
 
b/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV1.java
index 5a6a87dc4b7..a04f688ec7a 100644
--- 
a/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV1.java
+++ 
b/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV1.java
@@ -50,6 +50,7 @@ import com.google.cloud.hadoop.util.AsyncWriteChannelOptions;
 import com.google.cloud.hadoop.util.ResilientOperation;
 import com.google.cloud.hadoop.util.RetryDeterminer;
 import edu.umd.cs.findbugs.annotations.SuppressFBWarnings;
+import java.io.Closeable;
 import java.io.FileNotFoundException;
 import java.io.IOException;
 import java.nio.channels.SeekableByteChannel;
@@ -85,13 +86,20 @@ import 
org.apache.beam.sdk.extensions.gcp.util.channels.CountingWritableByteChan
 import org.apache.beam.sdk.extensions.gcp.util.gcsfs.GcsPath;
 import org.apache.beam.sdk.io.fs.MoveOptions;
 import org.apache.beam.sdk.io.fs.MoveOptions.StandardMoveOptions;
+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.sdk.util.FluentBackoff;
 import org.apache.beam.sdk.util.MoreFutures;
 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.base.Preconditions;
+import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.cache.Cache;
+import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.cache.CacheBuilder;
+import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.cache.RemovalNotification;
 import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList;
 import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Lists;
 import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Sets;
@@ -107,19 +115,35 @@ import org.slf4j.LoggerFactory;
 })
 class GcsUtilV1 {
 
+  /** Describes which GCS counters this {@link GcsUtilV1} emits. */
   @AutoValue
   public abstract static class GcsCountersOptions {
     public abstract @Nullable String getReadCounterPrefix();
 
     public abstract @Nullable String getWriteCounterPrefix();
 
+    /**
+     * Whether to emit the {@code gcs_*} performance counters, which are 
reported under {@link
+     * GcsUtil#METRIC_NAMESPACE} and are not per bucket. Set from {@link
+     * GcsOptions#getGcsPerformanceMetrics()}.
+     */
+    public abstract boolean getPerformanceMetricsEnabled();
+
     public boolean hasAnyPrefix() {
       return getWriteCounterPrefix() != null || getReadCounterPrefix() != null;
     }
 
     public static GcsCountersOptions create(
         @Nullable String readCounterPrefix, @Nullable String 
writeCounterPrefix) {
-      return new AutoValue_GcsUtilV1_GcsCountersOptions(readCounterPrefix, 
writeCounterPrefix);
+      return create(readCounterPrefix, writeCounterPrefix, false);
+    }
+
+    public static GcsCountersOptions create(
+        @Nullable String readCounterPrefix,
+        @Nullable String writeCounterPrefix,
+        boolean performanceMetricsEnabled) {
+      return new AutoValue_GcsUtilV1_GcsCountersOptions(
+          readCounterPrefix, writeCounterPrefix, performanceMetricsEnabled);
     }
   }
 
@@ -153,7 +177,8 @@ class GcsUtilV1 {
                   : null,
               gcsOptions.getEnableBucketWriteMetricCounter()
                   ? gcsOptions.getGcsWriteCounterPrefix()
-                  : null),
+                  : null,
+              Boolean.TRUE.equals(gcsOptions.getGcsPerformanceMetrics())),
           gcsOptions.getGoogleCloudStorageReadOptions());
     }
   }
@@ -213,6 +238,17 @@ class GcsUtilV1 {
 
   private GoogleCloudStorage googleCloudStorage;
   private GoogleCloudStorageOptions googleCloudStorageOptions;
+  private final Cache<MetricsContainer, GoogleCloudStorage> 
readStorageByContainer =
+      CacheBuilder.newBuilder()
+          .weakKeys()
+          .removalListener(
+              (RemovalNotification<MetricsContainer, GoogleCloudStorage> 
notification) -> {
+                GoogleCloudStorage storage = notification.getValue();
+                if (storage != null) {
+                  storage.close();
+                }
+              })
+          .build();
 
   private final int rewriteDataOpBatchLimit;
 
@@ -223,29 +259,6 @@ class GcsUtilV1 {
 
   @VisibleForTesting @Nullable AtomicInteger numRewriteTokensUsed;
 
-  @VisibleForTesting
-  GcsUtilV1(
-      Storage storageClient,
-      HttpRequestInitializer httpRequestInitializer,
-      ExecutorService executorService,
-      Boolean shouldUseGrpc,
-      Credentials credentials,
-      @Nullable Integer uploadBufferSizeBytes,
-      @Nullable Integer rewriteDataOpBatchLimit,
-      GcsCountersOptions gcsCountersOptions,
-      GcsOptions gcsOptions) {
-    this(
-        storageClient,
-        httpRequestInitializer,
-        executorService,
-        shouldUseGrpc,
-        credentials,
-        uploadBufferSizeBytes,
-        rewriteDataOpBatchLimit,
-        gcsCountersOptions,
-        gcsOptions.getGoogleCloudStorageReadOptions());
-  }
-
   @VisibleForTesting
   GcsUtilV1(
       Storage storageClient,
@@ -526,54 +539,79 @@ class GcsUtilV1 {
   }
 
   private WritableByteChannel wrapInCounting(
-      WritableByteChannel writableByteChannel, String bucket) {
+      WritableByteChannel writableByteChannel,
+      String bucket,
+      @Nullable MetricsContainer container) {
     if (writableByteChannel instanceof CountingWritableByteChannel) {
       return writableByteChannel;
     }
-    return Optional.ofNullable(gcsCountersOptions.getWriteCounterPrefix())
-        .<WritableByteChannel>map(
-            prefix -> {
-              LOG.debug(
-                  "wrapping writable byte channel using counter name prefix {} 
and bucket {}",
-                  prefix,
-                  bucket);
-              return new CountingWritableByteChannel(
-                  writableByteChannel, createCounterConsumer(prefix, bucket));
-            })
-        .orElse(writableByteChannel);
-  }
-
-  private SeekableByteChannel wrapInCounting(
-      SeekableByteChannel seekableByteChannel, String bucket) {
-    if (seekableByteChannel instanceof CountingSeekableByteChannel
-        || !gcsCountersOptions.hasAnyPrefix()) {
-      return seekableByteChannel;
-    }
 
-    return new CountingSeekableByteChannel(
-        seekableByteChannel,
-        Optional.ofNullable(gcsCountersOptions.getReadCounterPrefix())
+    Consumer<Integer> writeConsumer =
+        Optional.ofNullable(gcsCountersOptions.getWriteCounterPrefix())
             .map(
                 prefix -> {
                   LOG.debug(
-                      "wrapping seekable byte channel with \"bytes read\" 
counter name prefix {}"
-                          + " and bucket {}",
+                      "wrapping writable byte channel using counter name 
prefix {} and bucket {}",
                       prefix,
                       bucket);
                   return createCounterConsumer(prefix, bucket);
                 })
-            .orElse(null),
-        Optional.ofNullable(gcsCountersOptions.getWriteCounterPrefix())
+            .orElse(null);
+
+    if (gcsCountersOptions.getPerformanceMetricsEnabled() && 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);
+  }
+
+  private SeekableByteChannel wrapInCounting(
+      SeekableByteChannel seekableByteChannel,
+      String bucket,
+      @Nullable MetricsContainer container) {
+    if (seekableByteChannel instanceof CountingSeekableByteChannel) {
+      return seekableByteChannel;
+    }
+
+    // SeekableByteChannel is only returned by GcsUtilV1.open(...) for reading 
immutable GCS objects
+    // (GoogleCloudStorageReadChannel throws NonWritableChannelException on 
write). All GCS writes
+    // go through GcsUtilV1.create(...), which returns a WritableByteChannel. 
Therefore, only a read
+    // counter consumer is needed here.
+    Consumer<Integer> readConsumer =
+        Optional.ofNullable(gcsCountersOptions.getReadCounterPrefix())
             .map(
                 prefix -> {
                   LOG.debug(
-                      "wrapping seekable byte channel with \"bytes written\" 
counter name prefix {}"
+                      "wrapping seekable byte channel with \"bytes read\" 
counter name prefix {}"
                           + " and bucket {}",
                       prefix,
                       bucket);
                   return createCounterConsumer(prefix, bucket);
                 })
-            .orElse(null));
+            .orElse(null);
+
+    if (gcsCountersOptions.getPerformanceMetricsEnabled() && 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);
   }
 
   /**
@@ -615,11 +653,38 @@ class GcsUtilV1 {
     ServiceCallMetric serviceCallMetric =
         new ServiceCallMetric(MonitoringInfoConstants.Urns.API_REQUEST_COUNT, 
baseLabels);
     try {
+      GoogleCloudStorage gcpStorage = this.googleCloudStorage;
+      MetricsContainer container = null;
+      if (gcsCountersOptions.getPerformanceMetricsEnabled()) {
+        container = MetricsEnvironment.getCurrentContainer();
+        if (container != null) {
+          final MetricsContainer currentContainer = container;
+          try {
+            gcpStorage =
+                readStorageByContainer.get(
+                    currentContainer,
+                    () -> {
+                      HttpRequestInitializer scopedInitializer =
+                          Transport.withMetricsContainer(
+                              this.httpRequestInitializer, currentContainer, 
false);
+                      return createGoogleCloudStorage(
+                          googleCloudStorageOptions,
+                          this.storageClient,
+                          this.credentials,
+                          scopedInitializer);
+                    });
+          } catch (ExecutionException e) {
+            if (e.getCause() instanceof IOException) {
+              throw (IOException) e.getCause();
+            }
+            throw new IOException(e);
+          }
+        }
+      }
       SeekableByteChannel channel =
-          googleCloudStorage.open(
-              new StorageResourceId(path.getBucket(), path.getObject()), 
readOptions);
+          gcpStorage.open(new StorageResourceId(path.getBucket(), 
path.getObject()), readOptions);
       serviceCallMetric.call("ok");
-      return wrapInCounting(channel, path.getBucket());
+      return wrapInCounting(channel, path.getBucket(), container);
     } catch (IOException e) {
       if (e.getCause() instanceof GoogleJsonResponseException) {
         serviceCallMetric.call(((GoogleJsonResponseException) 
e.getCause()).getDetails().getCode());
@@ -701,9 +766,18 @@ class GcsUtilV1 {
     }
     GoogleCloudStorageOptions newGoogleCloudStorageOptions =
         
googleCloudStorageOptions.toBuilder().setWriteChannelOptions(wcOptions).build();
+    HttpRequestInitializer scopedInitializer = this.httpRequestInitializer;
+    MetricsContainer container = null;
+    if (gcsCountersOptions.getPerformanceMetricsEnabled()) {
+      container = MetricsEnvironment.getCurrentContainer();
+      if (container != null) {
+        scopedInitializer =
+            Transport.withMetricsContainer(this.httpRequestInitializer, 
container, true);
+      }
+    }
     GoogleCloudStorage gcpStorage =
         createGoogleCloudStorage(
-            newGoogleCloudStorageOptions, this.storageClient, 
this.credentials);
+            newGoogleCloudStorageOptions, this.storageClient, 
this.credentials, scopedInitializer);
     StorageResourceId resourceId =
         new StorageResourceId(
             path.getBucket(),
@@ -735,7 +809,7 @@ class GcsUtilV1 {
     try {
       WritableByteChannel channel = gcpStorage.create(resourceId, 
createBuilder.build());
       serviceCallMetric.call("ok");
-      return wrapInCounting(channel, path.getBucket());
+      return wrapInCounting(channel, path.getBucket(), container);
     } catch (IOException e) {
       if (e.getCause() instanceof GoogleJsonResponseException) {
         serviceCallMetric.call(((GoogleJsonResponseException) 
e.getCause()).getDetails().getCode());
@@ -744,10 +818,21 @@ class GcsUtilV1 {
     }
   }
 
-  @SuppressFBWarnings("LG_LOST_LOGGER_DUE_TO_WEAK_REFERENCE")
+  @VisibleForTesting
   GoogleCloudStorage createGoogleCloudStorage(
       GoogleCloudStorageOptions options, Storage storage, Credentials 
credentials)
       throws IOException {
+    return createGoogleCloudStorage(options, storage, credentials, 
this.httpRequestInitializer);
+  }
+
+  @VisibleForTesting
+  @SuppressFBWarnings("LG_LOST_LOGGER_DUE_TO_WEAK_REFERENCE")
+  GoogleCloudStorage createGoogleCloudStorage(
+      GoogleCloudStorageOptions options,
+      Storage storage,
+      Credentials credentials,
+      @Nullable HttpRequestInitializer httpRequestInitializer)
+      throws IOException {
     // Suppress log spams in gcsio 3.0
     if (overwriteLog.compareAndSet(false, true)) {
       
java.util.logging.Logger.getLogger("com.google.cloud.hadoop.gcsio.GoogleCloudStorageImpl")
@@ -949,9 +1034,21 @@ class GcsUtilV1 {
                 TimeUnit.MILLISECONDS,
                 new LinkedBlockingQueue<>()));
 
+    MetricsContainer container = MetricsEnvironment.getCurrentContainer();
     List<CompletionStage<Void>> futures = new ArrayList<>();
     for (final BatchInterface batch : batches) {
-      futures.add(MoreFutures.runAsync(batch::execute, executor));
+      futures.add(
+          MoreFutures.runAsync(
+              () -> {
+                if (container != null) {
+                  try (Closeable scope = 
MetricsEnvironment.scopedMetricsContainer(container)) {
+                    batch.execute();
+                  }
+                } else {
+                  batch.execute();
+                }
+              },
+              executor));
     }
 
     try {
diff --git 
a/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/Transport.java
 
b/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/Transport.java
index ea31e6c9180..6e100a6148d 100644
--- 
a/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/Transport.java
+++ 
b/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/Transport.java
@@ -21,7 +21,12 @@ import static 
org.apache.beam.sdk.extensions.gcp.options.GcsOptions.GcsCustomAud
 import static 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Strings.isNullOrEmpty;
 
 import com.google.api.client.googleapis.javanet.GoogleNetHttpTransport;
+import com.google.api.client.http.HttpExecuteInterceptor;
+import com.google.api.client.http.HttpIOExceptionHandler;
+import com.google.api.client.http.HttpRequest;
 import com.google.api.client.http.HttpRequestInitializer;
+import com.google.api.client.http.HttpResponse;
+import com.google.api.client.http.HttpResponseInterceptor;
 import com.google.api.client.http.HttpTransport;
 import com.google.api.client.json.JsonFactory;
 import com.google.api.client.json.gson.GsonFactory;
@@ -39,6 +44,9 @@ import java.util.Optional;
 import javax.annotation.Nullable;
 import org.apache.beam.sdk.extensions.gcp.auth.NullCredentialInitializer;
 import org.apache.beam.sdk.extensions.gcp.options.GcsOptions;
+import org.apache.beam.sdk.metrics.Counter;
+import org.apache.beam.sdk.metrics.MetricName;
+import org.apache.beam.sdk.metrics.MetricsContainer;
 import org.apache.beam.sdk.util.ReleaseInfo;
 import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList;
 import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap;
@@ -119,6 +127,106 @@ public class Transport {
     return storageBuilder;
   }
 
+  /**
+   * Wraps an {@link HttpRequestInitializer} so that HTTP execute and response 
interceptors
+   * increment {@link Counter} instances pre-bound to the given {@link 
MetricsContainer}. This
+   * guarantees that GCS HTTP metrics are attributed directly to the step that 
created the channel,
+   * even when requests execute on background worker threads.
+   *
+   * <ul>
+   *   <li>{@code request_count} counts every attempt, retries included, 
because the request
+   *       interceptor runs once per attempt.
+   *   <li>{@code status_2xx}, {@code status_3xx}, {@code status_4xx}, {@code 
status_5xx}, {@code
+   *       status_other} (1xx) and {@code request_no_response} add up to 
{@code request_count} when
+   *       no retries occur. If their sum is smaller than {@code 
request_count}, it indicates that
+   *       retries happened, as only the final response is recorded while 
multiple requests are
+   *       counted.
+   *   <li>For reads, every attempt is also classified by shape into {@code 
request_count_ranged} (a
+   *       GET with a Range header), {@code request_count_unbounded} (a GET 
without one) or {@code
+   *       request_count_other} (anything that is not a GET, e.g. a batched 
metadata POST). Their
+   *       sum equals {@code request_count}. Writes are not classified this 
way, as they are POSTs
+   *       and PUTs by construction.
+   * </ul>
+   *
+   * <p>Note that {@code request_count_unbounded} counts metadata GETs as well 
as full object reads,
+   * since neither carries a Range header.
+   */
+  public static HttpRequestInitializer withMetricsContainer(
+      HttpRequestInitializer base, @Nullable MetricsContainer container, 
boolean isWrite) {
+    if (container == null) {
+      return base;
+    }
+
+    String prefix = isWrite ? "gcs_http_write_" : "gcs_http_read_";
+
+    // Pre-resolve counters on the calling thread (e.g. DoFn thread) while it 
is in the target step
+    Counter requestCount =
+        container.getCounter(MetricName.named(GcsUtil.METRIC_NAMESPACE, prefix 
+ "request_count"));
+    Counter rangeRequestCount =
+        isWrite
+            ? null
+            : container.getCounter(
+                MetricName.named(GcsUtil.METRIC_NAMESPACE, prefix + 
"request_count_ranged"));
+    Counter unboundedStreamCount =
+        isWrite
+            ? null
+            : container.getCounter(
+                MetricName.named(GcsUtil.METRIC_NAMESPACE, prefix + 
"request_count_unbounded"));
+    // Requests that are not a GET, so that the shape counters above add up to 
request_count.
+    Counter otherRequestCount =
+        isWrite
+            ? null
+            : container.getCounter(
+                MetricName.named(GcsUtil.METRIC_NAMESPACE, prefix + 
"request_count_other"));
+    Counter status2xx =
+        container.getCounter(MetricName.named(GcsUtil.METRIC_NAMESPACE, prefix 
+ "status_2xx"));
+    // 3xx is not an error for GCS: a resumable upload answers 308 to every 
chunk but the last.
+    Counter status3xx =
+        container.getCounter(MetricName.named(GcsUtil.METRIC_NAMESPACE, prefix 
+ "status_3xx"));
+    Counter status4xx =
+        container.getCounter(MetricName.named(GcsUtil.METRIC_NAMESPACE, prefix 
+ "status_4xx"));
+    Counter status5xx =
+        container.getCounter(MetricName.named(GcsUtil.METRIC_NAMESPACE, prefix 
+ "status_5xx"));
+    Counter statusOther =
+        container.getCounter(MetricName.named(GcsUtil.METRIC_NAMESPACE, prefix 
+ "status_other"));
+    Counter noResponse =
+        container.getCounter(
+            MetricName.named(GcsUtil.METRIC_NAMESPACE, prefix + 
"request_no_response"));
+
+    return request -> {
+      base.initialize(request);
+      HttpExecuteInterceptor existingExecuteInterceptor = 
request.getInterceptor();
+      request.setInterceptor(
+          req -> {
+            if (existingExecuteInterceptor != null) {
+              existingExecuteInterceptor.intercept(req);
+            }
+            recordRequestMetrics(
+                req, requestCount, rangeRequestCount, unboundedStreamCount, 
otherRequestCount);
+          });
+
+      HttpResponseInterceptor existingResponseInterceptor = 
request.getResponseInterceptor();
+      request.setResponseInterceptor(
+          res -> {
+            if (existingResponseInterceptor != null) {
+              existingResponseInterceptor.interceptResponse(res);
+            }
+            recordResponseMetrics(res, status2xx, status3xx, status4xx, 
status5xx, statusOther);
+          });
+
+      // An attempt that throws before a response is received never reaches 
the response
+      // interceptor, so it is counted here instead. The existing handler 
decides whether the
+      // request is retried, this only observes it.
+      HttpIOExceptionHandler existingIOExceptionHandler = 
request.getIOExceptionHandler();
+      request.setIOExceptionHandler(
+          (req, supportsRetry) -> {
+            noResponse.inc();
+            return existingIOExceptionHandler != null
+                && existingIOExceptionHandler.handleIOException(req, 
supportsRetry);
+          });
+    };
+  }
+
   private static HttpRequestInitializer 
httpRequestInitializerFromOptions(GcsOptions options) {
     // Do not log the code 404. Code up the stack will deal with 404's if 
needed,
     // and logging it by default clutters the output during file staging.
@@ -148,12 +256,59 @@ public class Transport {
       retryHttpRequestInitializer.setWriteTimeout(writeTimeout);
     }
     Credentials credential = options.getGcpCredential();
-    if (credential == null) {
-      return new ChainingHttpRequestInitializer(
-          new NullCredentialInitializer(), retryHttpRequestInitializer);
+    HttpRequestInitializer credentialsInitializer =
+        credential == null
+            ? new NullCredentialInitializer()
+            : new HttpCredentialsAdapter(credential);
+
+    return new ChainingHttpRequestInitializer(credentialsInitializer, 
retryHttpRequestInitializer);
+  }
+
+  private static void recordRequestMetrics(
+      HttpRequest req,
+      Counter requestCount,
+      @Nullable Counter rangeRequestCount,
+      @Nullable Counter unboundedStreamCount,
+      @Nullable Counter otherRequestCount) {
+    String method = req.getRequestMethod();
+    requestCount.inc();
+    if ("GET".equalsIgnoreCase(method)) {
+      String range = req.getHeaders() != null ? req.getHeaders().getRange() : 
null;
+      if (range != null) {
+        if (rangeRequestCount != null) {
+          rangeRequestCount.inc();
+        }
+      } else {
+        if (unboundedStreamCount != null) {
+          unboundedStreamCount.inc();
+        }
+      }
+    } else if (otherRequestCount != null) {
+      // Not a GET, e.g. the POST of a batched metadata lookup. Counted so 
that the three shape
+      // counters add up to requestCount.
+      otherRequestCount.inc();
+    }
+  }
+
+  private static void recordResponseMetrics(
+      HttpResponse res,
+      Counter status2xx,
+      Counter status3xx,
+      Counter status4xx,
+      Counter status5xx,
+      Counter statusOther) {
+    int code = res.getStatusCode();
+    if (code >= 200 && code < 300) {
+      status2xx.inc();
+    } else if (code >= 300 && code < 400) {
+      // Not an error: a resumable upload answers 308 Resume Incomplete to 
every chunk but the last.
+      status3xx.inc();
+    } else if (code >= 400 && code < 500) {
+      status4xx.inc();
+    } else if (code >= 500 && code < 600) {
+      status5xx.inc();
     } else {
-      return new ChainingHttpRequestInitializer(
-          new HttpCredentialsAdapter(credential), retryHttpRequestInitializer);
+      statusOther.inc();
     }
   }
 }
diff --git 
a/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/storage/GcsFileSystemTest.java
 
b/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/storage/GcsFileSystemTest.java
index daa419abb57..c4f716df082 100644
--- 
a/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/storage/GcsFileSystemTest.java
+++ 
b/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/storage/GcsFileSystemTest.java
@@ -20,9 +20,12 @@ package org.apache.beam.sdk.extensions.gcp.storage;
 import static org.hamcrest.MatcherAssert.assertThat;
 import static org.hamcrest.Matchers.contains;
 import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.fail;
+import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.ArgumentMatchers.anyString;
 import static org.mockito.ArgumentMatchers.eq;
 import static org.mockito.ArgumentMatchers.isNull;
+import static org.mockito.Mockito.doThrow;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.times;
 import static org.mockito.Mockito.verify;
@@ -30,11 +33,13 @@ import static org.mockito.Mockito.when;
 
 import com.google.api.services.storage.model.Objects;
 import com.google.api.services.storage.model.StorageObject;
+import java.io.Closeable;
 import java.io.FileNotFoundException;
 import java.io.IOException;
 import java.math.BigInteger;
 import java.util.ArrayList;
 import java.util.List;
+import org.apache.beam.runners.core.metrics.MetricsContainerImpl;
 import org.apache.beam.sdk.extensions.gcp.options.GcsOptions;
 import org.apache.beam.sdk.extensions.gcp.util.GcsUtil;
 import 
org.apache.beam.sdk.extensions.gcp.util.GcsUtil.StorageObjectOrIOException;
@@ -42,6 +47,8 @@ import org.apache.beam.sdk.extensions.gcp.util.gcsfs.GcsPath;
 import org.apache.beam.sdk.io.fs.MatchResult;
 import org.apache.beam.sdk.io.fs.MatchResult.Status;
 import org.apache.beam.sdk.metrics.Lineage;
+import org.apache.beam.sdk.metrics.MetricName;
+import org.apache.beam.sdk.metrics.MetricsEnvironment;
 import org.apache.beam.sdk.options.PipelineOptionsFactory;
 import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.FluentIterable;
 import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList;
@@ -268,4 +275,58 @@ public class GcsFileSystemTest {
         .transform(metadata -> ((GcsResourceId) 
metadata.resourceId()).getGcsPath().toString())
         .toList();
   }
+
+  private GcsFileSystem fileSystemWithMetrics(boolean enabled) {
+    GcsOptions gcsOptions = PipelineOptionsFactory.as(GcsOptions.class);
+    gcsOptions.setGcsUtil(mockGcsUtil);
+    gcsOptions.setGcsPerformanceMetrics(enabled);
+    return new GcsFileSystem(gcsOptions);
+  }
+
+  private static long counter(MetricsContainerImpl container, String name) {
+    return container.getCounter(MetricName.named(GcsUtil.METRIC_NAMESPACE, 
name)).getCumulative();
+  }
+
+  private static List<GcsResourceId> resourceIds(String... uris) {
+    return FluentIterable.from(uris)
+        .transform(uri -> GcsResourceId.fromGcsPath(GcsPath.fromUri(uri)))
+        .toList();
+  }
+
+  @Test
+  public void testOperationMetricsAreOffByDefault() throws Exception {
+    MetricsContainerImpl container = new MetricsContainerImpl(null);
+    try (Closeable ignored = 
MetricsEnvironment.scopedMetricsContainer(container)) {
+      fileSystemWithMetrics(false)
+          .rename(resourceIds("gs://bucket/from"), 
resourceIds("gs://bucket/to"));
+    }
+    assertEquals(0L, counter(container, "gcs_op_rename_count"));
+    assertEquals(0L, counter(container, "gcs_op_rename_msec"));
+  }
+
+  @Test
+  public void testRenameRecordsObjectCountWhenEnabled() throws Exception {
+    MetricsContainerImpl container = new MetricsContainerImpl(null);
+    try (Closeable ignored = 
MetricsEnvironment.scopedMetricsContainer(container)) {
+      fileSystemWithMetrics(true)
+          .rename(
+              resourceIds("gs://bucket/a", "gs://bucket/b"),
+              resourceIds("gs://bucket/c", "gs://bucket/d"));
+    }
+    assertEquals(2L, counter(container, "gcs_op_rename_count"));
+  }
+
+  @Test
+  public void testFailedOperationIsStillRecorded() throws Exception {
+    doThrow(new IOException("boom")).when(mockGcsUtil).copy(any(), any());
+
+    MetricsContainerImpl container = new MetricsContainerImpl(null);
+    try (Closeable ignored = 
MetricsEnvironment.scopedMetricsContainer(container)) {
+      fileSystemWithMetrics(true).copy(resourceIds("gs://bucket/a"), 
resourceIds("gs://bucket/b"));
+      fail("Expected the copy to propagate the IOException");
+    } catch (IOException expected) {
+      // The point of the test is that the finally block still recorded the 
attempt.
+    }
+    assertEquals(1L, counter(container, "gcs_op_copy_count"));
+  }
 }
diff --git 
a/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilTest.java
 
b/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilTest.java
index 2f77f15dcff..52bb877bb96 100644
--- 
a/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilTest.java
+++ 
b/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilTest.java
@@ -1703,7 +1703,8 @@ public class GcsUtilTest {
                   : null,
               gcsOptions.getEnableBucketWriteMetricCounter()
                   ? gcsOptions.getGcsWriteCounterPrefix()
-                  : null),
+                  : null,
+              Boolean.TRUE.equals(gcsOptions.getGcsPerformanceMetrics())),
           gcsOptions.getGoogleCloudStorageReadOptions());
     }
 
@@ -1731,7 +1732,10 @@ public class GcsUtilTest {
 
     @Override
     GoogleCloudStorage createGoogleCloudStorage(
-        GoogleCloudStorageOptions options, Storage storage, Credentials 
credentials) {
+        GoogleCloudStorageOptions options,
+        Storage storage,
+        Credentials credentials,
+        @Nullable HttpRequestInitializer httpRequestInitializer) {
       return googleCloudStorage;
     }
   }
diff --git 
a/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/TransportTest.java
 
b/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/TransportTest.java
index a290d1d78b6..49c636eaa41 100644
--- 
a/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/TransportTest.java
+++ 
b/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/TransportTest.java
@@ -21,15 +21,22 @@ import static org.hamcrest.MatcherAssert.assertThat;
 import static org.hamcrest.Matchers.greaterThan;
 import static org.hamcrest.Matchers.greaterThanOrEqualTo;
 import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertSame;
 
+import com.google.api.client.http.GenericUrl;
 import com.google.api.client.http.HttpRequest;
+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.api.services.storage.Storage;
 import java.io.IOException;
 import java.util.Arrays;
 import java.util.Collections;
+import org.apache.beam.runners.core.metrics.MetricsContainerImpl;
 import org.apache.beam.sdk.extensions.gcp.auth.TestCredential;
 import org.apache.beam.sdk.extensions.gcp.options.GcsOptions;
 import 
org.apache.beam.sdk.extensions.gcp.options.GcsOptions.GcsCustomAuditEntries;
+import org.apache.beam.sdk.metrics.MetricName;
 import org.apache.beam.sdk.options.PipelineOptionsFactory;
 import org.apache.beam.sdk.util.ReleaseInfo;
 import org.junit.Test;
@@ -90,4 +97,103 @@ public class TransportTest {
         
request.getHeaders().getHeaderStringValues("x-goog-custom-audit-status"),
         Collections.singletonList("ok"));
   }
+
+  private static final String READ_PREFIX = "gcs_http_read_";
+  private static final String WRITE_PREFIX = "gcs_http_write_";
+
+  private static long counter(MetricsContainerImpl container, String name) {
+    return container.getCounter(MetricName.named(GcsUtil.METRIC_NAMESPACE, 
name)).getCumulative();
+  }
+
+  /**
+   * Executes one request through a metrics-wrapped initializer against a mock 
transport, and
+   * returns the container the counters were recorded against. An empty {@code 
range} sends no Range
+   * header.
+   */
+  private static MetricsContainerImpl executeRequest(
+      boolean isWrite, String method, int statusCode, String range) throws 
IOException {
+    MetricsContainerImpl container = new MetricsContainerImpl(null);
+    MockHttpTransport transport =
+        new MockHttpTransport.Builder()
+            .setLowLevelHttpResponse(new 
MockLowLevelHttpResponse().setStatusCode(statusCode))
+            .build();
+    HttpRequest request =
+        transport
+            .createRequestFactory(Transport.withMetricsContainer(req -> {}, 
container, isWrite))
+            .buildRequest(method, new 
GenericUrl("https://storage.googleapis.com/test";), null);
+    if (!range.isEmpty()) {
+      request.getHeaders().setRange(range);
+    }
+    // Observe the status that was actually returned, rather than throwing or 
following it.
+    request.setThrowExceptionOnExecuteError(false);
+    request.setFollowRedirects(false);
+    request.execute();
+    return container;
+  }
+
+  @Test
+  public void testReadMetricsCountUnboundedGets() throws IOException {
+    MetricsContainerImpl container = executeRequest(false, "GET", 200, "");
+
+    assertEquals(1, counter(container, READ_PREFIX + "request_count"));
+    assertEquals(1, counter(container, READ_PREFIX + 
"request_count_unbounded"));
+    assertEquals(0, counter(container, READ_PREFIX + "request_count_ranged"));
+    assertEquals(0, counter(container, READ_PREFIX + "request_count_other"));
+    assertEquals(1, counter(container, READ_PREFIX + "status_2xx"));
+    assertEquals(0, counter(container, READ_PREFIX + "request_no_response"));
+  }
+
+  @Test
+  public void testReadMetricsCountRangedGets() throws IOException {
+    MetricsContainerImpl container = executeRequest(false, "GET", 206, 
"bytes=0-9");
+
+    assertEquals(1, counter(container, READ_PREFIX + "request_count"));
+    assertEquals(1, counter(container, READ_PREFIX + "request_count_ranged"));
+    assertEquals(0, counter(container, READ_PREFIX + 
"request_count_unbounded"));
+    assertEquals(0, counter(container, READ_PREFIX + "request_count_other"));
+    // 206 Partial Content is still a success.
+    assertEquals(1, counter(container, READ_PREFIX + "status_2xx"));
+  }
+
+  @Test
+  public void testReadMetricsClassifyNonGetsAsOther() throws IOException {
+    MetricsContainerImpl container = executeRequest(false, "POST", 200, "");
+
+    assertEquals(1, counter(container, READ_PREFIX + "request_count"));
+    assertEquals(1, counter(container, READ_PREFIX + "request_count_other"));
+    assertEquals(0, counter(container, READ_PREFIX + "request_count_ranged"));
+    assertEquals(0, counter(container, READ_PREFIX + 
"request_count_unbounded"));
+  }
+
+  @Test
+  public void testWriteMetricsAreNotClassifiedByRequestShape() throws 
IOException {
+    MetricsContainerImpl container = executeRequest(true, "POST", 200, "");
+
+    assertEquals(1, counter(container, WRITE_PREFIX + "request_count"));
+    assertEquals(1, counter(container, WRITE_PREFIX + "status_2xx"));
+    // Writes are POSTs and PUTs by construction, so the shape counters are 
never allocated.
+    assertEquals(0, counter(container, WRITE_PREFIX + "request_count_ranged"));
+    assertEquals(0, counter(container, WRITE_PREFIX + 
"request_count_unbounded"));
+    assertEquals(0, counter(container, WRITE_PREFIX + "request_count_other"));
+  }
+
+  @Test
+  public void testResponsesAreCountedByStatusClass() throws IOException {
+    // A resumable upload answers 308 to every chunk but the last, so 3xx is 
not an error.
+    assertEquals(1, counter(executeRequest(true, "PUT", 308, ""), WRITE_PREFIX 
+ "status_3xx"));
+    assertEquals(1, counter(executeRequest(false, "GET", 404, ""), READ_PREFIX 
+ "status_4xx"));
+    assertEquals(1, counter(executeRequest(false, "GET", 503, ""), READ_PREFIX 
+ "status_5xx"));
+
+    // Each of those is still exactly one request, and none of them lands in 
2xx.
+    MetricsContainerImpl notFound = executeRequest(false, "GET", 404, "");
+    assertEquals(1, counter(notFound, READ_PREFIX + "request_count"));
+    assertEquals(0, counter(notFound, READ_PREFIX + "status_2xx"));
+  }
+
+  @Test
+  public void testInitializerIsUnchangedWithoutAContainer() {
+    HttpRequestInitializer base = request -> {};
+    assertSame(base, Transport.withMetricsContainer(base, null, false));
+    assertSame(base, Transport.withMetricsContainer(base, null, true));
+  }
 }
diff --git 
a/sdks/java/io/file-based-io-tests/src/test/java/org/apache/beam/sdk/io/text/TextIOIT.java
 
b/sdks/java/io/file-based-io-tests/src/test/java/org/apache/beam/sdk/io/text/TextIOIT.java
index d0ea19ffdf8..35155c8cec6 100644
--- 
a/sdks/java/io/file-based-io-tests/src/test/java/org/apache/beam/sdk/io/text/TextIOIT.java
+++ 
b/sdks/java/io/file-based-io-tests/src/test/java/org/apache/beam/sdk/io/text/TextIOIT.java
@@ -216,10 +216,10 @@ public class TextIOIT {
     if (gatherGcsPerformanceMetrics) {
       metricSuppliers.add(
           reader -> {
-            MetricsReader actualReader =
-                
reader.withNamespace("org.apache.beam.sdk.extensions.gcp.storage.GcsFileSystem");
-            long numRenames = actualReader.getCounterMetric("num_renames");
-            long renameTimeMsec = 
actualReader.getCounterMetric("rename_time_msec");
+            // Namespace and names are defined by GcsUtil.METRIC_NAMESPACE / 
GcsFileSystem.
+            MetricsReader actualReader = reader.withNamespace("Gcs");
+            long numRenames = 
actualReader.getCounterMetric("gcs_op_rename_count");
+            long renameTimeMsec = 
actualReader.getCounterMetric("gcs_op_rename_msec");
             double remamePerSec =
                 (numRenames < 0 || renameTimeMsec < 0) ? -1 : numRenames / 
(renameTimeMsec / 1e3);
             return NamedTestResult.create(uuid, timestamp, "rename_per_sec", 
remamePerSec);

Reply via email to