This is an automated email from the ASF dual-hosted git repository. shunping pushed a commit to branch fadvice in repository https://gitbox.apache.org/repos/asf/beam.git
commit 064aa6ca086a1d877c72d91db183901af059256d Author: Shunping Huang <[email protected]> AuthorDate: Wed Jul 22 16:12:38 2026 -0400 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. --- .../sdk/extensions/gcp/options/GcsOptions.java | 53 +++++++++++++++++++++- .../sdk/extensions/gcp/GcpCoreApiSurfaceTest.java | 2 + 2 files changed, 54 insertions(+), 1 deletion(-) 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..35cba02f7c1 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,8 +18,19 @@ 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.ObjectMapper; +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.lang.reflect.Method; import java.util.HashMap; import java.util.concurrent.ExecutorService; import org.apache.beam.sdk.extensions.gcp.storage.GcsPathValidator; @@ -53,8 +64,48 @@ public interface GcsOptions extends ApplicationNameOptions, GcpOptions, Pipeline } } + class GcsReadOptionsSerializer extends JsonSerializer<GoogleCloudStorageReadOptions> { + @Override + public void serialize( + GoogleCloudStorageReadOptions value, JsonGenerator gen, SerializerProvider serializers) + throws java.io.IOException { + serializers.defaultSerializeValue(value, gen); + } + } + + class GcsReadOptionsDeserializer extends JsonDeserializer<GoogleCloudStorageReadOptions> { + @Override + public GoogleCloudStorageReadOptions deserialize(JsonParser p, DeserializationContext ctxt) + throws java.io.IOException { + ObjectMapper mapper = (ObjectMapper) p.getCodec(); + JsonNode root = mapper.readTree(p); + GoogleCloudStorageReadOptions.Builder builder = GoogleCloudStorageReadOptions.builder(); + + for (Method method : GoogleCloudStorageReadOptions.Builder.class.getMethods()) { + if (method.getName().startsWith("set") && method.getParameterCount() == 1) { + String propName = + Character.toLowerCase(method.getName().charAt(3)) + method.getName().substring(4); + JsonNode node = root.get(propName); + if (node != null && !node.isNull()) { + try { + Class<?> paramType = method.getParameterTypes()[0]; + Object val = mapper.treeToValue(node, paramType); + if (val != null) { + method.invoke(builder, java.util.Objects.requireNonNull(val)); + } + } catch (Exception ignored) { + // Ignore incompatible or unmappable properties + } + } + } + } + return builder.build(); + } + } + /** @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) 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"),
