From 03f56f60b6ec71c64eba2a3bfe8d5375807533b1 Mon Sep 17 00:00:00 2001 From: Radek Stankiewicz Date: Wed, 22 Apr 2026 14:11:17 +0000 Subject: [PATCH 1/2] JsonPayloadSerializerProvider --- .../JsonPayloadSerializerProvider.java | 46 ++++++++++++++++--- 1 file changed, 39 insertions(+), 7 deletions(-) diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/schemas/io/payloads/JsonPayloadSerializerProvider.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/schemas/io/payloads/JsonPayloadSerializerProvider.java index 750cc3418820..c1e253502a79 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/schemas/io/payloads/JsonPayloadSerializerProvider.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/schemas/io/payloads/JsonPayloadSerializerProvider.java @@ -27,6 +27,8 @@ import org.apache.beam.sdk.util.RowJson.RowJsonDeserializer; import org.apache.beam.sdk.util.RowJson.RowJsonSerializer; import org.apache.beam.sdk.util.RowJsonUtils; +import org.apache.beam.sdk.values.Row; +import org.checkerframework.checker.nullness.qual.Nullable; @Internal @AutoService(PayloadSerializerProvider.class) @@ -38,12 +40,42 @@ public String identifier() { @Override public PayloadSerializer getSerializer(Schema schema, Map tableParams) { - ObjectMapper deserializeMapper = - RowJsonUtils.newObjectMapperWith(RowJsonDeserializer.forSchema(schema)); - ObjectMapper serializeMapper = - RowJsonUtils.newObjectMapperWith(RowJsonSerializer.forSchema(schema)); - return PayloadSerializer.of( - row -> RowJsonUtils.rowToJson(serializeMapper, row).getBytes(UTF_8), - bytes -> RowJsonUtils.jsonToRow(deserializeMapper, new String(bytes, UTF_8))); + return new JsonPayloadSerializer(schema); + } + + private static class JsonPayloadSerializer implements PayloadSerializer { + private final Schema schema; + private transient @Nullable ObjectMapper deserializeMapper; + private transient @Nullable ObjectMapper serializeMapper; + + public JsonPayloadSerializer(Schema schema) { + this.schema = schema; + this.deserializeMapper = null; + this.serializeMapper = null; + } + + private synchronized ObjectMapper getDeserializeMapper() { + if (deserializeMapper == null) { + deserializeMapper = RowJsonUtils.newObjectMapperWith(RowJsonDeserializer.forSchema(schema)); + } + return deserializeMapper; + } + + private synchronized ObjectMapper getSerializeMapper() { + if (serializeMapper == null) { + serializeMapper = RowJsonUtils.newObjectMapperWith(RowJsonSerializer.forSchema(schema)); + } + return serializeMapper; + } + + @Override + public byte[] serialize(Row row) { + return RowJsonUtils.rowToJson(getSerializeMapper(), row).getBytes(UTF_8); + } + + @Override + public Row deserialize(byte[] bytes) { + return RowJsonUtils.jsonToRow(getDeserializeMapper(), new String(bytes, UTF_8)); + } } } From ba2ce605ddf355a2d239da11a78feb01bccd2247 Mon Sep 17 00:00:00 2001 From: Radek Stankiewicz Date: Fri, 24 Apr 2026 16:25:29 +0000 Subject: [PATCH 2/2] Apply recommended changes to JsonPayloadSerializer: mark fields volatile and implement double-checked locking --- .../JsonPayloadSerializerProvider.java | 23 ++++++++++++------- 1 file changed, 15 insertions(+), 8 deletions(-) diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/schemas/io/payloads/JsonPayloadSerializerProvider.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/schemas/io/payloads/JsonPayloadSerializerProvider.java index c1e253502a79..d3054e3689a5 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/schemas/io/payloads/JsonPayloadSerializerProvider.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/schemas/io/payloads/JsonPayloadSerializerProvider.java @@ -45,25 +45,32 @@ public PayloadSerializer getSerializer(Schema schema, Map tableP private static class JsonPayloadSerializer implements PayloadSerializer { private final Schema schema; - private transient @Nullable ObjectMapper deserializeMapper; - private transient @Nullable ObjectMapper serializeMapper; + private transient volatile @Nullable ObjectMapper deserializeMapper; + private transient volatile @Nullable ObjectMapper serializeMapper; public JsonPayloadSerializer(Schema schema) { this.schema = schema; - this.deserializeMapper = null; - this.serializeMapper = null; } - private synchronized ObjectMapper getDeserializeMapper() { + private ObjectMapper getDeserializeMapper() { if (deserializeMapper == null) { - deserializeMapper = RowJsonUtils.newObjectMapperWith(RowJsonDeserializer.forSchema(schema)); + synchronized (this) { + if (deserializeMapper == null) { + deserializeMapper = + RowJsonUtils.newObjectMapperWith(RowJsonDeserializer.forSchema(schema)); + } + } } return deserializeMapper; } - private synchronized ObjectMapper getSerializeMapper() { + private ObjectMapper getSerializeMapper() { if (serializeMapper == null) { - serializeMapper = RowJsonUtils.newObjectMapperWith(RowJsonSerializer.forSchema(schema)); + synchronized (this) { + if (serializeMapper == null) { + serializeMapper = RowJsonUtils.newObjectMapperWith(RowJsonSerializer.forSchema(schema)); + } + } } return serializeMapper; }