From 2d96274409b61b62616b4cdff9ca8cf58c40058c Mon Sep 17 00:00:00 2001 From: Shunping Huang Date: Mon, 10 Aug 2026 10:19:07 -0400 Subject: [PATCH 1/4] Revert "Fix gcs endpoint pipeline option wiring (#39435)" This reverts commit adbb11ddb465fc68ff6d09351ba1b5fe7c1ba12f. --- .../beam/sdk/extensions/gcp/util/GcsUtilV1.java | 17 +++-------------- .../sdk/extensions/gcp/util/GcsUtilTest.java | 13 ------------- 2 files changed, 3 insertions(+), 27 deletions(-) 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 cfeb12dcae5c..1ad08f0ba1a2 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 @@ -261,18 +261,12 @@ public boolean shouldRetry(IOException e) { this.credentials = credentials; this.maxBytesRewrittenPerCall = null; this.numRewriteTokensUsed = null; - GoogleCloudStorageOptions.Builder optionsBuilder = + googleCloudStorageOptions = GoogleCloudStorageOptions.builder() .setAppName("Beam") .setReadChannelOptions(gcsReadOptions) - .setGrpcEnabled(shouldUseGrpc); - if (storageClient.getRootUrl() != null) { - optionsBuilder.setStorageRootUrl(storageClient.getRootUrl()); - } - if (storageClient.getServicePath() != null) { - optionsBuilder.setStorageServicePath(storageClient.getServicePath()); - } - googleCloudStorageOptions = optionsBuilder.build(); + .setGrpcEnabled(shouldUseGrpc) + .build(); try { googleCloudStorage = createGoogleCloudStorage(googleCloudStorageOptions, storageClient, credentials); @@ -501,11 +495,6 @@ private Long toFileSize(StorageObjectOrIOException storageObjectOrIOException) } } - @VisibleForTesting - GoogleCloudStorage getGoogleCloudStorage() { - return googleCloudStorage; - } - @VisibleForTesting void setCloudStorageImpl(GoogleCloudStorage g) { googleCloudStorage = g; 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 2f77f15dcffc..d32ca162e3fd 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 @@ -63,7 +63,6 @@ import com.google.auth.Credentials; 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; @@ -1867,18 +1866,6 @@ public void testReadMetricsAreNotCollectedWhenNotEnabledOpenWithOptions() throws testReadMetrics(false, GoogleCloudStorageReadOptions.DEFAULT); } - @Test - public void testGcsEndpoint() throws IOException { - GcsOptions pipelineOptions = PipelineOptionsFactory.as(GcsOptions.class); - pipelineOptions.setGcsEndpoint("http://localhost:4443/storage/v1/"); - - GcsUtil gcsUtil = pipelineOptions.getGcsUtil(); - GoogleCloudStorageImpl gcsImpl = - (GoogleCloudStorageImpl) gcsUtil.delegate.getGoogleCloudStorage(); - assertEquals("http://localhost:4443/", gcsImpl.getOptions().getStorageRootUrl()); - assertEquals("storage/v1/", gcsImpl.getOptions().getStorageServicePath()); - } - /** A helper to wrap a {@link GenericJson} object in a content stream. */ private static InputStream toStream(String content) throws IOException { return new ByteArrayInputStream(content.getBytes(StandardCharsets.UTF_8)); From 4e9e8311e8a3f50878ff5ee0f57aa198e38bcbc1 Mon Sep 17 00:00:00 2001 From: Shunping Huang Date: Mon, 10 Aug 2026 10:19:24 -0400 Subject: [PATCH 2/4] Revert "Suppress log spams in gcsio 3.0 (#38588)" This reverts commit 9d307e559eb078e8cccfd3c1e411f74563a33876. --- .../apache/beam/sdk/extensions/gcp/util/GcsUtilV1.java | 9 --------- 1 file changed, 9 deletions(-) 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 1ad08f0ba1a2..97778ac4e1df 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 @@ -71,7 +71,6 @@ import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.TimeUnit; -import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; import java.util.function.Consumer; import java.util.function.Supplier; @@ -185,7 +184,6 @@ public boolean shouldRetry(IOException e) { return RetryDeterminer.SOCKET_ERRORS.shouldRetry(e); } }; - private static final AtomicBoolean overwriteLog = new AtomicBoolean(false); ///////////////////////////////////////////////////////////////////////////// @@ -728,16 +726,9 @@ public WritableByteChannel create(GcsPath path, CreateOptions options) throws IO } } - @SuppressFBWarnings("LG_LOST_LOGGER_DUE_TO_WEAK_REFERENCE") GoogleCloudStorage createGoogleCloudStorage( GoogleCloudStorageOptions options, Storage storage, Credentials credentials) 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") - .setLevel(java.util.logging.Level.SEVERE); - } - return GoogleCloudStorageImpl.builder() .setOptions(options) .setHttpTransport(storage.getRequestFactory().getTransport()) From 0d040066d3075c125c66db80ed969fd3f90e98a3 Mon Sep 17 00:00:00 2001 From: Shunping Huang Date: Mon, 10 Aug 2026 10:19:42 -0400 Subject: [PATCH 3/4] Revert "Upgrade gcsio to 3.1.16 (#38419)" This reverts commit 89dde7c7a1b3deb4b2ee9d56b29fb3d5e3de6aab. --- .../beam/gradle/BeamModulePlugin.groovy | 6 +- .../sdk/extensions/gcp/util/GcsUtilV1.java | 62 ++++++++++++++----- .../sdk/extensions/gcp/util/GcsUtilTest.java | 19 +++--- 3 files changed, 56 insertions(+), 31 deletions(-) diff --git a/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy b/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy index 225201a5c0d8..847b2196afe9 100644 --- a/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy +++ b/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy @@ -623,7 +623,7 @@ class BeamModulePlugin implements Plugin { def gax_version = "2.82.0" def google_ads_version = "33.0.0" def google_clients_version = "2.0.0" - def google_cloud_bigdataoss_version = "3.1.16" + def google_cloud_bigdataoss_version = "2.2.26" def google_code_gson_version = "2.10.1" def google_oauth_clients_version = "1.34.1" // [bomupgrader] determined by: io.grpc:grpc-netty, consistent with: google_cloud_platform_libraries_bom @@ -712,9 +712,9 @@ class BeamModulePlugin implements Plugin { aws_java_sdk2_profiles : "software.amazon.awssdk:profiles:$aws_java_sdk2_version", azure_sdk_bom : "com.azure:azure-sdk-bom:1.2.14", bigdataoss_gcsio : "com.google.cloud.bigdataoss:gcsio:$google_cloud_bigdataoss_version", - bigdataoss_gcs_connector : "com.google.cloud.bigdataoss:gcs-connector:$google_cloud_bigdataoss_version", + bigdataoss_gcs_connector : "com.google.cloud.bigdataoss:gcs-connector:hadoop2-$google_cloud_bigdataoss_version", bigdataoss_util : "com.google.cloud.bigdataoss:util:$google_cloud_bigdataoss_version", - bigdataoss_util_hadoop : "com.google.cloud.bigdataoss:util-hadoop:$google_cloud_bigdataoss_version", + bigdataoss_util_hadoop : "com.google.cloud.bigdataoss:util-hadoop:hadoop2-$google_cloud_bigdataoss_version", byte_buddy : "net.bytebuddy:byte-buddy:1.17.7", cassandra_driver_core : "com.datastax.cassandra:cassandra-driver-core:$cassandra_driver_version", cassandra_driver_mapping : "com.datastax.cassandra:cassandra-driver-mapping:$cassandra_driver_version", 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 97778ac4e1df..1ade4be6fdb5 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 @@ -30,6 +30,7 @@ import com.google.api.client.http.HttpHeaders; import com.google.api.client.http.HttpRequestInitializer; import com.google.api.client.http.HttpStatusCodes; +import com.google.api.client.http.HttpTransport; import com.google.api.client.util.BackOff; import com.google.api.client.util.Sleeper; import com.google.api.services.storage.Storage; @@ -52,6 +53,7 @@ import edu.umd.cs.findbugs.annotations.SuppressFBWarnings; import java.io.FileNotFoundException; import java.io.IOException; +import java.lang.reflect.Method; import java.nio.channels.SeekableByteChannel; import java.nio.channels.WritableByteChannel; import java.nio.file.AccessDeniedException; @@ -265,12 +267,8 @@ public boolean shouldRetry(IOException e) { .setReadChannelOptions(gcsReadOptions) .setGrpcEnabled(shouldUseGrpc) .build(); - try { - googleCloudStorage = - createGoogleCloudStorage(googleCloudStorageOptions, storageClient, credentials); - } catch (IOException e) { - throw new RuntimeException(e); - } + googleCloudStorage = + createGoogleCloudStorage(googleCloudStorageOptions, storageClient, credentials); this.batchRequestSupplier = () -> { // Capture reference to this so that the most recent storageClient and initializer @@ -727,16 +725,48 @@ public WritableByteChannel create(GcsPath path, CreateOptions options) throws IO } GoogleCloudStorage createGoogleCloudStorage( - GoogleCloudStorageOptions options, Storage storage, Credentials credentials) - throws IOException { - return GoogleCloudStorageImpl.builder() - .setOptions(options) - .setHttpTransport(storage.getRequestFactory().getTransport()) - .setCredentials(credentials) - // gcsio 3 expects httpRequestInitializer to be either absent or - // com.google.cloud.hadoop.util.RetryHttpInitializer when credentials not provided - .setHttpRequestInitializer(credentials != null ? httpRequestInitializer : null) - .build(); + GoogleCloudStorageOptions options, Storage storage, Credentials credentials) { + try { + return new GoogleCloudStorageImpl(options, storage, credentials); + } catch (NoSuchMethodError e) { + // gcs-connector 3.x drops the direct constructor and exclusively uses Builder + // TODO eliminate reflection once Beam drops Java 8 support and upgrades to gcsio 3.x + try { + final Method builderMethod = GoogleCloudStorageImpl.class.getMethod("builder"); + Object builder = builderMethod.invoke(null); + final Class builderClass = + Class.forName( + "com.google.cloud.hadoop.gcsio.AutoBuilder_GoogleCloudStorageImpl_Builder"); + + final Method setOptionsMethod = + builderClass.getMethod("setOptions", GoogleCloudStorageOptions.class); + setOptionsMethod.setAccessible(true); + builder = setOptionsMethod.invoke(builder, options); + + final Method setHttpTransportMethod = + builderClass.getMethod("setHttpTransport", HttpTransport.class); + setHttpTransportMethod.setAccessible(true); + builder = + setHttpTransportMethod.invoke(builder, storage.getRequestFactory().getTransport()); + + final Method setCredentialsMethod = + builderClass.getMethod("setCredentials", Credentials.class); + setCredentialsMethod.setAccessible(true); + builder = setCredentialsMethod.invoke(builder, credentials); + + final Method setHttpRequestInitializerMethod = + builderClass.getMethod("setHttpRequestInitializer", HttpRequestInitializer.class); + setHttpRequestInitializerMethod.setAccessible(true); + builder = setHttpRequestInitializerMethod.invoke(builder, httpRequestInitializer); + + final Method buildMethod = builderClass.getMethod("build"); + buildMethod.setAccessible(true); + return (GoogleCloudStorage) buildMethod.invoke(builder); + } catch (Exception reflectionError) { + throw new RuntimeException( + "Failed to construct GoogleCloudStorageImpl from gcsio 3.x Builder", reflectionError); + } + } } /** 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 d32ca162e3fd..a2b0e0af502b 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 @@ -184,8 +184,8 @@ public void testCreationWithExplicitGoogleCloudStorageReadOptions() throws Excep GoogleCloudStorageReadOptions readOptions = GoogleCloudStorageReadOptions.builder() .setFadvise(GoogleCloudStorageReadOptions.Fadvise.AUTO) - .setGzipEncodingSupportEnabled(true) - .setFastFailOnNotFoundEnabled(false) + .setSupportGzipEncoding(true) + .setFastFailOnNotFound(false) .build(); GcsOptions pipelineOptions = PipelineOptionsFactory.as(GcsOptions.class); @@ -193,10 +193,7 @@ public void testCreationWithExplicitGoogleCloudStorageReadOptions() throws Excep GcsUtil gcsUtil = pipelineOptions.getGcsUtil(); GoogleCloudStorage googleCloudStorageMock = Mockito.spy(GoogleCloudStorage.class); - Mockito.when( - googleCloudStorageMock.open( - Mockito.any(StorageResourceId.class), - Mockito.any(GoogleCloudStorageReadOptions.class))) + Mockito.when(googleCloudStorageMock.open(Mockito.any(), Mockito.any())) .thenReturn(Mockito.mock(SeekableByteChannel.class)); gcsUtil.delegate.setCloudStorageImpl(googleCloudStorageMock); @@ -1009,7 +1006,7 @@ public void testGCSChannelCloseIdempotent() throws IOException { GcsOptions pipelineOptions = gcsOptionsWithTestCredential(); GcsUtil gcsUtil = pipelineOptions.getGcsUtil(); GoogleCloudStorageReadOptions readOptions = - GoogleCloudStorageReadOptions.builder().setFastFailOnNotFoundEnabled(false).build(); + GoogleCloudStorageReadOptions.builder().setFastFailOnNotFound(false).build(); gcsUtil.delegate.setCloudStorageImpl( GoogleCloudStorageOptions.builder() @@ -1029,7 +1026,7 @@ public void testGCSReadMetricsIsSet() { GcsOptions pipelineOptions = gcsOptionsWithTestCredential(); GcsUtil gcsUtil = pipelineOptions.getGcsUtil(); GoogleCloudStorageReadOptions readOptions = - GoogleCloudStorageReadOptions.builder().setFastFailOnNotFoundEnabled(true).build(); + GoogleCloudStorageReadOptions.builder().setFastFailOnNotFound(true).build(); gcsUtil.delegate.setCloudStorageImpl( GoogleCloudStorageOptions.builder() .setAppName("Beam") @@ -1676,10 +1673,8 @@ public static GcsUtilV1Mock createMockWithMockStorage( .thenReturn(Channels.newChannel(new ByteArrayOutputStream())); } else { SeekableByteChannel seekableByteChannel = new SeekableInMemoryByteChannel(readPayload); - Mockito.when(googleCloudStorageMock.open(Mockito.any(StorageResourceId.class))) - .thenReturn(seekableByteChannel); - Mockito.when( - googleCloudStorageMock.open(Mockito.any(StorageResourceId.class), Mockito.any())) + Mockito.when(googleCloudStorageMock.open(Mockito.any())).thenReturn(seekableByteChannel); + Mockito.when(googleCloudStorageMock.open(Mockito.any(), Mockito.any())) .thenReturn(seekableByteChannel); } return gcsUtilMock; From f1c6de3cf3941e6712780d7680739d485fb2de45 Mon Sep 17 00:00:00 2001 From: Shunping Huang Date: Mon, 10 Aug 2026 10:38:30 -0400 Subject: [PATCH 4/4] Revert "Override default fadvise to fix regression from gcs-connector v3 upgrade (#39445)" This reverts commit eb9e7a3dacdeea074b22d833cb1fda789d09d6f4. --- .../google-cloud-platform-core/build.gradle | 1 - .../extensions/gcp/options/GcsOptions.java | 85 +------------------ .../extensions/gcp/GcpCoreApiSurfaceTest.java | 2 - .../gcp/options/GcsOptionsTest.java | 47 ---------- 4 files changed, 2 insertions(+), 133 deletions(-) diff --git a/sdks/java/extensions/google-cloud-platform-core/build.gradle b/sdks/java/extensions/google-cloud-platform-core/build.gradle index 78cfe4739ec5..f1bfb63c7a3a 100644 --- a/sdks/java/extensions/google-cloud-platform-core/build.gradle +++ b/sdks/java/extensions/google-cloud-platform-core/build.gradle @@ -56,7 +56,6 @@ dependencies { implementation library.java.http_core implementation library.java.http_client implementation library.java.jackson_annotations - implementation library.java.jackson_core implementation library.java.jackson_databind permitUnusedDeclared library.java.jackson_databind // BEAM-11761 testImplementation project(path: ":sdks:java:core", configuration: "shadowTest") diff --git a/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/options/GcsOptions.java b/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/options/GcsOptions.java index 134c4cb3f281..2da382a5b674 100644 --- a/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/options/GcsOptions.java +++ b/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/options/GcsOptions.java @@ -18,15 +18,6 @@ package org.apache.beam.sdk.extensions.gcp.options; import com.fasterxml.jackson.annotation.JsonIgnore; -import com.fasterxml.jackson.core.JsonGenerator; -import com.fasterxml.jackson.core.JsonParser; -import com.fasterxml.jackson.databind.DeserializationContext; -import com.fasterxml.jackson.databind.JsonDeserializer; -import com.fasterxml.jackson.databind.JsonNode; -import com.fasterxml.jackson.databind.JsonSerializer; -import com.fasterxml.jackson.databind.SerializerProvider; -import com.fasterxml.jackson.databind.annotation.JsonDeserialize; -import com.fasterxml.jackson.databind.annotation.JsonSerialize; import com.google.cloud.hadoop.gcsio.GoogleCloudStorageReadOptions; import com.google.cloud.hadoop.util.AsyncWriteChannelOptions; import java.util.HashMap; @@ -58,16 +49,12 @@ public interface GcsOptions extends ApplicationNameOptions, GcpOptions, Pipeline class GcsReadOptionsFactory implements DefaultValueFactory { @Override public GoogleCloudStorageReadOptions create(PipelineOptions options) { - // In gcs-connector v3, GoogleCloudStorageReadOptions.DEFAULT changed fadvise from SEQUENTIAL - // to AUTO. Beam workloads default to SEQUENTIAL to preserve expected sequential read - // throughput and caching behavior. - return GcsReadOptionsSerializer.DEFAULT_OPTIONS; + return GoogleCloudStorageReadOptions.DEFAULT; } } /** @deprecated This option will be removed in a future release. */ - @JsonSerialize(using = GcsReadOptionsSerializer.class) - @JsonDeserialize(using = GcsReadOptionsDeserializer.class) + @JsonIgnore @Description( "The GoogleCloudStorageReadOptions instance that should be used to read from Google Cloud Storage.") @Default.InstanceFactory(GcsReadOptionsFactory.class) @@ -299,71 +286,3 @@ boolean exceedsEntryLimit() { } } } - -class GcsReadOptionsSerializer extends JsonSerializer { - static final GoogleCloudStorageReadOptions DEFAULT_OPTIONS = - GoogleCloudStorageReadOptions.DEFAULT - .toBuilder() - .setFadvise(GoogleCloudStorageReadOptions.Fadvise.SEQUENTIAL) - .build(); - - @Override - public void serialize( - GoogleCloudStorageReadOptions value, JsonGenerator gen, SerializerProvider serializers) - throws java.io.IOException { - // Note: We only support a partial set of options to propagate to remote - // workers. Setting the unsupported ones will not have any effect and will fall - // back to default in remote workers. Support for additional options can be - // added on a need basis. The full list of options can be seen in - // com.google.cloud.hadoop.gcsio.GoogleCloudStorageReadOptions. - gen.writeStartObject(); - if (value.getFadvise() != null && value.getFadvise() != DEFAULT_OPTIONS.getFadvise()) { - gen.writeStringField("fadvise", value.getFadvise().name()); - } - if (value.isFastFailOnNotFoundEnabled() != DEFAULT_OPTIONS.isFastFailOnNotFoundEnabled()) { - gen.writeBooleanField("fastFailOnNotFoundEnabled", value.isFastFailOnNotFoundEnabled()); - } - if (value.getMinRangeRequestSize() != DEFAULT_OPTIONS.getMinRangeRequestSize()) { - gen.writeNumberField("minRangeRequestSize", value.getMinRangeRequestSize()); - } - if (value.getInplaceSeekLimit() != DEFAULT_OPTIONS.getInplaceSeekLimit()) { - gen.writeNumberField("inplaceSeekLimit", value.getInplaceSeekLimit()); - } - if (value.isGrpcReadZeroCopyEnabled() != DEFAULT_OPTIONS.isGrpcReadZeroCopyEnabled()) { - gen.writeBooleanField("grpcReadZeroCopyEnabled", value.isGrpcReadZeroCopyEnabled()); - } - gen.writeEndObject(); - } -} - -class GcsReadOptionsDeserializer extends JsonDeserializer { - @Override - public GoogleCloudStorageReadOptions deserialize(JsonParser p, DeserializationContext ctxt) - throws java.io.IOException { - JsonNode root = p.readValueAsTree(); - GoogleCloudStorageReadOptions.Builder builder = - GcsReadOptionsSerializer.DEFAULT_OPTIONS.toBuilder(); - - if (root != null && root.isObject()) { - if (root.hasNonNull("fadvise")) { - builder.setFadvise( - GoogleCloudStorageReadOptions.Fadvise.valueOf(root.get("fadvise").asText())); - } - if (root.hasNonNull("fastFailOnNotFoundEnabled")) { - builder.setFastFailOnNotFoundEnabled(root.get("fastFailOnNotFoundEnabled").asBoolean()); - } else if (root.hasNonNull("fastFailOnNotFound")) { - builder.setFastFailOnNotFoundEnabled(root.get("fastFailOnNotFound").asBoolean()); - } - if (root.hasNonNull("minRangeRequestSize")) { - builder.setMinRangeRequestSize(root.get("minRangeRequestSize").asLong()); - } - if (root.hasNonNull("inplaceSeekLimit")) { - builder.setInplaceSeekLimit(root.get("inplaceSeekLimit").asLong()); - } - if (root.hasNonNull("grpcReadZeroCopyEnabled")) { - builder.setGrpcReadZeroCopyEnabled(root.get("grpcReadZeroCopyEnabled").asBoolean()); - } - } - return builder.build(); - } -} diff --git a/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/GcpCoreApiSurfaceTest.java b/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/GcpCoreApiSurfaceTest.java index cd51d4fd9d28..8af5e2260fc5 100644 --- a/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/GcpCoreApiSurfaceTest.java +++ b/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/GcpCoreApiSurfaceTest.java @@ -51,8 +51,6 @@ public void testGcpCoreApiSurface() throws Exception { final Set>> allowedClasses = ImmutableSet.of( classesInPackage("com.fasterxml.jackson.annotation"), - classesInPackage("com.fasterxml.jackson.core"), - classesInPackage("com.fasterxml.jackson.databind"), classesInPackage("com.google.api.client.googleapis"), classesInPackage("com.google.api.client.http"), classesInPackage("com.google.api.client.json"), diff --git a/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/options/GcsOptionsTest.java b/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/options/GcsOptionsTest.java index 912f0110e9ce..c499290b851d 100644 --- a/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/options/GcsOptionsTest.java +++ b/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/options/GcsOptionsTest.java @@ -18,16 +18,11 @@ package org.apache.beam.sdk.extensions.gcp.options; import static org.junit.Assert.assertEquals; -import static org.junit.Assert.assertFalse; -import static org.junit.Assert.assertNotNull; import static org.junit.Assert.assertThrows; -import com.fasterxml.jackson.databind.ObjectMapper; -import com.google.cloud.hadoop.gcsio.GoogleCloudStorageReadOptions; import java.util.Collections; import java.util.HashMap; import java.util.Map; -import org.apache.beam.sdk.options.PipelineOptions; import org.apache.beam.sdk.options.PipelineOptionsFactory; import org.junit.Test; import org.junit.runner.RunWith; @@ -80,46 +75,4 @@ public void testEntriesWithErrors() throws Exception { IllegalArgumentException.class, () -> PipelineOptionsFactory.fromArgs(TOO_MANY_ENTRIES_WITH_JOB).as(GcsOptions.class)); } - - @Test - public void testGoogleCloudStorageReadOptionsSerialization() throws Exception { - GcsOptions options = PipelineOptionsFactory.as(GcsOptions.class); - GoogleCloudStorageReadOptions readOptions = - GoogleCloudStorageReadOptions.builder() - .setFadvise(GoogleCloudStorageReadOptions.Fadvise.RANDOM) - .setFastFailOnNotFoundEnabled(false) - .setMinRangeRequestSize(12345L) - .build(); - options.setGoogleCloudStorageReadOptions(readOptions); - - ObjectMapper mapper = new ObjectMapper(); - String serialized = mapper.writeValueAsString(options); - GcsOptions deserialized = - mapper.readValue(serialized, PipelineOptions.class).as(GcsOptions.class); - - GoogleCloudStorageReadOptions deserializedReadOptions = - deserialized.getGoogleCloudStorageReadOptions(); - - assertNotNull(deserializedReadOptions); - assertEquals( - GoogleCloudStorageReadOptions.Fadvise.RANDOM, deserializedReadOptions.getFadvise()); - assertFalse(deserializedReadOptions.isFastFailOnNotFoundEnabled()); - assertEquals(12345L, deserializedReadOptions.getMinRangeRequestSize()); - } - - @Test - public void testDefaultGoogleCloudStorageReadOptionsSerialization() throws Exception { - GcsOptions options = PipelineOptionsFactory.as(GcsOptions.class); - ObjectMapper mapper = new ObjectMapper(); - String serialized = mapper.writeValueAsString(options); - GcsOptions deserialized = - mapper.readValue(serialized, PipelineOptions.class).as(GcsOptions.class); - - GoogleCloudStorageReadOptions deserializedReadOptions = - deserialized.getGoogleCloudStorageReadOptions(); - - assertNotNull(deserializedReadOptions); - assertEquals( - GoogleCloudStorageReadOptions.Fadvise.SEQUENTIAL, deserializedReadOptions.getFadvise()); - } }