This is an automated email from the ASF dual-hosted git repository.
stankiewicz 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 7ada585d853 [Java] Preserve CoderTranslatorRegistrar binary
compatibility (#39919)
7ada585d853 is described below
commit 7ada585d853c419efcaf4f9054eeeb0269638f03
Author: Bruno Volpato <[email protected]>
AuthorDate: Mon Aug 31 12:32:03 2026 -0400
[Java] Preserve CoderTranslatorRegistrar binary compatibility (#39919)
---
CHANGES.md | 1 +
.../construction/CoderTranslatorRegistrar.java | 20 ++++++++---
.../util/construction/CoderTranslationTest.java | 40 ++++++++++++++++++++++
3 files changed, 56 insertions(+), 5 deletions(-)
diff --git a/CHANGES.md b/CHANGES.md
index 045b95c4e76..0c519076f0d 100644
--- a/CHANGES.md
+++ b/CHANGES.md
@@ -88,6 +88,7 @@
## Bugfixes
+* (Java) Restored binary compatibility for `CoderTranslatorRegistrar`
implementations compiled against Beam 2.76 and earlier
([#38714](https://github.com/apache/beam/issues/38714)).
* (Java) Fixed the Spark runner firing processing-time timers in reverse
timestamp order ([#39824](https://github.com/apache/beam/issues/39824)).
* (Python) Fixed incorrect profiler options handling on portable runners
([#39613](https://github.com/apache/beam/issues/39613)).
* (Java) KafkaIO dynamic reads no longer require the obsolete `beam_fn_api`
experiment ([#29998](https://github.com/apache/beam/issues/29998)).
diff --git
a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/CoderTranslatorRegistrar.java
b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/CoderTranslatorRegistrar.java
index 44e8c2956ae..b781db08409 100644
---
a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/CoderTranslatorRegistrar.java
+++
b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/CoderTranslatorRegistrar.java
@@ -41,13 +41,23 @@ public interface CoderTranslatorRegistrar {
* Returns whether the given Coder is known to this
CoderTranslatorRegistrar. If the Coder is
* known, then getCoderTranslator() will return a non-null CoderTranslator.
*/
- boolean isKnownCoder(Coder<?> coder, PipelineOptions options);
+ default boolean isKnownCoder(Coder<?> coder, PipelineOptions options) {
+ return getCoderURNs().containsKey(coder.getClass());
+ }
/** Returns the CoderTranslator to use for this Coder, or null if the Coder
is not known. */
- @Nullable
- CoderTranslator<? extends Coder> getCoderTranslator(Class<? extends Coder>
coderClass);
+ default @Nullable CoderTranslator<? extends Coder> getCoderTranslator(
+ Class<? extends Coder> coderClass) {
+ return getCoderTranslators().get(coderClass);
+ }
/** Returns the Coder to use for the given Urn, or null if the Urn is for an
unknown Coder. */
- @Nullable
- Class<? extends Coder> getCoderForUrn(String coderUrn);
+ default @Nullable Class<? extends Coder> getCoderForUrn(String coderUrn) {
+ for (Map.Entry<Class<? extends Coder>, String> coderUrnEntry :
getCoderURNs().entrySet()) {
+ if (coderUrnEntry.getValue().equals(coderUrn)) {
+ return coderUrnEntry.getKey();
+ }
+ }
+ return null;
+ }
}
diff --git
a/sdks/java/core/src/test/java/org/apache/beam/sdk/util/construction/CoderTranslationTest.java
b/sdks/java/core/src/test/java/org/apache/beam/sdk/util/construction/CoderTranslationTest.java
index 1ec0a74f5be..f743b9a0b60 100644
---
a/sdks/java/core/src/test/java/org/apache/beam/sdk/util/construction/CoderTranslationTest.java
+++
b/sdks/java/core/src/test/java/org/apache/beam/sdk/util/construction/CoderTranslationTest.java
@@ -28,6 +28,7 @@ import java.io.InputStream;
import java.io.OutputStream;
import java.io.Serializable;
import java.util.HashSet;
+import java.util.Map;
import java.util.Set;
import org.apache.beam.model.pipeline.v1.RunnerApi;
import org.apache.beam.model.pipeline.v1.RunnerApi.Components;
@@ -46,6 +47,7 @@ import org.apache.beam.sdk.coders.SerializableCoder;
import org.apache.beam.sdk.coders.StringUtf8Coder;
import org.apache.beam.sdk.coders.TimestampPrefixingWindowCoder;
import org.apache.beam.sdk.coders.VarLongCoder;
+import org.apache.beam.sdk.options.PipelineOptionsFactory;
import org.apache.beam.sdk.schemas.AutoValueSchema;
import org.apache.beam.sdk.schemas.NoSuchSchemaException;
import org.apache.beam.sdk.schemas.Schema;
@@ -63,6 +65,7 @@ import org.apache.beam.sdk.values.TypeDescriptor;
import org.apache.beam.sdk.values.WindowedValues;
import org.apache.beam.sdk.values.WindowedValues.FullWindowedValueCoder;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList;
+import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableSet;
import org.hamcrest.Matchers;
import org.junit.Test;
@@ -173,6 +176,43 @@ public class CoderTranslationTest {
CoderTranslation.getKnownTranslators().keySet(),
hasItems(new
ModelCoderRegistrar().getCoderTranslators().keySet().toArray(new Class[0])));
}
+
+ @Test
+ public void legacyRegistrarUsesDefaultLookupMethods() {
+ CoderTranslatorRegistrar registrar = new
LegacyCoderTranslatorRegistrar();
+
+ assertThat(
+ registrar.isKnownCoder(StringUtf8Coder.of(),
PipelineOptionsFactory.create()),
+ equalTo(true));
+ assertThat(
+ registrar.getCoderTranslator(StringUtf8Coder.class),
+ equalTo(LegacyCoderTranslatorRegistrar.TRANSLATOR));
+ assertThat(
+ registrar.getCoderForUrn(LegacyCoderTranslatorRegistrar.URN),
+ equalTo(StringUtf8Coder.class));
+ assertThat(
+ registrar.isKnownCoder(VarLongCoder.of(),
PipelineOptionsFactory.create()),
+ equalTo(false));
+ assertThat(registrar.getCoderTranslator(VarLongCoder.class),
Matchers.nullValue());
+ assertThat(registrar.getCoderForUrn("unknown"), Matchers.nullValue());
+ }
+
+ /** Implements only the registrar methods available before Beam 2.77.0. */
+ private static class LegacyCoderTranslatorRegistrar implements
CoderTranslatorRegistrar {
+ private static final String URN = "beam:coder:legacy_test:v1";
+ private static final CoderTranslator<StringUtf8Coder> TRANSLATOR =
+ CoderTranslators.atomic(StringUtf8Coder.class);
+
+ @Override
+ public Map<Class<? extends Coder>, String> getCoderURNs() {
+ return ImmutableMap.of(StringUtf8Coder.class, URN);
+ }
+
+ @Override
+ public Map<Class<? extends Coder>, CoderTranslator<? extends Coder>>
getCoderTranslators() {
+ return ImmutableMap.of(StringUtf8Coder.class, TRANSLATOR);
+ }
+ }
}
/** Tests round-trip coder encodings for both known and unknown {@link Coder
coders}. */