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 eb9e7a3dacd Override default fadvise to fix regression from
gcs-connector v3 upgrade (#39445)
eb9e7a3dacd is described below
commit eb9e7a3dacdeea074b22d833cb1fda789d09d6f4
Author: Shunping Huang <[email protected]>
AuthorDate: Thu Jul 23 14:34:15 2026 -0400
Override default fadvise to fix regression from gcs-connector v3 upgrade
(#39445)
* Support serialization and deserialization for
GoogleCloudStorageReadOptions
Fixes an issue where GoogleCloudStorageReadOptions set on GcsOptions
was annotated with @JsonIgnore, causing it to be omitted when serializing
PipelineOptions for remote workers (e.g., Dataflow) and falling back to
defaults.
* Add a unit test
* Override default fadvice to SEQUENTIAL
This fixes the regression introduced by gcs-connector v3.
* Spotless
* Only support a subset of fields in GoogleCloudStorageReadOptions.
* Refactor
---
.../google-cloud-platform-core/build.gradle | 1 +
.../sdk/extensions/gcp/options/GcsOptions.java | 85 +++++++++++++++++++++-
.../sdk/extensions/gcp/GcpCoreApiSurfaceTest.java | 2 +
.../sdk/extensions/gcp/options/GcsOptionsTest.java | 47 ++++++++++++
4 files changed, 133 insertions(+), 2 deletions(-)
diff --git a/sdks/java/extensions/google-cloud-platform-core/build.gradle
b/sdks/java/extensions/google-cloud-platform-core/build.gradle
index f1bfb63c7a3..78cfe4739ec 100644
--- a/sdks/java/extensions/google-cloud-platform-core/build.gradle
+++ b/sdks/java/extensions/google-cloud-platform-core/build.gradle
@@ -56,6 +56,7 @@ 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 2da382a5b67..134c4cb3f28 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,6 +18,15 @@
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;
@@ -49,12 +58,16 @@ public interface GcsOptions extends ApplicationNameOptions,
GcpOptions, Pipeline
class GcsReadOptionsFactory implements
DefaultValueFactory<GoogleCloudStorageReadOptions> {
@Override
public GoogleCloudStorageReadOptions create(PipelineOptions options) {
- return GoogleCloudStorageReadOptions.DEFAULT;
+ // 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;
}
}
/** @deprecated This option will be removed in a future release. */
- @JsonIgnore
+ @JsonSerialize(using = GcsReadOptionsSerializer.class)
+ @JsonDeserialize(using = GcsReadOptionsDeserializer.class)
@Description(
"The GoogleCloudStorageReadOptions instance that should be used to read
from Google Cloud Storage.")
@Default.InstanceFactory(GcsReadOptionsFactory.class)
@@ -286,3 +299,71 @@ public interface GcsOptions extends
ApplicationNameOptions, GcpOptions, Pipeline
}
}
}
+
+class GcsReadOptionsSerializer extends
JsonSerializer<GoogleCloudStorageReadOptions> {
+ 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<GoogleCloudStorageReadOptions> {
+ @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 8af5e2260fc..cd51d4fd9d2 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,6 +51,8 @@ public class GcpCoreApiSurfaceTest {
final Set<Matcher<Class<?>>> 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 c499290b851..912f0110e9c 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,11 +18,16 @@
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;
@@ -75,4 +80,46 @@ public class GcsOptionsTest {
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());
+ }
}