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