Skip to content

Commit b03dd52

Browse files
committed
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.
1 parent 83fcb2f commit b03dd52

2 files changed

Lines changed: 54 additions & 1 deletion

File tree

sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/options/GcsOptions.java

Lines changed: 52 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -18,8 +18,19 @@
1818
package org.apache.beam.sdk.extensions.gcp.options;
1919

2020
import com.fasterxml.jackson.annotation.JsonIgnore;
21+
import com.fasterxml.jackson.core.JsonGenerator;
22+
import com.fasterxml.jackson.core.JsonParser;
23+
import com.fasterxml.jackson.databind.DeserializationContext;
24+
import com.fasterxml.jackson.databind.JsonDeserializer;
25+
import com.fasterxml.jackson.databind.JsonNode;
26+
import com.fasterxml.jackson.databind.JsonSerializer;
27+
import com.fasterxml.jackson.databind.ObjectMapper;
28+
import com.fasterxml.jackson.databind.SerializerProvider;
29+
import com.fasterxml.jackson.databind.annotation.JsonDeserialize;
30+
import com.fasterxml.jackson.databind.annotation.JsonSerialize;
2131
import com.google.cloud.hadoop.gcsio.GoogleCloudStorageReadOptions;
2232
import com.google.cloud.hadoop.util.AsyncWriteChannelOptions;
33+
import java.lang.reflect.Method;
2334
import java.util.HashMap;
2435
import java.util.concurrent.ExecutorService;
2536
import org.apache.beam.sdk.extensions.gcp.storage.GcsPathValidator;
@@ -53,8 +64,48 @@ public GoogleCloudStorageReadOptions create(PipelineOptions options) {
5364
}
5465
}
5566

67+
class GcsReadOptionsSerializer extends JsonSerializer<GoogleCloudStorageReadOptions> {
68+
@Override
69+
public void serialize(
70+
GoogleCloudStorageReadOptions value, JsonGenerator gen, SerializerProvider serializers)
71+
throws java.io.IOException {
72+
serializers.defaultSerializeValue(value, gen);
73+
}
74+
}
75+
76+
class GcsReadOptionsDeserializer extends JsonDeserializer<GoogleCloudStorageReadOptions> {
77+
@Override
78+
public GoogleCloudStorageReadOptions deserialize(JsonParser p, DeserializationContext ctxt)
79+
throws java.io.IOException {
80+
ObjectMapper mapper = (ObjectMapper) p.getCodec();
81+
JsonNode root = mapper.readTree(p);
82+
GoogleCloudStorageReadOptions.Builder builder = GoogleCloudStorageReadOptions.builder();
83+
84+
for (Method method : GoogleCloudStorageReadOptions.Builder.class.getMethods()) {
85+
if (method.getName().startsWith("set") && method.getParameterCount() == 1) {
86+
String propName =
87+
Character.toLowerCase(method.getName().charAt(3)) + method.getName().substring(4);
88+
JsonNode node = root.get(propName);
89+
if (node != null && !node.isNull()) {
90+
try {
91+
Class<?> paramType = method.getParameterTypes()[0];
92+
Object val = mapper.treeToValue(node, paramType);
93+
if (val != null) {
94+
method.invoke(builder, java.util.Objects.requireNonNull(val));
95+
}
96+
} catch (Exception ignored) {
97+
// Ignore incompatible or unmappable properties
98+
}
99+
}
100+
}
101+
}
102+
return builder.build();
103+
}
104+
}
105+
56106
/** @deprecated This option will be removed in a future release. */
57-
@JsonIgnore
107+
@JsonSerialize(using = GcsReadOptionsSerializer.class)
108+
@JsonDeserialize(using = GcsReadOptionsDeserializer.class)
58109
@Description(
59110
"The GoogleCloudStorageReadOptions instance that should be used to read from Google Cloud Storage.")
60111
@Default.InstanceFactory(GcsReadOptionsFactory.class)

sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/GcpCoreApiSurfaceTest.java

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -51,6 +51,8 @@ public void testGcpCoreApiSurface() throws Exception {
5151
final Set<Matcher<Class<?>>> allowedClasses =
5252
ImmutableSet.of(
5353
classesInPackage("com.fasterxml.jackson.annotation"),
54+
classesInPackage("com.fasterxml.jackson.core"),
55+
classesInPackage("com.fasterxml.jackson.databind"),
5456
classesInPackage("com.google.api.client.googleapis"),
5557
classesInPackage("com.google.api.client.http"),
5658
classesInPackage("com.google.api.client.json"),

0 commit comments

Comments
 (0)