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 94972054878 Route the remaining GCS calls through GcsUtilV2 when
enabled (#40273)
94972054878 is described below
commit 949720548789dbedd751fab7eb8f0be5221ec233
Author: Shunping Huang <[email protected]>
AuthorDate: Mon Sep 28 12:22:49 2026 -0400
Route the remaining GCS calls through GcsUtilV2 when enabled (#40273)
* Route the remaining GCS calls through GcsUtilV2
The facade only had V2 branches on its GcsPath-typed methods, but
GcsFileSystem and GcpOptions call the String-typed and legacy-model ones,
so match, list, copy, rename, delete and bucket creation all stayed on V1
even with use_gcsutil_v2 set. Add V2 branches to those, converting Blob
and BucketInfo back to the JSON API model so their callers need no change.
Strategies are picked to match V1: copy and rename rewrite without a
destination precondition, remove ignores a 404, and buckets are created
with projectPrivate ACLs.
* Match GcsUtilV1's upload chunk size in GcsUtilV2
The formula is copied rather than read from AsyncWriteChannelOptions
because GcsUtilV2 has no other gcsio reference and the migration ends
with that dependency deleted. The test keeps the two values in sync.
* Fix a bug of use_gcsutil_v2 flag not passing to experiments in ParquetIOLT
* Fix GcsUtilV2 parity gaps.
- Translate StorageException thrown on write channel close.
- Record a 404 request metric when open() hits a missing object.
* Fix bucketAccessible in V2 to match V1
* Route the deprecated getBucket and create(path, type) through GcsUtilV2
* Add some more routing tests and trim unnecessary ones.
* Make GcsUtilV2's storage client mockable in tests
Route all client access through a @VisibleForTesting storage() accessor.
In GcsUtilTest, merge the V2 helpers into one, test bucketAccessible
through a mocked client, and reset metrics containers in tearDown.
* Add reference link for default upload chunk size
---
.../apache/beam/it/gcp/storage/ParquetIOLT.java | 47 +-
.../beam/sdk/extensions/gcp/util/GcsUtil.java | 218 +++++++-
.../beam/sdk/extensions/gcp/util/GcsUtilV2.java | 94 +++-
.../beam/sdk/extensions/gcp/util/GcsUtilTest.java | 559 +++++++++++++++++++++
.../sdk/extensions/gcp/util/GcsUtilV2Test.java | 12 +
5 files changed, 885 insertions(+), 45 deletions(-)
diff --git
a/it/google-cloud-platform/src/test/java/org/apache/beam/it/gcp/storage/ParquetIOLT.java
b/it/google-cloud-platform/src/test/java/org/apache/beam/it/gcp/storage/ParquetIOLT.java
index 80eb7886c95..52b12372dcd 100644
---
a/it/google-cloud-platform/src/test/java/org/apache/beam/it/gcp/storage/ParquetIOLT.java
+++
b/it/google-cloud-platform/src/test/java/org/apache/beam/it/gcp/storage/ParquetIOLT.java
@@ -31,7 +31,9 @@ import java.nio.ByteBuffer;
import java.time.Duration;
import java.time.ZoneOffset;
import java.time.format.DateTimeFormatter;
+import java.util.ArrayList;
import java.util.Arrays;
+import java.util.List;
import java.util.Map;
import java.util.Random;
import java.util.UUID;
@@ -120,9 +122,9 @@ import org.junit.runners.MethodSorters;
*
* <p>Whichever is chosen, the job is launched with an explicit experiment,
{@code use_runner_v2} or
* {@code disable_runner_v2}. Leaving the choice to the service is not an
option here because the
- * container image is resolved on the client, see {@link
#dataflowWorkerExperiment()}: a job
- * submitted without an experiment ends up asking for an image tag that does
not exist and hangs
- * with the workers in ImagePullBackOff.
+ * container image is resolved on the client, see {@link
#dataflowExperiments()}: a job submitted
+ * without an experiment ends up asking for an image tag that does not exist
and hangs with the
+ * workers in ImagePullBackOff.
*
* <p>Note that Runner v2 stages the locally built SDK jars, so a local SDK
change is measured as
* is, while the legacy worker runs the Beam code baked into its container
image.
@@ -198,7 +200,7 @@ public final class ParquetIOLT extends GcsIOLoadTestBase {
/**
* Experiment that keeps the job on the legacy worker. Needed even though
the legacy worker is
* what a job without experiments is submitted as, because the service
upgrades such a job to
- * Runner v2 on its own, see {@link #dataflowWorkerExperiment()}.
+ * Runner v2 on its own, see {@link #dataflowExperiments()}.
*/
private static final String LEGACY_WORKER_EXPERIMENT = "disable_runner_v2";
@@ -499,8 +501,8 @@ public final class ParquetIOLT extends GcsIOLoadTestBase {
// the gcs_* metrics measure.
// maxNumWorkers is deliberately not set, it only bounds an autoscaling
pool.
builder
- // Picks the worker, see dataflowWorkerExperiment().
- .addParameter("experiments", dataflowWorkerExperiment())
+ // Picks the worker and the GcsUtil version, see
dataflowExperiments().
+ .addParameter("experiments", dataflowExperiments())
.addParameter("autoscalingAlgorithm", "NONE")
.addParameter("numWorkers", String.valueOf(DATAFLOW_NUM_WORKERS))
.addParameter("workerMachineType", DATAFLOW_MACHINE_TYPE);
@@ -510,14 +512,19 @@ public final class ParquetIOLT extends GcsIOLoadTestBase {
}
/**
- * Experiment that selects the Dataflow worker, {@code use_runner_v2} or
{@code
- * disable_runner_v2}.
+ * Every experiment the Dataflow job is submitted with, as the comma
separated list the {@code
+ * experiments} pipeline option parses.
*
- * <p>The worker is always selected explicitly, even though Runner v2 is
what the service picks on
- * its own, because the container image is resolved on the client: {@code
- * DataflowRunner.getDefaultContainerImageUrl} takes the image name and the
image tag from the
- * same branch of its {@code useUnifiedWorker()} check, and the two tags
({@code
- * dataflowFnapiContainerVersion} and {@code dataflowLegacyContainerVersion}
in {@code
+ * <p>They have to be joined into one value because the launcher carries the
parameters in a map,
+ * so {@code experiments} can only be given once. That map is also the only
channel that reaches
+ * the job: the launcher rebuilds the pipeline options from these parameters
alone, so an
+ * experiment set on the options of the pipeline would be dropped.
+ *
+ * <p>The list always names the worker, {@code use_runner_v2} or {@code
disable_runner_v2}, even
+ * though Runner v2 is what the service picks on its own, because the
container image is resolved
+ * on the client: {@code DataflowRunner.getDefaultContainerImageUrl} takes
the image name and the
+ * image tag from the same branch of its {@code useUnifiedWorker()} check,
and the two tags
+ * ({@code dataflowFnapiContainerVersion} and {@code
dataflowLegacyContainerVersion} in {@code
* runners/google-cloud-dataflow-java/build.gradle}) are bumped
independently.
*
* <p>A job submitted without an experiment therefore resolves the legacy
pair {@code
@@ -529,8 +536,13 @@ public final class ParquetIOLT extends GcsIOLoadTestBase {
* gives {@code beam_javaNN_sdk:<fnapi tag>}, {@code disable_runner_v2}
keeps the service from
* upgrading the job so {@code beam-javaNN-batch:<legacy tag>} stays correct.
*/
- private static String dataflowWorkerExperiment() {
- return configuration.useRunnerV2 ? RUNNER_V2_EXPERIMENT :
LEGACY_WORKER_EXPERIMENT;
+ private static String dataflowExperiments() {
+ List<String> experiments = new ArrayList<>();
+ experiments.add(configuration.useRunnerV2 ? RUNNER_V2_EXPERIMENT :
LEGACY_WORKER_EXPERIMENT);
+ if (configuration.useGcsUtilV2) {
+ experiments.add(GCS_UTIL_V2_EXPERIMENT);
+ }
+ return String.join(",", experiments);
}
/**
@@ -669,7 +681,7 @@ public final class ParquetIOLT extends GcsIOLoadTestBase {
DATAFLOW_RUNNER.equalsIgnoreCase(configuration.runner)
? String.format(
"%s (--experiments=%s)",
- configuration.useRunnerV2 ? "Runner v2" : "legacy",
dataflowWorkerExperiment())
+ configuration.useRunnerV2 ? "Runner v2" : "legacy",
dataflowExperiments())
: "n/a",
DATAFLOW_RUNNER.equalsIgnoreCase(configuration.runner)
? String.format("%d x %s, autoscaling off", DATAFLOW_NUM_WORKERS,
DATAFLOW_MACHINE_TYPE)
@@ -875,8 +887,7 @@ public final class ParquetIOLT extends GcsIOLoadTestBase {
/**
* Dataflow only. {@code true} runs the job on Runner v2, i.e. the unified
worker, {@code false}
* on the legacy worker. Either way the choice is sent to the service as
an explicit experiment,
- * see {@link ParquetIOLT#dataflowWorkerExperiment()} for why it must not
be left to the
- * service.
+ * see {@link ParquetIOLT#dataflowExperiments()} for why it must not be
left to the service.
*/
@JsonProperty public boolean useRunnerV2 = true;
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 9e2d82ae1eb..382ac6856bf 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
@@ -17,6 +17,7 @@
*/
package org.apache.beam.sdk.extensions.gcp.util;
+import com.google.api.client.util.DateTime;
import com.google.api.gax.paging.Page;
import com.google.api.services.storage.model.Bucket;
import com.google.api.services.storage.model.Objects;
@@ -28,10 +29,17 @@ import com.google.cloud.storage.Storage.BlobListOption;
import com.google.cloud.storage.Storage.BlobSourceOption;
import com.google.cloud.storage.Storage.BlobWriteOption;
import com.google.cloud.storage.Storage.BucketGetOption;
+import com.google.cloud.storage.Storage.BucketTargetOption;
+import com.google.cloud.storage.Storage.PredefinedAcl;
+import com.google.cloud.storage.StorageClass;
import edu.umd.cs.findbugs.annotations.SuppressFBWarnings;
import java.io.IOException;
+import java.math.BigInteger;
import java.nio.channels.SeekableByteChannel;
import java.nio.channels.WritableByteChannel;
+import java.time.Duration;
+import java.time.OffsetDateTime;
+import java.util.ArrayList;
import java.util.Collection;
import java.util.List;
import java.util.Set;
@@ -151,6 +159,9 @@ public class GcsUtil {
*/
@Deprecated
public StorageObject getObject(GcsPath gcsPath) throws IOException {
+ if (delegateV2 != null) {
+ return toStorageObject(delegateV2.getBlob(gcsPath));
+ }
return delegate.getObject(gcsPath);
}
@@ -166,6 +177,21 @@ public class GcsUtil {
*/
@Deprecated
public List<StorageObjectOrIOException> getObjects(List<GcsPath> gcsPaths)
throws IOException {
+ if (delegateV2 != null) {
+ List<StorageObjectOrIOException> results = new ArrayList<>();
+ for (BlobResult blobResult : delegateV2.getBlobs(gcsPaths)) {
+ Blob blob = blobResult.blob();
+ IOException ioException = blobResult.ioException();
+ if (blob != null) {
+
results.add(StorageObjectOrIOException.create(toStorageObject(blob)));
+ } else if (ioException != null) {
+ results.add(StorageObjectOrIOException.create(ioException));
+ } else {
+ throw new IOException("Invalid blob result: it holds neither a blob
nor an error.");
+ }
+ }
+ return results;
+ }
List<GcsUtilV1.StorageObjectOrIOException> legacy =
delegate.getObjects(gcsPaths);
return legacy.stream()
.map(StorageObjectOrIOException::fromLegacy)
@@ -186,6 +212,9 @@ public class GcsUtil {
@Deprecated
public Objects listObjects(String bucket, String prefix, @Nullable String
pageToken)
throws IOException {
+ if (delegateV2 != null) {
+ return toObjects(delegateV2.listBlobs(bucket, prefix, pageToken));
+ }
return delegate.listObjects(bucket, prefix, pageToken);
}
@@ -196,6 +225,9 @@ public class GcsUtil {
public Objects listObjects(
String bucket, String prefix, @Nullable String pageToken, @Nullable
String delimiter)
throws IOException {
+ if (delegateV2 != null) {
+ return toObjects(delegateV2.listBlobs(bucket, prefix, pageToken,
delimiter));
+ }
return delegate.listObjects(bucket, prefix, pageToken, delimiter);
}
@@ -240,7 +272,9 @@ public class GcsUtil {
*/
@Deprecated
public WritableByteChannel create(GcsPath path, String type) throws
IOException {
- return delegate.create(path, type);
+ // Built the same way as GcsUtilV1#create(GcsPath, String), but through
this class so that it
+ // follows GcsUtilV2 when enabled.
+ return create(path, CreateOptions.builder().setContentType(type).build());
}
/**
@@ -249,7 +283,12 @@ public class GcsUtil {
@Deprecated
public WritableByteChannel create(GcsPath path, String type, Integer
uploadBufferSizeBytes)
throws IOException {
- return delegate.create(path, type, uploadBufferSizeBytes);
+ return create(
+ path,
+ CreateOptions.builder()
+ .setContentType(type)
+ .setUploadBufferSizeBytes(uploadBufferSizeBytes)
+ .build());
}
public static class CreateOptions {
@@ -341,16 +380,27 @@ public class GcsUtil {
}
/**
- * @deprecated use {@link #createBucket(BucketInfo)}.
+ * @deprecated use {@link #createBucket(BucketInfo, BucketTargetOption...)}.
*/
@Deprecated
public void createBucket(String projectId, Bucket bucket) throws IOException
{
+ if (delegateV2 != null) {
+ // GcsUtilV1 always creates buckets with projectPrivate ACLs, which
java-storage does not do
+ // on its own, so they have to be requested explicitly to keep the same
access.
+ delegateV2.createBucket(
+ projectId,
+ toBucketInfo(bucket),
+ BucketTargetOption.predefinedAcl(PredefinedAcl.PROJECT_PRIVATE),
+
BucketTargetOption.predefinedDefaultObjectAcl(PredefinedAcl.PROJECT_PRIVATE));
+ return;
+ }
delegate.createBucket(projectId, bucket);
}
- public void createBucket(BucketInfo bucketInfo) throws IOException {
+ public void createBucket(BucketInfo bucketInfo, BucketTargetOption...
options)
+ throws IOException {
if (delegateV2 != null) {
- delegateV2.createBucket(bucketInfo);
+ delegateV2.createBucket(bucketInfo, options);
} else {
throw new IOException("GcsUtil V2 not initialized.");
}
@@ -361,6 +411,9 @@ public class GcsUtil {
*/
@Deprecated
public @Nullable Bucket getBucket(GcsPath path) throws IOException {
+ if (delegateV2 != null) {
+ return toBucket(delegateV2.getBucket(path));
+ }
return delegate.getBucket(path);
}
@@ -377,6 +430,10 @@ public class GcsUtil {
*/
@Deprecated
public void removeBucket(Bucket bucket) throws IOException {
+ if (delegateV2 != null) {
+ delegateV2.removeBucket(toBucketInfo(bucket));
+ return;
+ }
delegate.removeBucket(bucket);
}
@@ -390,6 +447,14 @@ public class GcsUtil {
public void copy(Iterable<String> srcFilenames, Iterable<String>
destFilenames)
throws IOException {
+ if (delegateV2 != null) {
+ // GcsUtilV1 issues a rewrite without any destination precondition, so
ALWAYS_OVERWRITE is
+ // the strategy that preserves its behavior. The strategies that inspect
the destination
+ // would also cost an extra GET per file.
+ delegateV2.copy(
+ toGcsPaths(srcFilenames), toGcsPaths(destFilenames),
OverwriteStrategy.ALWAYS_OVERWRITE);
+ return;
+ }
delegate.copy(srcFilenames, destFilenames);
}
@@ -412,6 +477,22 @@ public class GcsUtil {
public void rename(
Iterable<String> srcFilenames, Iterable<String> destFilenames,
MoveOptions... moveOptions)
throws IOException {
+ GcsUtilV2 v2 = delegateV2;
+ if (v2 != null) {
+ Set<MoveOptions> moveOptionSet = Sets.newHashSet(moveOptions);
+ // Note this differs from renameV2, which defaults to SAFE_OVERWRITE.
GcsUtilV1 rewrites
+ // without a destination precondition, so ALWAYS_OVERWRITE is the
behavior preserving choice.
+ v2.move(
+ toGcsPaths(srcFilenames),
+ toGcsPaths(destFilenames),
+ moveOptionSet.contains(StandardMoveOptions.IGNORE_MISSING_FILES)
+ ? MissingStrategy.SKIP_IF_MISSING
+ : MissingStrategy.FAIL_IF_MISSING,
+
moveOptionSet.contains(StandardMoveOptions.SKIP_IF_DESTINATION_EXISTS)
+ ? OverwriteStrategy.SKIP_IF_EXISTS
+ : OverwriteStrategy.ALWAYS_OVERWRITE);
+ return;
+ }
delegate.rename(srcFilenames, destFilenames, moveOptions);
}
@@ -453,6 +534,11 @@ public class GcsUtil {
}
public void remove(Collection<String> filenames) throws IOException {
+ if (delegateV2 != null) {
+ // GcsUtilV1 ignores a 404 on delete, which is SKIP_IF_MISSING.
+ delegateV2.remove(toGcsPaths(filenames),
MissingStrategy.SKIP_IF_MISSING);
+ return;
+ }
delegate.remove(filenames);
}
@@ -470,6 +556,128 @@ public class GcsUtil {
}
}
+ private static List<GcsPath> toGcsPaths(Iterable<String> filenames) {
+ List<GcsPath> paths = new ArrayList<>();
+ for (String filename : filenames) {
+ paths.add(GcsPath.fromUri(filename));
+ }
+ return paths;
+ }
+
+ /**
+ * Converts a JSON API {@link Bucket} into the java-storage {@link
BucketInfo} model.
+ *
+ * <p>Only the properties that callers of the deprecated {@link
#createBucket(String, Bucket)} set
+ * are carried over. Like {@link #toStorageObject}, this is expected to go
away with the
+ * deprecated methods it serves.
+ */
+ private static BucketInfo toBucketInfo(Bucket bucket) {
+ BucketInfo.Builder builder = BucketInfo.newBuilder(bucket.getName());
+ if (bucket.getLocation() != null) {
+ builder.setLocation(bucket.getLocation());
+ }
+ if (bucket.getStorageClass() != null) {
+ builder.setStorageClass(StorageClass.valueOf(bucket.getStorageClass()));
+ }
+ Bucket.SoftDeletePolicy softDeletePolicy = bucket.getSoftDeletePolicy();
+ if (softDeletePolicy != null &&
softDeletePolicy.getRetentionDurationSeconds() != null) {
+ builder.setSoftDeletePolicy(
+ BucketInfo.SoftDeletePolicy.newBuilder()
+ .setRetentionDuration(
+
Duration.ofSeconds(softDeletePolicy.getRetentionDurationSeconds()))
+ .build());
+ }
+ return builder.build();
+ }
+
+ /**
+ * Converts a java-storage {@link com.google.cloud.storage.Bucket} back into
the JSON API {@link
+ * Bucket} model, for the deprecated {@link #getBucket(GcsPath)}.
+ *
+ * <p>The reverse of {@link #toBucketInfo}, plus the owning project number.
Like it, this is
+ * expected to go away with the deprecated methods it serves.
+ */
+ private static Bucket toBucket(com.google.cloud.storage.Bucket bucketInfo) {
+ Bucket bucket =
+ new Bucket()
+ .setName(bucketInfo.getName())
+ .setLocation(bucketInfo.getLocation())
+ .setProjectNumber(bucketInfo.getProject());
+ StorageClass storageClass = bucketInfo.getStorageClass();
+ if (storageClass != null) {
+ bucket.setStorageClass(storageClass.name());
+ }
+ BucketInfo.SoftDeletePolicy softDeletePolicy =
bucketInfo.getSoftDeletePolicy();
+ Duration retention = softDeletePolicy == null ? null :
softDeletePolicy.getRetentionDuration();
+ if (retention != null) {
+ bucket.setSoftDeletePolicy(
+ new
Bucket.SoftDeletePolicy().setRetentionDurationSeconds(retention.getSeconds()));
+ }
+ return bucket;
+ }
+
+ /**
+ * Converts a java-storage {@link Blob} back into the JSON API {@link
StorageObject} model.
+ *
+ * <p>This lets the deprecated, legacy typed methods of this class be served
by {@link GcsUtilV2}
+ * without their callers having to change. It is expected to go away once
those methods do.
+ */
+ private static StorageObject toStorageObject(Blob blob) {
+ StorageObject storageObject =
+ new StorageObject()
+ .setBucket(blob.getBucket())
+ .setName(blob.getName())
+ .setGeneration(blob.getGeneration())
+ .setMetageneration(blob.getMetageneration())
+ .setContentType(blob.getContentType())
+ .setContentEncoding(blob.getContentEncoding())
+ .setMd5Hash(blob.getMd5())
+ .setCrc32c(blob.getCrc32c())
+ .setEtag(blob.getEtag());
+ Long size = blob.getSize();
+ if (size != null) {
+ storageObject.setSize(BigInteger.valueOf(size));
+ }
+ OffsetDateTime updated = blob.getUpdateTimeOffsetDateTime();
+ if (updated != null) {
+ storageObject.setUpdated(new
DateTime(updated.toInstant().toEpochMilli()));
+ }
+ OffsetDateTime created = blob.getCreateTimeOffsetDateTime();
+ if (created != null) {
+ storageObject.setTimeCreated(new
DateTime(created.toInstant().toEpochMilli()));
+ }
+ return storageObject;
+ }
+
+ /** Converts a single page of java-storage {@link Blob}s into the JSON API
{@link Objects}. */
+ private static Objects toObjects(Page<Blob> page) {
+ List<StorageObject> items = new ArrayList<>();
+ List<String> prefixes = new ArrayList<>();
+ for (Blob blob : page.getValues()) {
+ // A delimited listing reports each common prefix as a directory
placeholder.
+ if (blob.isDirectory()) {
+ prefixes.add(blob.getName());
+ } else {
+ items.add(toStorageObject(blob));
+ }
+ }
+ Objects objects = new Objects();
+ // Leave items and prefixes null when empty, as the JSON API does, so that
callers looping on
+ // getItems() != null keep terminating.
+ if (!items.isEmpty()) {
+ objects.setItems(items);
+ }
+ if (!prefixes.isEmpty()) {
+ objects.setPrefixes(prefixes);
+ }
+ // Page.getNextPageToken() may be an empty string rather than null on the
last page.
+ String nextPageToken = page.hasNextPage() ? page.getNextPageToken() : null;
+ if (nextPageToken != null) {
+ objects.setNextPageToken(nextPageToken);
+ }
+ return objects;
+ }
+
@SuppressFBWarnings("NM_CLASS_NOT_EXCEPTION")
public static class StorageObjectOrIOException {
final GcsUtilV1.StorageObjectOrIOException delegate;
diff --git
a/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV2.java
b/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV2.java
index 9d223f81050..20260aa3c79 100644
---
a/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV2.java
+++
b/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV2.java
@@ -45,6 +45,7 @@ import com.google.cloud.storage.Storage.BlobSourceOption;
import com.google.cloud.storage.Storage.BlobWriteOption;
import com.google.cloud.storage.Storage.BucketField;
import com.google.cloud.storage.Storage.BucketGetOption;
+import com.google.cloud.storage.Storage.BucketTargetOption;
import com.google.cloud.storage.Storage.CopyRequest;
import com.google.cloud.storage.StorageBatch;
import com.google.cloud.storage.StorageBatchResult;
@@ -100,7 +101,16 @@ class GcsUtilV2 {
}
}
- private Storage storage;
+ private final Storage storage;
+
+ /**
+ * The shared client. Every operation reaches it through this method, so
that tests can substitute
+ * a mocked client for a spied instance.
+ */
+ @VisibleForTesting
+ Storage storage() {
+ return storage;
+ }
private final @Nullable Integer uploadBufferSizeBytes;
@@ -116,6 +126,15 @@ class GcsUtilV2 {
/** Maximum number of requests permitted in a GCS batch request. */
private static final int MAX_REQUESTS_PER_BATCH = 100;
+ /**
+ * Upload chunk size applied when the pipeline does not ask for one. Mirrors
gcsio's {@code
+ * AsyncWriteChannelOptions} default, which java-storage does not share.
Reference link:
+ *
https://github.com/GoogleCloudDataproc/hadoop-connectors/blob/v3.1.14/util/src/main/java/com/google/cloud/hadoop/util/AsyncWriteChannelOptions.java#L72
+ */
+ @VisibleForTesting
+ static final int DEFAULT_UPLOAD_CHUNK_SIZE_BYTES =
+ Runtime.getRuntime().maxMemory() < 512 * 1024 * 1024 ? 8 * 1024 * 1024 :
3 * 8 * 1024 * 1024;
+
/**
* Limit the number of bytes Cloud Storage will attempt to copy before
responding to an individual
* request. If you see Read Timeout errors, try reducing this value.
@@ -294,13 +313,13 @@ class GcsUtilV2 {
@VisibleForTesting
Storage storageWithHttpMetrics(@Nullable MetricsContainer container, boolean
isWrite) {
if (container == null) {
- return storage;
+ return storage();
}
- StorageOptions options = storage.getOptions();
+ StorageOptions options = storage().getOptions();
TransportOptions transportOptions = options.getTransportOptions();
if (!(transportOptions instanceof HttpTransportOptions)) {
// A non-HTTP transport (e.g. gRPC) has no HttpRequestInitializer to
wrap.
- return storage;
+ return storage();
}
return options.toBuilder()
.setTransportOptions(
@@ -346,7 +365,7 @@ class GcsUtilV2 {
}
public Blob getBlob(GcsPath gcsPath, BlobGetOption... options) throws
IOException {
- return getBlob(storage, gcsPath, options);
+ return getBlob(storage(), gcsPath, options);
}
/** As {@link #getBlob(GcsPath, BlobGetOption...)}, but issued through a
specific client. */
@@ -399,7 +418,7 @@ class GcsUtilV2 {
Lists.partition(Lists.newArrayList(gcsPaths), MAX_REQUESTS_PER_BATCH))
{
// Create a new empty batch every time
- StorageBatch batch = storage.batch();
+ StorageBatch batch = storage().batch();
List<StorageBatchResult<Blob>> batchResultFutures = new ArrayList<>();
for (GcsPath path : pathPartition) {
@@ -455,7 +474,7 @@ class GcsUtilV2 {
}
try {
- return storage.list(bucket, blobListOptions.toArray(new
BlobListOption[0]));
+ return storage().list(bucket, blobListOptions.toArray(new
BlobListOption[0]));
} catch (StorageException e) {
throw translateStorageException(bucket, prefix, e);
}
@@ -525,7 +544,7 @@ class GcsUtilV2 {
Lists.partition(Lists.newArrayList(paths), MAX_REQUESTS_PER_BATCH)) {
// Create a new empty batch every time
- StorageBatch batch = storage.batch();
+ StorageBatch batch = storage().batch();
List<StorageBatchResult<Boolean>> batchResultFutures = new ArrayList<>();
for (GcsPath path : pathPartition) {
@@ -592,7 +611,7 @@ class GcsUtilV2 {
// FAIL_IF_EXISTS, SKIP_IF_EXISTS and SAFE_OVERWRITE require checking
the target blob
BlobInfo existingTarget;
try {
- existingTarget = storage.get(dstId);
+ existingTarget = storage().get(dstId);
} catch (StorageException e) {
throw translateStorageException(dstPath, e);
}
@@ -622,11 +641,11 @@ class GcsUtilV2 {
}
try {
- CopyWriter copyWriter = storage.copy(copyRequestBuilder.build());
+ CopyWriter copyWriter = storage().copy(copyRequestBuilder.build());
copyWriter.getResult();
if (deleteSrc) {
- if (!storage.delete(srcId)) {
+ if (!storage().delete(srcId)) {
// This may happen if the source file is deleted by another
process after copy.
LOG.warn(
"Source file {} could not be deleted after move to {}. It may
not have existed.",
@@ -664,7 +683,7 @@ class GcsUtilV2 {
public Bucket getBucket(GcsPath path, BucketGetOption... options) throws
IOException {
String bucketName = path.getBucket();
try {
- Bucket bucket = storage.get(bucketName, options);
+ Bucket bucket = storage().get(bucketName, options);
if (bucket == null) {
throw new FileNotFoundException(
String.format("The specified bucket does not exist: gs://%s",
bucketName));
@@ -675,13 +694,17 @@ class GcsUtilV2 {
}
}
- /** Returns whether the GCS bucket exists and is accessible. */
- public boolean bucketAccessible(GcsPath path) {
+ /**
+ * Returns whether the GCS bucket exists and is accessible. This will return
false if the bucket
+ * does not exist or is inaccessible due to permissions; any other failure
is propagated, as
+ * {@link GcsUtilV1} does.
+ */
+ public boolean bucketAccessible(GcsPath path) throws IOException {
try {
// Fetch only the name field to minimize data transfer
getBucket(path, BucketGetOption.fields(BucketField.NAME));
return true;
- } catch (IOException e) {
+ } catch (AccessDeniedException | FileNotFoundException e) {
return false;
}
}
@@ -705,9 +728,27 @@ class GcsUtilV2 {
return bucket.getProject().longValue();
}
- public void createBucket(BucketInfo bucketInfo) throws IOException {
+ public void createBucket(BucketInfo bucketInfo, BucketTargetOption...
options)
+ throws IOException {
+ createBucket(null, bucketInfo, options);
+ }
+
+ /**
+ * As {@link #createBucket(BucketInfo, BucketTargetOption...)}, but creates
the bucket in the
+ * given project instead of the one this instance is configured with.
+ */
+ public void createBucket(
+ @Nullable String projectId, BucketInfo bucketInfo, BucketTargetOption...
options)
+ throws IOException {
+ Storage client = storage();
+ if (projectId != null && !projectId.equals(this.projectId)) {
+ // The owning project is a property of the client rather than of the
insert request, so
+ // asking for a different one means deriving a client for it. The
derived client shares the
+ // credentials, host and transport of the original.
+ client =
storage().getOptions().toBuilder().setProjectId(projectId).build().getService();
+ }
try {
- storage.create(bucketInfo);
+ client.create(bucketInfo, options);
} catch (StorageException e) {
throw translateStorageException(bucketInfo.getName(), null, e);
}
@@ -715,7 +756,7 @@ class GcsUtilV2 {
public void removeBucket(BucketInfo bucketInfo) throws IOException {
try {
- if (!storage.delete(bucketInfo.getName())) {
+ if (!storage().delete(bucketInfo.getName())) {
throw new FileNotFoundException(
String.format("The specified bucket does not exist: gs://%s",
bucketInfo.getName()));
}
@@ -801,6 +842,11 @@ class GcsUtilV2 {
serviceCallMetric.call("ok");
return wrapInCounting(
new GcsSeekableByteChannel(reader, blob.getSize()),
path.getBucket(), container);
+ } catch (FileNotFoundException e) {
+ // getBlob reports a missing object as a FileNotFoundException rather
than a
+ // StorageException, so record its status here like GcsUtilV1 does.
+ serviceCallMetric.call(404);
+ throw e;
} catch (StorageException e) {
serviceCallMetric.call(e.getCode());
throw translateStorageException(path, e);
@@ -833,7 +879,12 @@ class GcsUtilV2 {
@Override
public void close() throws IOException {
- writer.close();
+ // The upload is finalized here, so this is where a failed precondition
surfaces.
+ try {
+ writer.close();
+ } catch (StorageException e) {
+ throw translateStorageException(gcsPath, e);
+ }
}
}
@@ -873,9 +924,8 @@ class GcsUtilV2 {
options.getUploadBufferSizeBytes() != null
? options.getUploadBufferSizeBytes()
: this.uploadBufferSizeBytes;
- if (uploadBufferSizeBytes != null) {
- writer.setChunkSize(uploadBufferSizeBytes);
- }
+ writer.setChunkSize(
+ uploadBufferSizeBytes != null ? uploadBufferSizeBytes :
DEFAULT_UPLOAD_CHUNK_SIZE_BYTES);
serviceCallMetric.call("ok");
// Return the bridge wrapper
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 52bb877bb96..4beb56d6c20 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
@@ -55,18 +55,26 @@ import com.google.api.client.testing.http.MockHttpTransport;
import com.google.api.client.testing.http.MockLowLevelHttpRequest;
import com.google.api.client.testing.http.MockLowLevelHttpResponse;
import com.google.api.client.util.BackOff;
+import com.google.api.gax.paging.Page;
import com.google.api.services.storage.Storage;
import com.google.api.services.storage.model.Bucket;
import com.google.api.services.storage.model.Objects;
import com.google.api.services.storage.model.RewriteResponse;
import com.google.api.services.storage.model.StorageObject;
import com.google.auth.Credentials;
+import com.google.cloud.WriteChannel;
import com.google.cloud.hadoop.gcsio.CreateObjectOptions;
import com.google.cloud.hadoop.gcsio.GoogleCloudStorage;
import com.google.cloud.hadoop.gcsio.GoogleCloudStorageImpl;
import com.google.cloud.hadoop.gcsio.GoogleCloudStorageOptions;
import com.google.cloud.hadoop.gcsio.GoogleCloudStorageReadOptions;
import com.google.cloud.hadoop.gcsio.StorageResourceId;
+import com.google.cloud.storage.BucketInfo;
+import com.google.cloud.storage.Storage.BucketGetOption;
+import com.google.cloud.storage.Storage.BucketTargetOption;
+import com.google.cloud.storage.Storage.PredefinedAcl;
+import com.google.cloud.storage.StorageClass;
+import com.google.cloud.storage.StorageException;
import java.io.ByteArrayInputStream;
import java.io.ByteArrayOutputStream;
import java.io.FileNotFoundException;
@@ -114,9 +122,11 @@ import org.apache.beam.sdk.util.FluentBackoff;
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.checkerframework.checker.nullness.qual.Nullable;
+import org.junit.After;
import org.junit.Before;
import org.junit.Rule;
import org.junit.Test;
+import org.junit.function.ThrowingRunnable;
import org.junit.rules.ExpectedException;
import org.junit.runner.RunWith;
import org.junit.runners.JUnit4;
@@ -136,6 +146,13 @@ public class GcsUtilTest {
MetricsEnvironment.setCurrentContainer(testMetricsContainer);
}
+ @After
+ public void tearDown() {
+ // Don't leak the containers installed by setUp into later tests in the
same JVM.
+ MetricsEnvironment.setProcessWideContainer(null);
+ MetricsEnvironment.setCurrentContainer(null);
+ }
+
private static GcsOptions gcsOptionsWithTestCredential() {
GcsOptions pipelineOptions = PipelineOptionsFactory.as(GcsOptions.class);
pipelineOptions.setGcpCredential(new TestCredential());
@@ -1887,4 +1904,546 @@ public class GcsUtilTest {
private static InputStream toStream(String content) throws IOException {
return new ByteArrayInputStream(content.getBytes(StandardCharsets.UTF_8));
}
+
+ // The tests below cover the routing that the GcsUtil facade performs for
the legacy typed
+ // methods once the use_gcsutil_v2 experiment installs a GcsUtilV2 delegate.
Both delegates are
+ // mocked, so they assert both that V2 is used and that V1 is left alone.
+
+ private GcsUtilV1 mockDelegate;
+ private GcsUtilV2 mockDelegateV2;
+
+ private GcsUtil gcsUtilRoutingToV2() {
+ GcsUtil gcsUtil = gcsOptionsWithTestCredential().getGcsUtil();
+ mockDelegate = Mockito.mock(GcsUtilV1.class);
+ mockDelegateV2 = Mockito.mock(GcsUtilV2.class);
+ gcsUtil.delegate = mockDelegate;
+ gcsUtil.delegateV2 = mockDelegateV2;
+ return gcsUtil;
+ }
+
+ private static com.google.cloud.storage.Blob mockBlob(String bucket, String
object) {
+ com.google.cloud.storage.Blob blob =
Mockito.mock(com.google.cloud.storage.Blob.class);
+ when(blob.getBucket()).thenReturn(bucket);
+ when(blob.getName()).thenReturn(object);
+ return blob;
+ }
+
+ @Test
+ public void testCopyIsRoutedToV2AsAnUnconditionalOverwrite() throws
IOException {
+ GcsUtil gcsUtil = gcsUtilRoutingToV2();
+
+ gcsUtil.copy(ImmutableList.of("gs://bucket/from"),
ImmutableList.of("gs://bucket/to"));
+
+ verify(mockDelegateV2)
+ .copy(
+ ImmutableList.of(GcsPath.fromUri("gs://bucket/from")),
+ ImmutableList.of(GcsPath.fromUri("gs://bucket/to")),
+ GcsUtilV2.OverwriteStrategy.ALWAYS_OVERWRITE);
+ Mockito.verifyNoMoreInteractions(mockDelegate);
+ }
+
+ @Test
+ public void testRenameWithoutOptionsIsRoutedToV2AsAnUnconditionalOverwrite()
throws IOException {
+ GcsUtil gcsUtil = gcsUtilRoutingToV2();
+
+ gcsUtil.rename(ImmutableList.of("gs://bucket/from"),
ImmutableList.of("gs://bucket/to"));
+
+ verify(mockDelegateV2)
+ .move(
+ ImmutableList.of(GcsPath.fromUri("gs://bucket/from")),
+ ImmutableList.of(GcsPath.fromUri("gs://bucket/to")),
+ GcsUtilV2.MissingStrategy.FAIL_IF_MISSING,
+ GcsUtilV2.OverwriteStrategy.ALWAYS_OVERWRITE);
+ Mockito.verifyNoMoreInteractions(mockDelegate);
+ }
+
+ @Test
+ public void testRenameMoveOptionsAreTranslatedForV2() throws IOException {
+ GcsUtil gcsUtil = gcsUtilRoutingToV2();
+
+ gcsUtil.rename(
+ ImmutableList.of("gs://bucket/from"),
+ ImmutableList.of("gs://bucket/to"),
+ StandardMoveOptions.IGNORE_MISSING_FILES,
+ StandardMoveOptions.SKIP_IF_DESTINATION_EXISTS);
+
+ verify(mockDelegateV2)
+ .move(
+ ImmutableList.of(GcsPath.fromUri("gs://bucket/from")),
+ ImmutableList.of(GcsPath.fromUri("gs://bucket/to")),
+ GcsUtilV2.MissingStrategy.SKIP_IF_MISSING,
+ GcsUtilV2.OverwriteStrategy.SKIP_IF_EXISTS);
+ Mockito.verifyNoMoreInteractions(mockDelegate);
+ }
+
+ @Test
+ public void testRemoveIsRoutedToV2AndToleratesMissingFiles() throws
IOException {
+ GcsUtil gcsUtil = gcsUtilRoutingToV2();
+
+ gcsUtil.remove(ImmutableList.of("gs://bucket/one", "gs://bucket/two"));
+
+ verify(mockDelegateV2)
+ .remove(
+ ImmutableList.of(
+ GcsPath.fromUri("gs://bucket/one"),
GcsPath.fromUri("gs://bucket/two")),
+ GcsUtilV2.MissingStrategy.SKIP_IF_MISSING);
+ Mockito.verifyNoMoreInteractions(mockDelegate);
+ }
+
+ /**
+ * Blobs and errors are passed through in order. The field-by-field
conversion is covered by
+ * {@link #testGetObjectIsRoutedToV2AndKeepsAllFields}.
+ */
+ @Test
+ public void testGetObjectsIsRoutedToV2AndConvertsBlobs() throws IOException {
+ GcsUtil gcsUtil = gcsUtilRoutingToV2();
+ com.google.cloud.storage.Blob blob = mockBlob("bucket", "found");
+ FileNotFoundException notFound = new
FileNotFoundException("gs://bucket/missing");
+ List<GcsPath> paths =
+ ImmutableList.of(GcsPath.fromUri("gs://bucket/found"),
GcsPath.fromUri("gs://bucket/miss"));
+ when(mockDelegateV2.getBlobs(paths))
+ .thenReturn(
+ ImmutableList.of(
+ GcsUtilV2.BlobResult.create(blob),
GcsUtilV2.BlobResult.create(notFound)));
+
+ List<StorageObjectOrIOException> results = gcsUtil.getObjects(paths);
+
+ assertEquals(2, results.size());
+ StorageObject converted = results.get(0).storageObject();
+ assertNotNull(converted);
+ assertEquals("found", converted.getName());
+ assertNull(results.get(0).ioException());
+ assertSame(notFound, results.get(1).ioException());
+ assertNull(results.get(1).storageObject());
+ Mockito.verifyNoMoreInteractions(mockDelegate);
+ }
+
+ @Test
+ public void testListObjectsIsRoutedToV2AndConvertsAPage() throws IOException
{
+ GcsUtil gcsUtil = gcsUtilRoutingToV2();
+ com.google.cloud.storage.Blob object = mockBlob("bucket", "prefix/object");
+ com.google.cloud.storage.Blob directory = mockBlob("bucket",
"prefix/dir/");
+ when(directory.isDirectory()).thenReturn(true);
+ @SuppressWarnings("unchecked")
+ Page<com.google.cloud.storage.Blob> page = Mockito.mock(Page.class);
+ when(page.getValues()).thenReturn(ImmutableList.of(object, directory));
+ when(page.hasNextPage()).thenReturn(true);
+ when(page.getNextPageToken()).thenReturn("next");
+ when(mockDelegateV2.listBlobs("bucket", "prefix/", null)).thenReturn(page);
+
+ Objects objects = gcsUtil.listObjects("bucket", "prefix/", null);
+
+ assertEquals(1, objects.getItems().size());
+ assertEquals("prefix/object", objects.getItems().get(0).getName());
+ assertEquals(ImmutableList.of("prefix/dir/"), objects.getPrefixes());
+ assertEquals("next", objects.getNextPageToken());
+ Mockito.verifyNoMoreInteractions(mockDelegate);
+ }
+
+ @Test
+ public void testListObjectsReportsTheLastPageWithANullToken() throws
IOException {
+ GcsUtil gcsUtil = gcsUtilRoutingToV2();
+ @SuppressWarnings("unchecked")
+ Page<com.google.cloud.storage.Blob> page = Mockito.mock(Page.class);
+ when(page.getValues()).thenReturn(ImmutableList.of());
+ // A gax page reports an empty token rather than a null one once it is
exhausted. Callers of
+ // listObjects loop until the token is null, so it has to be normalized.
+ when(page.hasNextPage()).thenReturn(false);
+ when(page.getNextPageToken()).thenReturn("");
+ when(mockDelegateV2.listBlobs("bucket", "prefix/", null)).thenReturn(page);
+
+ Objects objects = gcsUtil.listObjects("bucket", "prefix/", null);
+
+ assertNull(objects.getItems());
+ assertNull(objects.getPrefixes());
+ assertNull(objects.getNextPageToken());
+ }
+
+ /**
+ * The storage class is checked here because the emulator ignores it, so
only a unit test can see
+ * it carried over.
+ */
+ @Test
+ public void testCreateBucketIsRoutedToV2WithProjectPrivateAcls() throws
IOException {
+ GcsUtil gcsUtil = gcsUtilRoutingToV2();
+ // This is the bucket that GcpOptions.tryCreateDefaultBucketWithPrefix
builds, plus a storage
+ // class.
+ Bucket bucket =
+ new Bucket()
+ .setName("bucket")
+ .setLocation("us-central1")
+ .setStorageClass("NEARLINE")
+ .setSoftDeletePolicy(new
Bucket.SoftDeletePolicy().setRetentionDurationSeconds(0L));
+
+ gcsUtil.createBucket("a-project", bucket);
+
+ verify(mockDelegateV2)
+ .createBucket(
+ "a-project",
+ BucketInfo.newBuilder("bucket")
+ .setLocation("us-central1")
+ .setStorageClass(StorageClass.NEARLINE)
+ .setSoftDeletePolicy(
+ BucketInfo.SoftDeletePolicy.newBuilder()
+ .setRetentionDuration(java.time.Duration.ZERO)
+ .build())
+ .build(),
+ BucketTargetOption.predefinedAcl(PredefinedAcl.PROJECT_PRIVATE),
+
BucketTargetOption.predefinedDefaultObjectAcl(PredefinedAcl.PROJECT_PRIVATE));
+ Mockito.verifyNoMoreInteractions(mockDelegate);
+ }
+
+ @Test
+ public void testRemoveBucketIsRoutedToV2() throws IOException {
+ GcsUtil gcsUtil = gcsUtilRoutingToV2();
+
+ gcsUtil.removeBucket(new Bucket().setName("bucket"));
+
+ verify(mockDelegateV2).removeBucket(BucketInfo.of("bucket"));
+ Mockito.verifyNoMoreInteractions(mockDelegate);
+ }
+
+ @Test
+ public void testBucketOwnerIsRoutedToV2BucketProject() throws IOException {
+ GcsUtil gcsUtil = gcsUtilRoutingToV2();
+ GcsPath path = GcsPath.fromUri("gs://bucket/object");
+ when(mockDelegateV2.bucketProject(path)).thenReturn(123L);
+
+ assertEquals(123L, gcsUtil.bucketOwner(path));
+ Mockito.verifyNoMoreInteractions(mockDelegate);
+ }
+
+ @Test
+ public void testGetObjectIsRoutedToV2AndKeepsAllFields() throws IOException {
+ GcsUtil gcsUtil = gcsUtilRoutingToV2();
+ GcsPath path = GcsPath.fromUri("gs://bucket/object");
+ com.google.cloud.storage.Blob blob = mockBlob("bucket", "object");
+ when(blob.getSize()).thenReturn(42L);
+ when(blob.getGeneration()).thenReturn(7L);
+ when(blob.getMetageneration()).thenReturn(3L);
+ when(blob.getContentType()).thenReturn("text/csv");
+ when(blob.getContentEncoding()).thenReturn("gzip");
+ when(blob.getMd5()).thenReturn("md5==");
+ when(blob.getCrc32c()).thenReturn("crc==");
+ when(blob.getEtag()).thenReturn("etag");
+ when(blob.getUpdateTimeOffsetDateTime())
+
.thenReturn(java.time.Instant.ofEpochMilli(1234L).atOffset(java.time.ZoneOffset.UTC));
+ when(blob.getCreateTimeOffsetDateTime())
+
.thenReturn(java.time.Instant.ofEpochMilli(1000L).atOffset(java.time.ZoneOffset.UTC));
+ when(mockDelegateV2.getBlob(path)).thenReturn(blob);
+
+ StorageObject object = gcsUtil.getObject(path);
+
+ assertEquals("bucket", object.getBucket());
+ assertEquals("object", object.getName());
+ assertEquals(BigInteger.valueOf(42L), object.getSize());
+ assertEquals(Long.valueOf(7L), object.getGeneration());
+ assertEquals(Long.valueOf(3L), object.getMetageneration());
+ assertEquals("text/csv", object.getContentType());
+ assertEquals("gzip", object.getContentEncoding());
+ assertEquals("md5==", object.getMd5Hash());
+ assertEquals("crc==", object.getCrc32c());
+ assertEquals("etag", object.getEtag());
+ assertEquals(1234L, object.getUpdated().getValue());
+ assertEquals(1000L, object.getTimeCreated().getValue());
+ Mockito.verifyNoMoreInteractions(mockDelegate);
+ }
+
+ /** Fields a blob may not carry (e.g. when fetched with a field mask) are
left unset. */
+ @Test
+ public void testGetObjectLeavesMissingFieldsUnsetForV2() throws IOException {
+ GcsUtil gcsUtil = gcsUtilRoutingToV2();
+ GcsPath path = GcsPath.fromUri("gs://bucket/object");
+ com.google.cloud.storage.Blob blob = mockBlob("bucket", "object");
+ // Mockito would otherwise answer 0 for the boxed size.
+ when(blob.getSize()).thenReturn(null);
+ when(mockDelegateV2.getBlob(path)).thenReturn(blob);
+
+ StorageObject object = gcsUtil.getObject(path);
+
+ assertEquals("object", object.getName());
+ assertNull(object.getSize());
+ assertNull(object.getUpdated());
+ assertNull(object.getTimeCreated());
+ }
+
+ @Test
+ public void testCreateWithTypeIsRoutedToV2() throws IOException {
+ GcsUtil gcsUtil = gcsUtilRoutingToV2();
+ GcsPath path = GcsPath.fromUri("gs://bucket/object");
+
+ gcsUtil.create(path, "text/plain");
+ gcsUtil.create(path, "text/plain", 1024);
+
+ verify(mockDelegateV2)
+ .create(path,
GcsUtilV1.CreateOptions.builder().setContentType("text/plain").build());
+ verify(mockDelegateV2)
+ .create(
+ path,
+ GcsUtilV1.CreateOptions.builder()
+ .setContentType("text/plain")
+ .setUploadBufferSizeBytes(1024)
+ .build());
+ Mockito.verifyNoMoreInteractions(mockDelegate);
+ }
+
+ /** GcpOptions reads the soft delete policy of the temp bucket through this
method. */
+ @Test
+ public void testGetBucketIsRoutedToV2AndConvertsTheBucket() throws
IOException {
+ GcsUtil gcsUtil = gcsUtilRoutingToV2();
+ GcsPath path = GcsPath.fromUri("gs://bucket/object");
+ com.google.cloud.storage.Bucket bucket =
Mockito.mock(com.google.cloud.storage.Bucket.class);
+ when(bucket.getName()).thenReturn("bucket");
+ when(bucket.getLocation()).thenReturn("US-CENTRAL1");
+ when(bucket.getProject()).thenReturn(BigInteger.valueOf(123L));
+ when(bucket.getStorageClass()).thenReturn(StorageClass.NEARLINE);
+ when(bucket.getSoftDeletePolicy())
+ .thenReturn(
+ BucketInfo.SoftDeletePolicy.newBuilder()
+ .setRetentionDuration(java.time.Duration.ofDays(7))
+ .build());
+ when(mockDelegateV2.getBucket(path)).thenReturn(bucket);
+
+ Bucket converted = gcsUtil.getBucket(path);
+
+ assertNotNull(converted);
+ assertEquals("bucket", converted.getName());
+ assertEquals("US-CENTRAL1", converted.getLocation());
+ assertEquals(BigInteger.valueOf(123L), converted.getProjectNumber());
+ assertEquals("NEARLINE", converted.getStorageClass());
+ assertEquals(
+ Long.valueOf(java.time.Duration.ofDays(7).getSeconds()),
+ converted.getSoftDeletePolicy().getRetentionDurationSeconds());
+ Mockito.verifyNoMoreInteractions(mockDelegate);
+ }
+
+ @Test
+ public void testGetBucketWithoutSoftDeletePolicyForV2() throws IOException {
+ GcsUtil gcsUtil = gcsUtilRoutingToV2();
+ GcsPath path = GcsPath.fromUri("gs://bucket/object");
+ com.google.cloud.storage.Bucket bucket =
Mockito.mock(com.google.cloud.storage.Bucket.class);
+ when(bucket.getName()).thenReturn("bucket");
+ when(mockDelegateV2.getBucket(path)).thenReturn(bucket);
+
+ Bucket converted = gcsUtil.getBucket(path);
+
+ assertNotNull(converted);
+ assertNull(converted.getSoftDeletePolicy());
+ assertNull(converted.getStorageClass());
+ }
+
+ /** Without the use_gcsutil_v2 experiment, the V2-only methods fail rather
than fall back. */
+ @Test
+ public void testV2OnlyMethodsFailWithoutV2() {
+ GcsUtil gcsUtil = gcsOptionsWithTestCredential().getGcsUtil();
+ assertNull(gcsUtil.delegateV2);
+ GcsUtilV1 v1 = Mockito.mock(GcsUtilV1.class);
+ gcsUtil.delegate = v1;
+ GcsPath path = GcsPath.fromUri("gs://bucket/object");
+ List<GcsPath> paths = ImmutableList.of(path);
+
+ List<ThrowingRunnable> calls =
+ ImmutableList.of(
+ () -> gcsUtil.getBlob(path),
+ () -> gcsUtil.getBlobs(paths),
+ () -> gcsUtil.listBlobs("bucket", "prefix", null),
+ () -> gcsUtil.listBlobs("bucket", "prefix", null, "/"),
+ () -> gcsUtil.openV2(path),
+ () -> gcsUtil.createV2(path,
GcsUtil.CreateOptions.builder().build()),
+ () -> gcsUtil.createBucket(BucketInfo.of("bucket")),
+ () -> gcsUtil.getBucketWithOptions(path),
+ () -> gcsUtil.removeBucket(BucketInfo.of("bucket")),
+ () -> gcsUtil.copyV2(paths, paths),
+ () -> gcsUtil.copy(paths, paths,
GcsUtilV2.OverwriteStrategy.ALWAYS_OVERWRITE),
+ () -> gcsUtil.renameV2(paths, paths),
+ () ->
+ gcsUtil.rename(
+ paths,
+ paths,
+ GcsUtilV2.MissingStrategy.FAIL_IF_MISSING,
+ GcsUtilV2.OverwriteStrategy.ALWAYS_OVERWRITE),
+ () -> gcsUtil.removeV2(paths),
+ () -> gcsUtil.remove(paths,
GcsUtilV2.MissingStrategy.FAIL_IF_MISSING));
+
+ for (ThrowingRunnable call : calls) {
+ IOException e = assertThrows(IOException.class, call);
+ assertEquals("GcsUtil V2 not initialized.", e.getMessage());
+ }
+ Mockito.verifyNoInteractions(v1);
+ }
+
+ @Test
+ public void testExpandIsRoutedToV2() throws IOException {
+ GcsUtil gcsUtil = gcsUtilRoutingToV2();
+ GcsPath pattern = GcsPath.fromUri("gs://bucket/prefix/*");
+ List<GcsPath> expanded =
ImmutableList.of(GcsPath.fromUri("gs://bucket/prefix/a"));
+ when(mockDelegateV2.expand(pattern)).thenReturn(expanded);
+
+ assertSame(expanded, gcsUtil.expand(pattern));
+ Mockito.verifyNoMoreInteractions(mockDelegate);
+ }
+
+ @Test
+ public void testFileSizeIsRoutedToV2() throws IOException {
+ GcsUtil gcsUtil = gcsUtilRoutingToV2();
+ GcsPath path = GcsPath.fromUri("gs://bucket/object");
+ when(mockDelegateV2.fileSize(path)).thenReturn(42L);
+
+ assertEquals(42L, gcsUtil.fileSize(path));
+ Mockito.verifyNoMoreInteractions(mockDelegate);
+ }
+
+ /**
+ * Only the routing of the delimiter overload is checked. The page
conversion is covered by {@link
+ * #testListObjectsIsRoutedToV2AndConvertsAPage}.
+ */
+ @Test
+ public void testListObjectsWithDelimiterIsRoutedToV2() throws IOException {
+ GcsUtil gcsUtil = gcsUtilRoutingToV2();
+ com.google.cloud.storage.Blob object = mockBlob("bucket", "prefix/object");
+ @SuppressWarnings("unchecked")
+ Page<com.google.cloud.storage.Blob> page = Mockito.mock(Page.class);
+ when(page.getValues()).thenReturn(ImmutableList.of(object));
+ when(mockDelegateV2.listBlobs("bucket", "prefix/", "token",
"/")).thenReturn(page);
+
+ Objects objects = gcsUtil.listObjects("bucket", "prefix/", "token", "/");
+
+ assertEquals("prefix/object", objects.getItems().get(0).getName());
+ Mockito.verifyNoMoreInteractions(mockDelegate);
+ }
+
+ @Test
+ public void testOpenIsRoutedToV2() throws IOException {
+ GcsUtil gcsUtil = gcsUtilRoutingToV2();
+ GcsPath path = GcsPath.fromUri("gs://bucket/object");
+ SeekableByteChannel channel = Mockito.mock(SeekableByteChannel.class);
+ when(mockDelegateV2.open(path)).thenReturn(channel);
+
+ assertSame(channel, gcsUtil.open(path));
+ Mockito.verifyNoMoreInteractions(mockDelegate);
+ }
+
+ @Test
+ public void testVerifyBucketAccessibleIsRoutedToV2() throws IOException {
+ GcsUtil gcsUtil = gcsUtilRoutingToV2();
+ GcsPath path = GcsPath.fromUri("gs://bucket/object");
+
+ gcsUtil.verifyBucketAccessible(path);
+
+ verify(mockDelegateV2).verifyBucketAccessible(path);
+ Mockito.verifyNoMoreInteractions(mockDelegate);
+ }
+
+ @Test
+ public void testBucketAccessibleIsRoutedToV2() throws IOException {
+ GcsUtil gcsUtil = gcsUtilRoutingToV2();
+ GcsPath accessible = GcsPath.fromUri("gs://accessible/object");
+ GcsPath inaccessible = GcsPath.fromUri("gs://inaccessible/object");
+ when(mockDelegateV2.bucketAccessible(accessible)).thenReturn(true);
+ when(mockDelegateV2.bucketAccessible(inaccessible)).thenReturn(false);
+
+ assertTrue(gcsUtil.bucketAccessible(accessible));
+ assertFalse(gcsUtil.bucketAccessible(inaccessible));
+ Mockito.verifyNoMoreInteractions(mockDelegate);
+ }
+
+ // The tests below exercise a real GcsUtilV2 delegate whose java-storage
client is mocked, to
+ // cover behavior that GcsUtilV2 must share with GcsUtilV1.
+ // TODO: Move these to a parity test that runs against both delegates.
+
+ /**
+ * Returns a {@link GcsUtil} backed by a real {@link GcsUtilV2} that issues
every call to {@code
+ * storage}. Performance metrics are off by default, so the per-operation
clients of {@link
+ * GcsUtilV2#storageWithHttpMetrics} resolve to this one as well.
+ */
+ private GcsUtil gcsUtilWithV2Storage(com.google.cloud.storage.Storage
storage) {
+ GcsOptions options = gcsOptionsWithTestCredential();
+ options.setProject("my_project");
+ GcsUtil gcsUtil = options.getGcsUtil();
+ GcsUtilV2 delegateV2 = Mockito.spy(new GcsUtilV2(options));
+ Mockito.doReturn(storage).when(delegateV2).storage();
+ gcsUtil.delegateV2 = delegateV2;
+ return gcsUtil;
+ }
+
+ /**
+ * Mirrors {@link #testGCSReadMetricsIsSet} for V2: opening a missing object
records a {@code
+ * not_found} request, even though V2 detects it through a lookup rather
than a {@link
+ * StorageException}.
+ */
+ @Test
+ public void testV2OpenMissingObjectRecordsNotFoundMetric() {
+ com.google.cloud.storage.Storage storage =
Mockito.mock(com.google.cloud.storage.Storage.class);
+ // An unstubbed get() returns null, which is how java-storage reports a
missing object.
+ GcsUtil gcsUtil = gcsUtilWithV2Storage(storage);
+
+ assertThrows(
+ FileNotFoundException.class,
+ () -> gcsUtil.open(GcsPath.fromComponents("testbucket",
"testobject")));
+
+ verifyMetricWasSet("my_project", "testbucket", "GcsGet", "not_found", 1);
+ verifyMetricWasSet("my_project", "testbucket", "GcsGet", "ok", 0);
+ }
+
+ /**
+ * An upload is finalized when its channel is closed, so a failed
precondition surfaces there. V2
+ * must report it as an {@link IOException}, as V1 does, rather than an
unchecked {@link
+ * StorageException}.
+ */
+ @Test
+ public void testV2WriteChannelCloseTranslatesStorageException() throws
IOException {
+ com.google.cloud.storage.Storage storage =
Mockito.mock(com.google.cloud.storage.Storage.class);
+ WriteChannel writer = Mockito.mock(WriteChannel.class);
+ StorageException preconditionFailed = new StorageException(412,
"Precondition Failed");
+ Mockito.doThrow(preconditionFailed).when(writer).close();
+ when(storage.writer(any(com.google.cloud.storage.BlobInfo.class),
any())).thenReturn(writer);
+ GcsUtil gcsUtil = gcsUtilWithV2Storage(storage);
+
+ WritableByteChannel channel =
+ gcsUtil.create(
+ GcsPath.fromComponents("testbucket", "testobject"),
+ CreateOptions.builder().setExpectFileToNotExist(true).build());
+
+ IOException thrown = assertThrows(IOException.class, channel::close);
+ assertSame(preconditionFailed, thrown.getCause());
+ }
+
+ /** Mirrors {@link #testBucketDoesNotExist} for V2. */
+ @Test
+ public void testV2BucketAccessibleIsFalseWhenBucketDoesNotExist() throws
IOException {
+ com.google.cloud.storage.Storage storage =
Mockito.mock(com.google.cloud.storage.Storage.class);
+ // An unstubbed get() returns null, which is how java-storage reports a
missing bucket.
+ GcsUtil gcsUtil = gcsUtilWithV2Storage(storage);
+
+ assertFalse(gcsUtil.bucketAccessible(GcsPath.fromComponents("testbucket",
"testobject")));
+ }
+
+ /** Mirrors {@link #testBucketDoesNotExistBecauseOfAccessError} for V2. */
+ @Test
+ public void testV2BucketAccessibleIsFalseWhenAccessIsDenied() throws
IOException {
+ com.google.cloud.storage.Storage storage =
Mockito.mock(com.google.cloud.storage.Storage.class);
+ when(storage.get(Mockito.eq("testbucket"), any(BucketGetOption.class)))
+ .thenThrow(new StorageException(403, "Forbidden"));
+ GcsUtil gcsUtil = gcsUtilWithV2Storage(storage);
+
+ assertFalse(gcsUtil.bucketAccessible(GcsPath.fromComponents("testbucket",
"testobject")));
+ }
+
+ /**
+ * Any other failure (e.g. a 5xx) says nothing about whether the bucket is
accessible, so it must
+ * propagate rather than be reported as an inaccessible bucket, as V1 does.
+ */
+ @Test
+ public void testV2BucketAccessiblePropagatesOtherFailures() {
+ com.google.cloud.storage.Storage storage =
Mockito.mock(com.google.cloud.storage.Storage.class);
+ StorageException serverError = new StorageException(503, "Service
Unavailable");
+ when(storage.get(Mockito.eq("testbucket"),
any(BucketGetOption.class))).thenThrow(serverError);
+ GcsUtil gcsUtil = gcsUtilWithV2Storage(storage);
+
+ IOException thrown =
+ assertThrows(
+ IOException.class,
+ () ->
gcsUtil.bucketAccessible(GcsPath.fromComponents("testbucket", "testobject")));
+ assertSame(serverError, thrown.getCause());
+ }
}
diff --git
a/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV2Test.java
b/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV2Test.java
index 2b9e55b98b8..8bd79860cc2 100644
---
a/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV2Test.java
+++
b/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV2Test.java
@@ -30,6 +30,7 @@ import com.google.api.client.testing.http.MockHttpTransport;
import com.google.api.client.testing.http.MockLowLevelHttpResponse;
import com.google.auth.Credentials;
import com.google.cloud.NoCredentials;
+import com.google.cloud.hadoop.util.AsyncWriteChannelOptions;
import com.google.cloud.http.HttpTransportOptions;
import com.google.cloud.storage.Storage;
import com.google.cloud.storage.StorageOptions;
@@ -415,4 +416,15 @@ public class GcsUtilV2Test {
assertNull(gcsCounter(other, "gcs_http_read_wire_bytes_received"));
assertNull(gcsCounter(other, "gcs_http_write_wire_bytes_sent"));
}
+
+ /**
+ * The chunk size decides how many requests a write costs, and the two
clients do not default to
+ * the same one, so a drift here is a silent throughput regression rather
than a test failure.
+ */
+ @Test
+ public void testDefaultUploadChunkSizeMatchesV1() {
+ assertEquals(
+ AsyncWriteChannelOptions.DEFAULT.getUploadChunkSize(),
+ GcsUtilV2.DEFAULT_UPLOAD_CHUNK_SIZE_BYTES);
+ }
}