From 194ee1ea04144667a5ea697bf2e6e64ffc73f269 Mon Sep 17 00:00:00 2001 From: Andrew Crites Date: Tue, 28 Jul 2026 15:04:47 -0700 Subject: [PATCH 1/5] Changes SplittableDoFn to call TruncateRestriction first when getting a timer caused by drain. We then pass the residual restriction (if present) to ProcessElement. --- .../SplittableParDoViaKeyedWorkItems.java | 70 ++++- .../core/SplittableParDoProcessFnTest.java | 282 ++++++++++++++++-- 2 files changed, 332 insertions(+), 20 deletions(-) diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/SplittableParDoViaKeyedWorkItems.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/SplittableParDoViaKeyedWorkItems.java index a750b01963f6..4eea2c0509dd 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/SplittableParDoViaKeyedWorkItems.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/SplittableParDoViaKeyedWorkItems.java @@ -469,8 +469,9 @@ public String getErrorContext() { restrictionState.readLater(); watermarkEstimatorState.readLater(); WindowedValue read = elementState.read(); + RestrictionT restriction = restrictionState.read(); if (timer.causedByDrain() == CausedByDrain.CAUSED_BY_DRAIN) { - read = + WindowedValue drainRead = WindowedValues.of( read.getValue(), read.getTimestamp(), @@ -481,8 +482,73 @@ public String getErrorContext() { CausedByDrain.CAUSED_BY_DRAIN, read.getOpenTelemetryContext(), read.getValueKind()); + RestrictionTracker.TruncateResult truncateResult = + invoker.invokeTruncateRestriction( + new BaseArgumentProvider() { + @Override + public InputT element(DoFn doFn) { + return drainRead.getValue(); + } + + @Override + public Object restriction() { + return restriction; + } + + @Override + public RestrictionTracker restrictionTracker() { + return invoker.invokeNewTracker(this); + } + + @Override + public Instant timestamp(DoFn doFn) { + return drainRead.getTimestamp(); + } + + @Override + public PipelineOptions pipelineOptions() { + return c.getPipelineOptions(); + } + + @Override + public PaneInfo paneInfo(DoFn doFn) { + return drainRead.getPaneInfo(); + } + + @Override + public BoundedWindow window() { + return Iterables.getOnlyElement(drainRead.getWindows()); + } + + @Override + public Object sideInput(String tagId) { + PCollectionView view = sideInputMapping.get(tagId); + if (view == null) { + throw new IllegalArgumentException( + "calling getSideInput() with unknown view"); + } + return sideInputReader.get( + view, view.getWindowMappingFn().getSideInputWindow(window())); + } + + @Override + public String getErrorContext() { + return ProcessFn.class.getSimpleName() + ".invokeTruncateRestriction"; + } + }); + if (truncateResult == null) { + elementState.clear(); + restrictionState.clear(); + watermarkEstimatorState.clear(); + holdState.clear(); + return; + } + RestrictionT truncatedRestriction = truncateResult.getTruncatedRestriction(); + elementAndRestriction = KV.of(drainRead, truncatedRestriction); + restrictionState.write(truncatedRestriction); + } else { + elementAndRestriction = KV.of(read, restriction); } - elementAndRestriction = KV.of(read, restrictionState.read()); watermarkEstimatorStateT = watermarkEstimatorState.read(); } diff --git a/runners/core-java/src/test/java/org/apache/beam/runners/core/SplittableParDoProcessFnTest.java b/runners/core-java/src/test/java/org/apache/beam/runners/core/SplittableParDoProcessFnTest.java index 381e41c98705..83c454492613 100644 --- a/runners/core-java/src/test/java/org/apache/beam/runners/core/SplittableParDoProcessFnTest.java +++ b/runners/core-java/src/test/java/org/apache/beam/runners/core/SplittableParDoProcessFnTest.java @@ -36,6 +36,7 @@ import java.util.Arrays; import java.util.Collections; import java.util.List; +import java.util.Map; import java.util.NoSuchElementException; import java.util.concurrent.Executors; import org.apache.beam.runners.core.SplittableParDoViaKeyedWorkItems.ProcessFn; @@ -48,12 +49,15 @@ import org.apache.beam.sdk.state.TimeDomain; import org.apache.beam.sdk.testing.ResetDateTimeProvider; import org.apache.beam.sdk.testing.TestPipeline; +import org.apache.beam.sdk.transforms.Create; import org.apache.beam.sdk.transforms.DoFn; import org.apache.beam.sdk.transforms.DoFnTester; +import org.apache.beam.sdk.transforms.View; import org.apache.beam.sdk.transforms.splittabledofn.HasDefaultTracker; import org.apache.beam.sdk.transforms.splittabledofn.ManualWatermarkEstimator; import org.apache.beam.sdk.transforms.splittabledofn.OffsetRangeTracker; import org.apache.beam.sdk.transforms.splittabledofn.RestrictionTracker; +import org.apache.beam.sdk.transforms.splittabledofn.RestrictionTracker.IsBounded; import org.apache.beam.sdk.transforms.splittabledofn.SplitResult; import org.apache.beam.sdk.transforms.splittabledofn.WatermarkEstimators; import org.apache.beam.sdk.transforms.windowing.BoundedWindow; @@ -152,6 +156,29 @@ private static class ProcessFnTester< int maxOutputsPerBundle, Duration maxBundleDuration) throws Exception { + this( + currentProcessingTime, + fn, + inputCoder, + restrictionCoder, + watermarkEstimatorStateCoder, + maxOutputsPerBundle, + maxBundleDuration, + Collections.emptyMap(), + NullSideInputReader.empty()); + } + + ProcessFnTester( + Instant currentProcessingTime, + final DoFn fn, + Coder inputCoder, + Coder restrictionCoder, + Coder watermarkEstimatorStateCoder, + int maxOutputsPerBundle, + Duration maxBundleDuration, + Map> sideInputMapping, + SideInputReader sideInputReader) + throws Exception { // The exact windowing strategy doesn't matter in this test, but it should be able to // encode IntervalWindow's because that's what all tests here use. WindowingStrategy windowingStrategy = @@ -163,35 +190,20 @@ private static class ProcessFnTester< restrictionCoder, watermarkEstimatorStateCoder, windowingStrategy, - Collections.emptyMap()); + sideInputMapping); this.tester = DoFnTester.of(processFn); this.timerInternals = new InMemoryTimerInternals(); this.stateInternals = new TestInMemoryStateInternals<>("dummy"); processFn.setStateInternalsFactory(key -> stateInternals); processFn.setTimerInternalsFactory(key -> timerInternals); - processFn.setSideInputReader(NullSideInputReader.empty()); + processFn.setSideInputReader(sideInputReader); processFn.setProcessElementInvoker( new OutputAndTimeBoundedSplittableProcessElementInvoker<>( fn, tester.getPipelineOptions(), new DoFnTesterWindowedValueReceiver(tester), tester.getMainOutputTag(), - new SideInputReader() { - @Override - public T get(PCollectionView view, BoundedWindow window) { - throw new NoSuchElementException(); - } - - @Override - public boolean contains(PCollectionView view) { - return false; - } - - @Override - public boolean isEmpty() { - return true; - } - }, + sideInputReader, Executors.newSingleThreadScheduledExecutor(Executors.defaultThreadFactory()), maxOutputsPerBundle, maxBundleDuration, @@ -790,4 +802,238 @@ public void testReportsBacklogWithoutGetSize() throws Exception { assertEquals(7.0, backlogs.get(0), 0.001); } } + + private static class TruncateFn extends DoFn { + private final boolean truncateToNull; + private final List calls = new ArrayList<>(); + + public TruncateFn(boolean truncateToNull) { + this.truncateToNull = truncateToNull; + } + + @ProcessElement + public ProcessContinuation process( + ProcessContext c, RestrictionTracker tracker) { + for (long i = tracker.currentRestriction().getFrom(); tracker.tryClaim(i); ++i) { + c.output(c.element() + ":" + i); + if (i == 2) { + return resume(); + } + } + return stop(); + } + + @GetInitialRestriction + public OffsetRange getInitialRestriction() { + return new OffsetRange(0, 10); + } + + @NewTracker + public OffsetRangeTracker newTracker(@Restriction OffsetRange range) { + return new OffsetRangeTracker(range); + } + + @TruncateRestriction + public RestrictionTracker.TruncateResult truncate( + @Restriction OffsetRange restriction, @Element Integer element) { + calls.add("truncate:" + element + ":" + restriction); + if (truncateToNull) { + return null; + } + // Truncate so that we only process one more element. + return RestrictionTracker.TruncateResult.of( + new OffsetRange(restriction.getFrom(), restriction.getFrom() + 1)); + } + } + + @Test + public void testTruncateRestrictionOnDrain() throws Exception { + TruncateFn fn = new TruncateFn(false); + Instant base = Instant.now(); + + try (ProcessFnTester tester = + new ProcessFnTester<>( + base, + fn, + BigEndianIntegerCoder.of(), + SerializableCoder.of(OffsetRange.class), + VoidCoder.of(), + MAX_OUTPUTS_PER_BUNDLE, + MAX_BUNDLE_DURATION)) { + tester.startElement(42, new OffsetRange(0, 10)); + assertThat(tester.takeOutputElements(), contains("42:0", "42:1", "42:2")); + + assertTrue(tester.advanceDrain()); + assertThat(tester.takeOutputElements(), contains("42:3")); + assertEquals(Collections.singletonList("truncate:42:[3, 10)"), fn.calls); + assertEquals(null, tester.getWatermarkHold()); + } + } + + @Test + public void testTruncateRestrictionReturnsNullOnDrain() throws Exception { + TruncateFn fn = new TruncateFn(true); + Instant base = Instant.now(); + + try (ProcessFnTester tester = + new ProcessFnTester<>( + base, + fn, + BigEndianIntegerCoder.of(), + SerializableCoder.of(OffsetRange.class), + VoidCoder.of(), + MAX_OUTPUTS_PER_BUNDLE, + MAX_BUNDLE_DURATION)) { + tester.startElement(42, new OffsetRange(0, 10)); + assertThat(tester.takeOutputElements(), contains("42:0", "42:1", "42:2")); + + assertTrue(tester.advanceDrain()); + assertTrue(tester.takeOutputElements().isEmpty()); + assertEquals(Collections.singletonList("truncate:42:[3, 10)"), fn.calls); + assertEquals(null, tester.getWatermarkHold()); + } + } + + private static class TruncateWithSideInputFn extends DoFn { + private final List calls = new ArrayList<>(); + + @ProcessElement + public ProcessContinuation process( + ProcessContext c, RestrictionTracker tracker) { + for (long i = tracker.currentRestriction().getFrom(); tracker.tryClaim(i); ++i) { + c.output(c.element() + ":" + i); + if (i == 2) { + return resume(); + } + } + return stop(); + } + + @GetInitialRestriction + public OffsetRange getInitialRestriction() { + return new OffsetRange(0, 10); + } + + @NewTracker + public OffsetRangeTracker newTracker(@Restriction OffsetRange range) { + return new OffsetRangeTracker(range); + } + + @TruncateRestriction + public RestrictionTracker.TruncateResult truncate( + @Restriction OffsetRange restriction, + @Element Integer element, + @SideInput("sideInput") String sideInput) { + calls.add("truncate:" + element + ":" + sideInput + ":" + restriction); + return RestrictionTracker.TruncateResult.of( + new OffsetRange(restriction.getFrom(), restriction.getFrom() + 1)); + } + } + + @Test + public void testTruncateRestrictionWithSideInputOnDrain() throws Exception { + TruncateWithSideInputFn fn = new TruncateWithSideInputFn(); + Instant base = Instant.now(); + PCollectionView view = + TestPipeline.create().apply(Create.of("sideValue")).apply(View.asSingleton()); + Map> sideInputMapping = Collections.singletonMap("sideInput", view); + SideInputReader sideInputReader = + new SideInputReader() { + @Override + public T get(PCollectionView v, BoundedWindow window) { + if (v.equals(view)) { + return (T) "sideValue"; + } + throw new NoSuchElementException(); + } + + @Override + public boolean contains(PCollectionView v) { + return v.equals(view); + } + + @Override + public boolean isEmpty() { + return false; + } + }; + + try (ProcessFnTester tester = + new ProcessFnTester<>( + base, + fn, + BigEndianIntegerCoder.of(), + SerializableCoder.of(OffsetRange.class), + VoidCoder.of(), + MAX_OUTPUTS_PER_BUNDLE, + MAX_BUNDLE_DURATION, + sideInputMapping, + sideInputReader)) { + tester.startElement(42, new OffsetRange(0, 10)); + assertThat(tester.takeOutputElements(), contains("42:0", "42:1", "42:2")); + + assertTrue(tester.advanceDrain()); + assertThat(tester.takeOutputElements(), contains("42:3")); + assertEquals(Collections.singletonList("truncate:42:sideValue:[3, 10)"), fn.calls); + } + } + + private static class UnboundedOffsetRangeTracker extends OffsetRangeTracker { + public UnboundedOffsetRangeTracker(OffsetRange range) { + super(range); + } + + @Override + public IsBounded isBounded() { + return IsBounded.UNBOUNDED; + } + } + + // Tests that if we don't override T + private static class DefaultTruncateUnboundedFn extends DoFn { + @ProcessElement + public ProcessContinuation process( + ProcessContext c, RestrictionTracker tracker) { + for (long i = tracker.currentRestriction().getFrom(); tracker.tryClaim(i); ++i) { + c.output(c.element() + ":" + i); + if (i == 2) { + return resume(); + } + } + return stop(); + } + + @GetInitialRestriction + public OffsetRange getInitialRestriction() { + return new OffsetRange(0, 10); + } + + @NewTracker + public RestrictionTracker newTracker(@Restriction OffsetRange range) { + return new UnboundedOffsetRangeTracker(range); + } + } + + @Test + public void testDefaultTruncateRestrictionUnboundedStopsOnDrain() throws Exception { + DefaultTruncateUnboundedFn fn = new DefaultTruncateUnboundedFn(); + Instant base = Instant.now(); + + try (ProcessFnTester tester = + new ProcessFnTester<>( + base, + fn, + BigEndianIntegerCoder.of(), + SerializableCoder.of(OffsetRange.class), + VoidCoder.of(), + MAX_OUTPUTS_PER_BUNDLE, + MAX_BUNDLE_DURATION)) { + tester.startElement(42, new OffsetRange(0, 10)); + assertThat(tester.takeOutputElements(), contains("42:0", "42:1", "42:2")); + + assertTrue(tester.advanceDrain()); + assertTrue(tester.takeOutputElements().isEmpty()); + assertEquals(null, tester.getWatermarkHold()); + } + } } From d4bcf9bc8dac2b5ca1f362a64c7587ab804294d5 Mon Sep 17 00:00:00 2001 From: Andrew Crites Date: Tue, 28 Jul 2026 15:12:57 -0700 Subject: [PATCH 2/5] Adds the rest of a comment string. --- .../apache/beam/runners/core/SplittableParDoProcessFnTest.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/runners/core-java/src/test/java/org/apache/beam/runners/core/SplittableParDoProcessFnTest.java b/runners/core-java/src/test/java/org/apache/beam/runners/core/SplittableParDoProcessFnTest.java index 83c454492613..247412a85dc4 100644 --- a/runners/core-java/src/test/java/org/apache/beam/runners/core/SplittableParDoProcessFnTest.java +++ b/runners/core-java/src/test/java/org/apache/beam/runners/core/SplittableParDoProcessFnTest.java @@ -989,7 +989,8 @@ public IsBounded isBounded() { } } - // Tests that if we don't override T + // Tests that if we don't override TruncateRestriction, the default TruncateRestriction + // implementation is used (which for unbounded restrictions stops processing immediately). private static class DefaultTruncateUnboundedFn extends DoFn { @ProcessElement public ProcessContinuation process( From 166d6a85ccb9ca079c0fe228a1c35dce922eb83c Mon Sep 17 00:00:00 2001 From: Andrew Crites Date: Thu, 30 Jul 2026 14:37:49 -0700 Subject: [PATCH 3/5] Revert "Revert #37631 and #38497 on HEAD (#38516)" This reverts commit 357fd2622114ccce99726baa5e9735895e5392ed. --- ...m_PostCommit_Java_ValidatesRunner_ULR.json | 3 +- CHANGES.md | 3 + .../dataflow/DataflowPipelineTranslator.java | 2 +- .../beam/runners/dataflow/DataflowRunner.java | 54 +++++++------ .../DataflowPipelineTranslatorTest.java | 77 ++++++++++++++---- .../control/ProcessBundleDescriptorsTest.java | 6 +- runners/portability/java/build.gradle | 5 ++ .../util/construction/CoderTranslation.java | 62 ++++++++++++++- .../util/construction/CoderTranslator.java | 1 + .../CoderTranslatorRegistrar.java | 16 ++++ .../util/construction/CoderTranslators.java | 79 +++++++++++++++++++ .../construction/ModelCoderRegistrar.java | 28 ++++++- .../sdk/util/construction/ModelCoders.java | 2 + .../construction/RehydratedComponents.java | 3 +- .../sdk/util/construction/SdkComponents.java | 39 +++++---- .../construction/CoderTranslationTest.java | 38 ++++++++- .../expansion/service/ExpansionService.java | 2 +- .../avro/AvroGenericCoderRegistrar.java | 18 +++++ .../fn/harness/state/StateBackedIterable.java | 22 ++++++ 19 files changed, 390 insertions(+), 70 deletions(-) diff --git a/.github/trigger_files/beam_PostCommit_Java_ValidatesRunner_ULR.json b/.github/trigger_files/beam_PostCommit_Java_ValidatesRunner_ULR.json index 6e2f429dd24e..fbd81891f93b 100644 --- a/.github/trigger_files/beam_PostCommit_Java_ValidatesRunner_ULR.json +++ b/.github/trigger_files/beam_PostCommit_Java_ValidatesRunner_ULR.json @@ -2,5 +2,6 @@ "https://github.com/apache/beam/pull/34902": "Introducing OutputBuilder", "comment": "Modify this file in a trivial way to cause this test suite to run", "https://github.com/apache/beam/pull/31156": "noting that PR #31156 should run this test", - "https://github.com/apache/beam/pull/35159": "moving WindowedValue and making an interface" + "https://github.com/apache/beam/pull/35159": "moving WindowedValue and making an interface", + "https://github.com/apache/beam/pull/38497": "sickbay two failed tests" } diff --git a/CHANGES.md b/CHANGES.md index a1fdede94494..e4db58147260 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -173,6 +173,9 @@ ## Breaking Changes +* Portable Java SDK now encodes SchemaCoders in a portable way ([#34672](https://github.com/apache/beam/issues/34672)). + - Original custom Java coder encoding can still be obtained using [StreamingOptions.setUpdateCompatibilityVersion("2.73")](https://github.com/apache/beam/blob/2cf0930e7ae1aa389c26ce6639b584877a3e31d9/sdks/java/core/src/main/java/org/apache/beam/sdk/options/StreamingOptions.java#L47) ([#34672](https://github.com/apache/beam/issues/34672)). + - Fixes ([#36496](https://github.com/apache/beam/issues/36496)), ([#30276](https://github.com/apache/beam/issues/30276)), ([#29245](https://github.com/apache/beam/issues/29245)). * (Python) Made Beartype the default fallback type checking tool. This can be disabled with the `--disable_beartype` pipeline option. ([#38275](https://github.com/apache/beam/issues/38275)) ## Deprecations diff --git a/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowPipelineTranslator.java b/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowPipelineTranslator.java index 4016f31a5475..1609cf6ea238 100644 --- a/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowPipelineTranslator.java +++ b/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowPipelineTranslator.java @@ -221,7 +221,7 @@ public Boolean visit(OrFinallyTrigger trigger) { private static byte[] serializeWindowingStrategy( WindowingStrategy windowingStrategy, PipelineOptions options) { try { - SdkComponents sdkComponents = SdkComponents.create(); + SdkComponents sdkComponents = SdkComponents.create(options); String workerHarnessContainerImageURL = DataflowRunner.getContainerImageForJob(options.as(DataflowPipelineOptions.class)); diff --git a/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowRunner.java b/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowRunner.java index 1b03cc10351e..ee749901ec62 100644 --- a/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowRunner.java +++ b/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowRunner.java @@ -1355,19 +1355,20 @@ public DataflowPipelineJob run(Pipeline pipeline) { // with the SDK harness image (which implements Fn API). // // The same Environment is used in different and contradictory ways, depending on whether - // it is a v1 or v2 job submission. + // it is a portable or non-portable job submission. RunnerApi.Environment defaultEnvironmentForDataflow = Environments.createDockerEnvironment(workerHarnessContainerImageURL); - // The SdkComponents for portable an non-portable job submission must be kept distinct. Both + // The SdkComponents for portable and non-portable job submission must be kept distinct. Both // need the default environment. - SdkComponents portableComponents = SdkComponents.create(); - portableComponents.registerEnvironment( - defaultEnvironmentForDataflow - .toBuilder() - .addAllDependencies(getDefaultArtifacts()) - .addAllCapabilities(Environments.getJavaCapabilities()) - .build()); + SdkComponents portableComponents = + SdkComponents.create( + options, + defaultEnvironmentForDataflow + .toBuilder() + .addAllDependencies(getDefaultArtifacts()) + .addAllCapabilities(Environments.getJavaCapabilities()) + .build()); RunnerApi.Pipeline portablePipelineProto = PipelineTranslation.toProto(pipeline, portableComponents, false); @@ -1400,27 +1401,28 @@ public DataflowPipelineJob run(Pipeline pipeline) { "Skipping Dataflow Streaming Java Runner transform replacements since job will run on Dataflow Portable Runner."); } else { // Now rewrite things to be as needed for Dataflow Streaming Java Runner (mutates the - // pipeline) + // pipeline). // This way the job submitted is valid for Dataflow Streaming Java Runner and Dataflow - // Portable Runner, simultaneously + // Portable Runner, simultaneously. replaceV1Transforms(pipeline); } - // Capture the SdkComponents for look up during step translations - SdkComponents dataflowV1Components = SdkComponents.create(); - dataflowV1Components.registerEnvironment( - defaultEnvironmentForDataflow - .toBuilder() - .addAllDependencies(getDefaultArtifacts()) - .addAllCapabilities(Environments.getJavaCapabilities()) - .build()); + // Capture the SdkComponents for look up during step translations. + SdkComponents dataflowNonPortableComponents = + SdkComponents.create( + options, + defaultEnvironmentForDataflow + .toBuilder() + .addAllDependencies(getDefaultArtifacts()) + .addAllCapabilities(Environments.getJavaCapabilities()) + .build()); // No need to perform transform upgrading for the Dataflow Streaming Java Runner proto. - RunnerApi.Pipeline dataflowV1PipelineProto = - PipelineTranslation.toProto(pipeline, dataflowV1Components, true, false); + RunnerApi.Pipeline dataflowNonPortablePipelineProto = + PipelineTranslation.toProto(pipeline, dataflowNonPortableComponents, true, false); if (LOG.isDebugEnabled()) { LOG.debug( - "Dataflow v1 pipeline proto:\n{}", - TextFormat.printer().printToString(dataflowV1PipelineProto)); + "Dataflow non-portable worker pipeline proto:\n{}", + TextFormat.printer().printToString(dataflowNonPortablePipelineProto)); } // Set a unique client_request_id in the CreateJob request. @@ -1440,7 +1442,11 @@ public DataflowPipelineJob run(Pipeline pipeline) { JobSpecification jobSpecification = translator.translate( - pipeline, dataflowV1PipelineProto, dataflowV1Components, this, packages); + pipeline, + dataflowNonPortablePipelineProto, + dataflowNonPortableComponents, + this, + packages); if (!isNullOrEmpty(dataflowOptions.getDataflowWorkerJar()) && !useUnifiedWorker(options)) { List experiments = diff --git a/runners/google-cloud-dataflow-java/src/test/java/org/apache/beam/runners/dataflow/DataflowPipelineTranslatorTest.java b/runners/google-cloud-dataflow-java/src/test/java/org/apache/beam/runners/dataflow/DataflowPipelineTranslatorTest.java index f8818931a68f..953f7a638ede 100644 --- a/runners/google-cloud-dataflow-java/src/test/java/org/apache/beam/runners/dataflow/DataflowPipelineTranslatorTest.java +++ b/runners/google-cloud-dataflow-java/src/test/java/org/apache/beam/runners/dataflow/DataflowPipelineTranslatorTest.java @@ -47,6 +47,7 @@ import com.google.api.services.dataflow.model.Job; import com.google.api.services.dataflow.model.Step; import com.google.api.services.dataflow.model.WorkerPool; +import com.google.auto.value.AutoValue; import java.io.File; import java.io.IOException; import java.io.Serializable; @@ -92,6 +93,8 @@ import org.apache.beam.sdk.options.PipelineOptionsFactory; import org.apache.beam.sdk.options.StreamingOptions; import org.apache.beam.sdk.options.ValueProvider; +import org.apache.beam.sdk.schemas.AutoValueSchema; +import org.apache.beam.sdk.schemas.annotations.DefaultSchema; import org.apache.beam.sdk.state.StateSpec; import org.apache.beam.sdk.state.StateSpecs; import org.apache.beam.sdk.state.ValueState; @@ -166,15 +169,11 @@ public class DataflowPipelineTranslatorTest implements Serializable { @Rule public transient ExpectedException thrown = ExpectedException.none(); private SdkComponents createSdkComponents(PipelineOptions options) { - SdkComponents sdkComponents = SdkComponents.create(); - String containerImageURL = DataflowRunner.getContainerImageForJob(options.as(DataflowPipelineOptions.class)); RunnerApi.Environment defaultEnvironmentForDataflow = Environments.createDockerEnvironment(containerImageURL); - - sdkComponents.registerEnvironment(defaultEnvironmentForDataflow); - return sdkComponents; + return SdkComponents.create(options, defaultEnvironmentForDataflow); } // A Custom Mockito matcher for an initial Job that checks that all @@ -1294,15 +1293,16 @@ public String apply(byte[] input) { file1.deleteOnExit(); File file2 = File.createTempFile("file2-", ".txt"); file2.deleteOnExit(); - SdkComponents sdkComponents = SdkComponents.create(); - sdkComponents.registerEnvironment( - Environments.createDockerEnvironment(DataflowRunner.getContainerImageForJob(options)) - .toBuilder() - .addAllDependencies( - Environments.getArtifacts( - ImmutableList.of("file1.txt=" + file1, "file2.txt=" + file2))) - .addAllCapabilities(Environments.getJavaCapabilities()) - .build()); + SdkComponents sdkComponents = + SdkComponents.create( + options, + Environments.createDockerEnvironment(DataflowRunner.getContainerImageForJob(options)) + .toBuilder() + .addAllDependencies( + Environments.getArtifacts( + ImmutableList.of("file1.txt=" + file1, "file2.txt=" + file2))) + .addAllCapabilities(Environments.getJavaCapabilities()) + .build()); RunnerApi.Pipeline pipelineProto = PipelineTranslation.toProto(pipeline, sdkComponents, true); @@ -1870,4 +1870,53 @@ public OffsetRange getInitialRange(@SuppressWarnings("unused") @Element String e return null; } } + + @AutoValue + @DefaultSchema(AutoValueSchema.class) + public abstract static class SimpleAutoValue { + public abstract String getString(); + + public abstract int getInt32(); + + public abstract long getInt64(); + + public static DataflowPipelineTranslatorTest.SimpleAutoValue of( + String string, int int32, long int64) { + return new AutoValue_DataflowPipelineTranslatorTest_SimpleAutoValue(string, int32, int64); + } + } + + @Test + public void testSchemaCoderTranslation() throws Exception { + DataflowPipelineOptions options = buildPipelineOptions(); + Pipeline pipeline = Pipeline.create(options); + pipeline + .apply(Impulse.create()) + .apply( + MapElements.via( + new SimpleFunction() { + @Override + public SimpleAutoValue apply(byte[] input) { + return SimpleAutoValue.of("foo", 5, 10L); + } + })) + .apply(Window.into(FixedWindows.of(Duration.standardMinutes(1)))); + { + SdkComponents sdkComponents = createSdkComponents(options); + RunnerApi.Pipeline pipelineProto = PipelineTranslation.toProto(pipeline, sdkComponents, true); + Map coders = pipelineProto.getComponents().getCodersMap(); + assertTrue(coders.containsKey("SchemaCoder")); + assertEquals("beam:coder:schema:v1", coders.get("SchemaCoder").getSpec().getUrn()); + } + + // Prior to version 2.74, SchemaCoders are translated as custom java coders. + { + options.as(StreamingOptions.class).setUpdateCompatibilityVersion("2.73"); + SdkComponents sdkComponents = createSdkComponents(options); + RunnerApi.Pipeline pipelineProto = PipelineTranslation.toProto(pipeline, sdkComponents, true); + Map coders = pipelineProto.getComponents().getCodersMap(); + assertTrue(coders.containsKey("SchemaCoder")); + assertEquals("beam:coders:javasdk:0.1", coders.get("SchemaCoder").getSpec().getUrn()); + } + } } diff --git a/runners/java-fn-execution/src/test/java/org/apache/beam/runners/fnexecution/control/ProcessBundleDescriptorsTest.java b/runners/java-fn-execution/src/test/java/org/apache/beam/runners/fnexecution/control/ProcessBundleDescriptorsTest.java index 21d7550c38b9..9ea7404053d6 100644 --- a/runners/java-fn-execution/src/test/java/org/apache/beam/runners/fnexecution/control/ProcessBundleDescriptorsTest.java +++ b/runners/java-fn-execution/src/test/java/org/apache/beam/runners/fnexecution/control/ProcessBundleDescriptorsTest.java @@ -78,7 +78,8 @@ public void testLengthPrefixingOfKeyCoderInStatefulExecutableStage() throws Exce // Add another stateful stage with a non-standard key coder Pipeline p = Pipeline.create(); Coder keycoder = VoidCoder.of(); - assertThat(ModelCoderRegistrar.isKnownCoder(keycoder), is(false)); + ModelCoderRegistrar coderRegistrar = new ModelCoderRegistrar(); + assertThat(coderRegistrar.isKnownCoder(keycoder, p.getOptions()), is(false)); p.apply("impulse", Impulse.create()) .apply( "create", @@ -165,7 +166,8 @@ public void onTimer() {} public void testLengthPrefixingOfInputCoderExecutableStage() throws Exception { Pipeline p = Pipeline.create(); Coder voidCoder = VoidCoder.of(); - assertThat(ModelCoderRegistrar.isKnownCoder(voidCoder), is(false)); + ModelCoderRegistrar coderRegistrar = new ModelCoderRegistrar(); + assertThat(coderRegistrar.isKnownCoder(voidCoder, p.getOptions()), is(false)); p.apply("impulse", Impulse.create()) .apply( ParDo.of( diff --git a/runners/portability/java/build.gradle b/runners/portability/java/build.gradle index 6e3b431e802b..aa147e8426c4 100644 --- a/runners/portability/java/build.gradle +++ b/runners/portability/java/build.gradle @@ -214,6 +214,11 @@ def createUlrValidatesRunnerTask = { name, environmentType, dockerImageTask = "" // TODO(https://github.com/apache/beam/issues/31231) excludeTestsMatching 'org.apache.beam.sdk.transforms.RedistributeTest.testRedistributePreservesMetadata' + // TODO(https://github.com/apache/beam/issues/33859): Failed with "KeyError: 'beam:coder:schema:v1'". + // New schema coder urn is not yet supported in runners other than dataflow + excludeTestsMatching 'org.apache.beam.sdk.transforms.PerKeyOrderingTest.testMultipleStatefulOrderingWithShuffle' + excludeTestsMatching 'org.apache.beam.sdk.transforms.PerKeyOrderingTest.testMultipleStatefulOrderingWithoutShuffle' + for (String test : sickbayTests) { excludeTestsMatching test } diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/CoderTranslation.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/CoderTranslation.java index 22859dc68b93..2cc4bf0c6a03 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/CoderTranslation.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/CoderTranslation.java @@ -25,11 +25,13 @@ import org.apache.beam.model.pipeline.v1.RunnerApi; import org.apache.beam.model.pipeline.v1.RunnerApi.FunctionSpec; import org.apache.beam.sdk.coders.Coder; +import org.apache.beam.sdk.options.PipelineOptions; import org.apache.beam.sdk.util.SerializableUtils; import org.apache.beam.vendor.grpc.v1p69p0.com.google.protobuf.ByteString; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.BiMap; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableBiMap; +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.checkerframework.checker.nullness.qual.MonotonicNonNull; import org.checkerframework.dataflow.qual.Deterministic; @@ -62,6 +64,8 @@ private static class DefaultTranslationContext implements TranslationContext {} private static @MonotonicNonNull BiMap, String> knownCoderUrns; + private static @MonotonicNonNull List coderTranslatorRegistrars; + private static @MonotonicNonNull Map, CoderTranslator> knownTranslators; @@ -80,6 +84,53 @@ static BiMap, String> getKnownCoderUrns() { return knownCoderUrns; } + private static void initializeCoderTranslatorRegistrars() { + ImmutableList.Builder registrars = ImmutableList.builder(); + for (CoderTranslatorRegistrar coderTranslatorRegistrar : + ServiceLoader.load(CoderTranslatorRegistrar.class)) { + registrars.add(coderTranslatorRegistrar); + } + coderTranslatorRegistrars = registrars.build(); + } + + static boolean isKnownCoder(Coder coder, PipelineOptions options) { + if (coderTranslatorRegistrars == null) { + initializeCoderTranslatorRegistrars(); + } + for (CoderTranslatorRegistrar registrar : coderTranslatorRegistrars) { + if (registrar.isKnownCoder(coder, options)) { + return true; + } + } + return false; + } + + static CoderTranslator getCoderTranslator(Class coderClass) { + if (coderTranslatorRegistrars == null) { + initializeCoderTranslatorRegistrars(); + } + for (CoderTranslatorRegistrar registrar : coderTranslatorRegistrars) { + CoderTranslator translator = registrar.getCoderTranslator(coderClass); + if (translator != null) { + return translator; + } + } + return null; + } + + static Class getCoderForUrn(String coderUrn) { + if (coderTranslatorRegistrars == null) { + initializeCoderTranslatorRegistrars(); + } + for (CoderTranslatorRegistrar registrar : coderTranslatorRegistrars) { + Class coder = registrar.getCoderForUrn(coderUrn); + if (coder != null) { + return coder; + } + } + return null; + } + @VisibleForTesting @Deterministic static Map, CoderTranslator> getKnownTranslators() { @@ -107,7 +158,7 @@ public static RunnerApi.MessageWithComponents toProto(Coder coder) throws IOE public static RunnerApi.Coder toProto(Coder coder, SdkComponents components) throws IOException { - if (getKnownCoderUrns().containsKey(coder.getClass())) { + if (isKnownCoder(coder, components.getPipelineOptions())) { return toKnownCoder(coder, components); } @@ -129,7 +180,10 @@ private static RunnerApi.Coder toUnknownCoderWrapper(UnknownCoderWrapper coder) private static RunnerApi.Coder toKnownCoder(Coder coder, SdkComponents components) throws IOException { - CoderTranslator translator = getKnownTranslators().get(coder.getClass()); + CoderTranslator translator = getCoderTranslator(coder.getClass()); + if (translator == null) { + throw new IOException("Unable to find CoderTranslator for known Coder"); + } List componentIds = registerComponents(coder, translator, components); return RunnerApi.Coder.newBuilder() .addAllComponentCoderIds(componentIds) @@ -186,8 +240,8 @@ private static Coder fromKnownCoder( components.getComponents().getCodersOrThrow(componentId), components, context); coderComponents.add(innerCoder); } - Class coderType = getKnownCoderUrns().inverse().get(coderUrn); - CoderTranslator translator = getKnownTranslators().get(coderType); + Class coderType = getCoderForUrn(coderUrn); + CoderTranslator translator = getCoderTranslator(coderType); if (translator != null) { return translator.fromComponents( coderComponents, coder.getSpec().getPayload().toByteArray(), context); diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/CoderTranslator.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/CoderTranslator.java index 3d89c4c7ff4a..78f5b61c0f0e 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/CoderTranslator.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/CoderTranslator.java @@ -28,6 +28,7 @@ * additional payload, which is not currently supported. This exists as a temporary measure. */ public interface CoderTranslator> { + /** Extract all component {@link Coder coders} within a coder. */ List> getComponents(T from); 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 b69d0290de52..44e8c2956aee 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 @@ -19,6 +19,8 @@ import java.util.Map; import org.apache.beam.sdk.coders.Coder; +import org.apache.beam.sdk.options.PipelineOptions; +import org.checkerframework.checker.nullness.qual.Nullable; /** A registrar of {@link Coder} URNs to the associated {@link CoderTranslator}. */ @SuppressWarnings({ @@ -34,4 +36,18 @@ public interface CoderTranslatorRegistrar { /** Returns a mapping of URN to {@link CoderTranslator}. */ Map, CoderTranslator> getCoderTranslators(); + + /** + * 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); + + /** Returns the CoderTranslator to use for this Coder, or null if the Coder is not known. */ + @Nullable + CoderTranslator getCoderTranslator(Class coderClass); + + /** Returns the Coder to use for the given Urn, or null if the Urn is for an unknown Coder. */ + @Nullable + Class getCoderForUrn(String coderUrn); } diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/CoderTranslators.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/CoderTranslators.java index 84a90721a983..a847bf780dff 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/CoderTranslators.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/CoderTranslators.java @@ -19,6 +19,7 @@ import static org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkArgument; +import java.io.IOException; import java.util.Collections; import java.util.List; import org.apache.beam.model.pipeline.v1.SchemaApi; @@ -30,12 +31,19 @@ import org.apache.beam.sdk.coders.RowCoder; import org.apache.beam.sdk.coders.TimestampPrefixingWindowCoder; import org.apache.beam.sdk.schemas.Schema; +import org.apache.beam.sdk.schemas.SchemaCoder; import org.apache.beam.sdk.schemas.SchemaTranslation; +import org.apache.beam.sdk.transforms.SerializableFunction; import org.apache.beam.sdk.transforms.windowing.BoundedWindow; import org.apache.beam.sdk.util.InstanceBuilder; +import org.apache.beam.sdk.util.SerializableUtils; import org.apache.beam.sdk.util.ShardedKey; +import org.apache.beam.sdk.util.construction.CoderTranslation.TranslationContext; +import org.apache.beam.sdk.values.Row; +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.grpc.v1p69p0.com.google.protobuf.ByteString; import org.apache.beam.vendor.grpc.v1p69p0.com.google.protobuf.InvalidProtocolBufferException; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList; @@ -177,6 +185,77 @@ public RowCoder fromComponents( }; } + static CoderTranslator> schema() { + return new CoderTranslator>() { + private static final String TO_ROW_FUNCTION_URN = "beam:torowfn:javasdk:v1"; + private static final String FROM_ROW_FUNCTION_URN = "beam:fromrowfn:javasdk:v1"; + private static final String TYPE_DESCRIPTOR_URN = "beam:typedescriptor:javasdk:v1"; + + @Override + public ImmutableList> getComponents(SchemaCoder from) { + return ImmutableList.of(); + } + + @Override + public byte[] getPayload(SchemaCoder from) { + SchemaApi.SchemaCoderPayload.Builder payload = SchemaApi.SchemaCoderPayload.newBuilder(); + payload.setSchema(SchemaTranslation.schemaToProto(from.getSchema(), true)); + payload + .getToRowFnBuilder() + .setUrn(TO_ROW_FUNCTION_URN) + .setPayload( + ByteString.copyFrom( + SerializableUtils.serializeToByteArray(from.getToRowFunction()))); + payload + .getFromRowFnBuilder() + .setUrn(FROM_ROW_FUNCTION_URN) + .setPayload( + ByteString.copyFrom( + SerializableUtils.serializeToByteArray(from.getFromRowFunction()))); + payload + .addAdditionalCoderInfosBuilder() + .setUrn(TYPE_DESCRIPTOR_URN) + .setPayload( + ByteString.copyFrom( + SerializableUtils.serializeToByteArray(from.getEncodedTypeDescriptor()))); + return payload.build().toByteArray(); + } + + @Override + public SchemaCoder fromComponents( + List> components, byte[] payload, TranslationContext context) { + checkArgument( + components.isEmpty(), "Expected empty component list, but received: %s", components); + try { + SchemaApi.SchemaCoderPayload schemaCoderPayload = + SchemaApi.SchemaCoderPayload.parseFrom(payload); + if (schemaCoderPayload.getAdditionalCoderInfosCount() == 0) { + throw new IllegalArgumentException("Missing serialized typeDescriptor"); + } + TypeDescriptor typeDescriptor = + (TypeDescriptor) + SerializableUtils.deserializeFromByteArray( + schemaCoderPayload.getAdditionalCoderInfos(0).getPayload().toByteArray(), + "typeDescriptor"); + SerializableFunction toRowFunction = + (SerializableFunction) + SerializableUtils.deserializeFromByteArray( + schemaCoderPayload.getToRowFn().getPayload().toByteArray(), "toRowFunction"); + SerializableFunction fromRowFunction = + (SerializableFunction) + SerializableUtils.deserializeFromByteArray( + schemaCoderPayload.getFromRowFn().getPayload().toByteArray(), + "fromRowFunction"); + + Schema schema = SchemaTranslation.schemaFromProto(schemaCoderPayload.getSchema()); + return SchemaCoder.of(schema, typeDescriptor, toRowFunction, fromRowFunction); + } catch (IOException | IllegalArgumentException e) { + throw new RuntimeException(e); + } + } + }; + } + static CoderTranslator> shardedKey() { return new SimpleStructuredCoderTranslator>() { @Override diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/ModelCoderRegistrar.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/ModelCoderRegistrar.java index 5b0d5aedd619..1f9f1eaafbed 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/ModelCoderRegistrar.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/ModelCoderRegistrar.java @@ -34,6 +34,9 @@ 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.PipelineOptions; +import org.apache.beam.sdk.options.StreamingOptions; +import org.apache.beam.sdk.schemas.SchemaCoder; import org.apache.beam.sdk.transforms.windowing.GlobalWindow; import org.apache.beam.sdk.transforms.windowing.IntervalWindow.IntervalWindowCoder; import org.apache.beam.sdk.util.ShardedKey; @@ -71,6 +74,7 @@ public class ModelCoderRegistrar implements CoderTranslatorRegistrar { ModelCoders.PARAM_WINDOWED_VALUE_CODER_URN) .put(DoubleCoder.class, ModelCoders.DOUBLE_CODER_URN) .put(RowCoder.class, ModelCoders.ROW_CODER_URN) + .put(SchemaCoder.class, ModelCoders.SCHEMA_CODER_URN) .put(ShardedKey.Coder.class, ModelCoders.SHARDED_KEY_CODER_URN) .put(TimestampPrefixingWindowCoder.class, ModelCoders.CUSTOM_WINDOW_CODER_URN) .put(NullableCoder.class, ModelCoders.NULLABLE_CODER_URN) @@ -96,6 +100,7 @@ public class ModelCoderRegistrar implements CoderTranslatorRegistrar { CoderTranslators.paramWindowedValue()) .put(DoubleCoder.class, CoderTranslators.atomic(DoubleCoder.class)) .put(RowCoder.class, CoderTranslators.row()) + .put(SchemaCoder.class, CoderTranslators.schema()) .put(ShardedKey.Coder.class, CoderTranslators.shardedKey()) .put(TimestampPrefixingWindowCoder.class, CoderTranslators.timestampPrefixingWindow()) .put(NullableCoder.class, CoderTranslators.nullable()) @@ -123,10 +128,6 @@ public class ModelCoderRegistrar implements CoderTranslatorRegistrar { Coder.class.getSimpleName()); } - public static boolean isKnownCoder(Coder coder) { - return BEAM_MODEL_CODER_URNS.containsKey(coder.getClass()); - } - @Override public Map, String> getCoderURNs() { return BEAM_MODEL_CODER_URNS; @@ -136,4 +137,23 @@ public Map, String> getCoderURNs() { public Map, CoderTranslator> getCoderTranslators() { return BEAM_MODEL_CODERS; } + + @Override + public boolean isKnownCoder(Coder coder, PipelineOptions options) { + if (coder.getClass() == SchemaCoder.class + && StreamingOptions.updateCompatibilityVersionLessThan(options, "2.74")) { + return false; + } + return BEAM_MODEL_CODER_URNS.containsKey(coder.getClass()); + } + + @Override + public CoderTranslator getCoderTranslator(Class coderClass) { + return BEAM_MODEL_CODERS.getOrDefault(coderClass, null); + } + + @Override + public Class getCoderForUrn(String coderUrn) { + return BEAM_MODEL_CODER_URNS.inverse().getOrDefault(coderUrn, null); + } } diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/ModelCoders.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/ModelCoders.java index 7b7546aceb61..5059cc1c6b83 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/ModelCoders.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/ModelCoders.java @@ -61,6 +61,7 @@ private ModelCoders() {} getUrn(StandardCoders.Enum.PARAM_WINDOWED_VALUE); public static final String ROW_CODER_URN = getUrn(StandardCoders.Enum.ROW); + public static final String SCHEMA_CODER_URN = getUrn(StandardCoders.Enum.SCHEMA); public static final String STATE_BACKED_ITERABLE_CODER_URN = "beam:coder:state_backed_iterable:v1"; @@ -90,6 +91,7 @@ private ModelCoders() {} WINDOWED_VALUE_CODER_URN, DOUBLE_CODER_URN, ROW_CODER_URN, + SCHEMA_CODER_URN, PARAM_WINDOWED_VALUE_CODER_URN, STATE_BACKED_ITERABLE_CODER_URN, SHARDED_KEY_CODER_URN, diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/RehydratedComponents.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/RehydratedComponents.java index f79696214368..64c7898a37b9 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/RehydratedComponents.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/RehydratedComponents.java @@ -189,6 +189,7 @@ public SdkComponents getSdkComponents(Collection requirements) { windowingStrategies.asMap(), coders.asMap(), Collections.emptyMap(), - requirements); + requirements, + pipeline.getOptions()); } } diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/SdkComponents.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/SdkComponents.java index 446697f24a81..6288649aba3d 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/SdkComponents.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/SdkComponents.java @@ -63,6 +63,7 @@ public class SdkComponents { private final BiMap environmentIds = HashBiMap.create(); private final BiMap coderProtoToId = HashBiMap.create(); private final Set requirements; + private final PipelineOptions pipelineOptions; private final Set reservedIds = new HashSet<>(); @@ -71,17 +72,7 @@ public class SdkComponents { /** Create a new {@link SdkComponents} with no components. */ public static SdkComponents create() { - return new SdkComponents(RunnerApi.Components.getDefaultInstance(), null, ""); - } - - /** - * Create new {@link SdkComponents} importing all items from provided {@link Components} object. - * - *

WARNING: This action might cause some of duplicate items created. - */ - public static SdkComponents create( - RunnerApi.Components components, Collection requirements) { - return new SdkComponents(components, requirements, ""); + return new SdkComponents(RunnerApi.Components.getDefaultInstance(), null, "", null); } /*package*/ static SdkComponents create( @@ -91,8 +82,9 @@ public static SdkComponents create( Map> windowingStrategies, Map> coders, Map environments, - Collection requirements) { - SdkComponents sdkComponents = SdkComponents.create(components, requirements); + Collection requirements, + PipelineOptions pipelineOptions) { + SdkComponents sdkComponents = new SdkComponents(components, requirements, "", pipelineOptions); sdkComponents.transformIds.inverse().putAll(transforms); sdkComponents.pCollectionIds.inverse().putAll(pCollections); sdkComponents.windowingStrategyIds.inverse().putAll(windowingStrategies); @@ -103,19 +95,28 @@ public static SdkComponents create( public static SdkComponents create(PipelineOptions options) { SdkComponents sdkComponents = - new SdkComponents(RunnerApi.Components.getDefaultInstance(), null, ""); + new SdkComponents(RunnerApi.Components.getDefaultInstance(), null, "", options); PortablePipelineOptions portablePipelineOptions = options.as(PortablePipelineOptions.class); sdkComponents.registerEnvironment( Environments.createOrGetDefaultEnvironment(portablePipelineOptions)); return sdkComponents; } + public static SdkComponents create(PipelineOptions options, Environment environment) { + SdkComponents sdkComponents = + new SdkComponents(RunnerApi.Components.getDefaultInstance(), null, "", options); + sdkComponents.registerEnvironment(environment); + return sdkComponents; + } + private SdkComponents( @Nullable Components components, @Nullable Collection requirements, - String newIdPrefix) { + String newIdPrefix, + @Nullable PipelineOptions pipelineOptions) { this.newIdPrefix = newIdPrefix; this.requirements = new HashSet<>(); + this.pipelineOptions = pipelineOptions; if (components == null) { if (requirements != null) { @@ -153,7 +154,7 @@ public void mergeFrom( */ public SdkComponents withNewIdPrefix(String newIdPrefix) { SdkComponents sdkComponents = - new SdkComponents(componentsBuilder.build(), requirements, newIdPrefix); + new SdkComponents(componentsBuilder.build(), requirements, newIdPrefix, pipelineOptions); sdkComponents.transformIds.putAll(transformIds); sdkComponents.pCollectionIds.putAll(pCollectionIds); sdkComponents.windowingStrategyIds.putAll(windowingStrategyIds); @@ -174,7 +175,7 @@ public String registerPTransform( throws IOException { String name = getApplicationName(appliedPTransform); // If this transform is present in the components, nothing to do. return the existing name. - // Otherwise the transform must be translated and added to the components. + // Otherwise, the transform must be translated and added to the components. if (componentsBuilder.getTransformsOrDefault(name, null) != null) { return name; } @@ -375,4 +376,8 @@ public RunnerApi.Components toComponents() { public Collection requirements() { return ImmutableSet.copyOf(requirements); } + + public PipelineOptions getPipelineOptions() { + return pipelineOptions; + } } 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 b8f92ff0053e..1ec0a74f5be1 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 @@ -22,6 +22,7 @@ import static org.hamcrest.Matchers.hasItems; import static org.hamcrest.Matchers.not; +import com.google.auto.value.AutoValue; import java.io.IOException; import java.io.InputStream; import java.io.OutputStream; @@ -45,14 +46,20 @@ 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.schemas.AutoValueSchema; +import org.apache.beam.sdk.schemas.NoSuchSchemaException; import org.apache.beam.sdk.schemas.Schema; import org.apache.beam.sdk.schemas.Schema.Field; import org.apache.beam.sdk.schemas.Schema.FieldType; +import org.apache.beam.sdk.schemas.SchemaCoder; +import org.apache.beam.sdk.schemas.SchemaRegistry; +import org.apache.beam.sdk.schemas.annotations.DefaultSchema; import org.apache.beam.sdk.schemas.logicaltypes.FixedBytes; import org.apache.beam.sdk.transforms.windowing.GlobalWindow; import org.apache.beam.sdk.transforms.windowing.IntervalWindow.IntervalWindowCoder; import org.apache.beam.sdk.util.ShardedKey; import org.apache.beam.sdk.util.construction.CoderTranslation.TranslationContext; +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; @@ -70,6 +77,34 @@ "rawtypes", // TODO(https://github.com/apache/beam/issues/20447) }) public class CoderTranslationTest { + @AutoValue + @DefaultSchema(AutoValueSchema.class) + public abstract static class SimpleAutoValue { + public abstract String getString(); + + public abstract int getInt32(); + + public abstract long getInt64(); + + public static SimpleAutoValue of(String string, Integer int32, Long int64) { + return new AutoValue_CoderTranslationTest_SimpleAutoValue(string, int32, int64); + } + } + + private static final SchemaRegistry REGISTRY = SchemaRegistry.createDefault(); + + private static SchemaCoder schemaCoderFrom(TypeDescriptor typeDescriptor) { + try { + return SchemaCoder.of( + REGISTRY.getSchema(typeDescriptor), + typeDescriptor, + REGISTRY.getToRowFunction(typeDescriptor), + REGISTRY.getFromRowFunction(typeDescriptor)); + } catch (NoSuchSchemaException e) { + throw new RuntimeException(e); + } + } + private static final Set> KNOWN_CODERS = ImmutableSet.>builder() .add(ByteArrayCoder.of()) @@ -94,6 +129,7 @@ public class CoderTranslationTest { Field.of("array", FieldType.array(FieldType.STRING)), Field.of("map", FieldType.map(FieldType.STRING, FieldType.INT32)), Field.of("bar", FieldType.logicalType(FixedBytes.of(123)))))) + .add(schemaCoderFrom(TypeDescriptor.of(SimpleAutoValue.class))) .add(ShardedKey.Coder.of(StringUtf8Coder.of())) .add(TimestampPrefixingWindowCoder.of(IntervalWindowCoder.of())) .add(NullableCoder.of(ByteArrayCoder.of())) @@ -127,7 +163,7 @@ public void validateKnownCoders() { } @Test - public void validateCoderTranslators() { + public void validateModelCoderTranslators() { assertThat( "Every Model Coder must have a Translator", new ModelCoderRegistrar().getCoderURNs().keySet(), diff --git a/sdks/java/expansion-service/src/main/java/org/apache/beam/sdk/expansion/service/ExpansionService.java b/sdks/java/expansion-service/src/main/java/org/apache/beam/sdk/expansion/service/ExpansionService.java index 6ecb029c5d97..c93de2014798 100644 --- a/sdks/java/expansion-service/src/main/java/org/apache/beam/sdk/expansion/service/ExpansionService.java +++ b/sdks/java/expansion-service/src/main/java/org/apache/beam/sdk/expansion/service/ExpansionService.java @@ -600,7 +600,7 @@ private Map loadRegisteredTransforms() { pipeline.getOptions().as(ExperimentalOptions.class), "use_sdf_read"); } else { LOG.warn( - "Using use_depreacted_read in portable runners is runner-dependent. The " + "Using use_deprecated_read in portable runners is runner-dependent. The " + "ExpansionService will respect that, but if your runner does not have support for " + "native Read transform, your Pipeline will fail during Pipeline submission."); } diff --git a/sdks/java/extensions/avro/src/main/java/org/apache/beam/sdk/extensions/avro/AvroGenericCoderRegistrar.java b/sdks/java/extensions/avro/src/main/java/org/apache/beam/sdk/extensions/avro/AvroGenericCoderRegistrar.java index 14ab48f66699..8bd18fd8e250 100644 --- a/sdks/java/extensions/avro/src/main/java/org/apache/beam/sdk/extensions/avro/AvroGenericCoderRegistrar.java +++ b/sdks/java/extensions/avro/src/main/java/org/apache/beam/sdk/extensions/avro/AvroGenericCoderRegistrar.java @@ -21,9 +21,11 @@ import java.util.Map; import org.apache.beam.sdk.coders.Coder; import org.apache.beam.sdk.extensions.avro.coders.AvroGenericCoder; +import org.apache.beam.sdk.options.PipelineOptions; import org.apache.beam.sdk.util.construction.CoderTranslator; import org.apache.beam.sdk.util.construction.CoderTranslatorRegistrar; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap; +import org.checkerframework.checker.nullness.qual.Nullable; /** Coder registrar for AvroGenericCoder. */ @AutoService(CoderTranslatorRegistrar.class) @@ -42,4 +44,20 @@ public Map, String> getCoderURNs() { public Map, CoderTranslator> getCoderTranslators() { return ImmutableMap.of(AvroGenericCoder.class, new AvroGenericCoderTranslator()); } + + @Override + public boolean isKnownCoder(Coder coder, PipelineOptions options) { + return coder.getClass() == AvroGenericCoder.class; + } + + @Override + public @Nullable CoderTranslator getCoderTranslator( + Class coderClass) { + return coderClass == AvroGenericCoder.class ? new AvroGenericCoderTranslator() : null; + } + + @Override + public @Nullable Class getCoderForUrn(String coderUrn) { + return AVRO_CODER_URN.equals(coderUrn) ? AvroGenericCoder.class : null; + } } diff --git a/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/state/StateBackedIterable.java b/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/state/StateBackedIterable.java index ef8d69bc1ec3..42a6f8d11c2a 100644 --- a/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/state/StateBackedIterable.java +++ b/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/state/StateBackedIterable.java @@ -38,6 +38,7 @@ import org.apache.beam.sdk.coders.IterableLikeCoder; import org.apache.beam.sdk.fn.stream.PrefetchableIterable; import org.apache.beam.sdk.fn.stream.PrefetchableIterators; +import org.apache.beam.sdk.options.PipelineOptions; import org.apache.beam.sdk.util.BufferedElementCountingOutputStream; import org.apache.beam.sdk.util.VarInt; import org.apache.beam.sdk.util.common.ElementByteSizeObservableIterable; @@ -52,6 +53,7 @@ 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.Iterables; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.io.ByteStreams; +import org.checkerframework.checker.nullness.qual.Nullable; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -300,6 +302,26 @@ public Map, String> getCoderUR getCoderTranslators() { return ImmutableMap.of(StateBackedIterable.Coder.class, new Translator()); } + + @Override + public boolean isKnownCoder( + org.apache.beam.sdk.coders.Coder coder, PipelineOptions options) { + return coder.getClass() == StateBackedIterable.Coder.class; + } + + @Override + public @Nullable CoderTranslator getCoderTranslator( + Class coderClass) { + return coderClass == StateBackedIterable.Coder.class ? new Translator() : null; + } + + @Override + public @Nullable Class getCoderForUrn( + String coderUrn) { + return STATE_BACKED_ITERABLE_CODER_URN.equals(coderUrn) + ? StateBackedIterable.Coder.class + : null; + } } /** From f09683548aa22e66729f4d9e57420048d51d8923 Mon Sep 17 00:00:00 2001 From: Andrew Crites Date: Mon, 3 Aug 2026 09:14:32 -0700 Subject: [PATCH 4/5] Changes usage of known Schema coder to depend on use_known_schema_coder experiment instead of version compatibility number. --- .../dataflow/DataflowPipelineTranslatorTest.java | 10 ++++++---- .../sdk/util/construction/ModelCoderRegistrar.java | 4 ++-- 2 files changed, 8 insertions(+), 6 deletions(-) diff --git a/runners/google-cloud-dataflow-java/src/test/java/org/apache/beam/runners/dataflow/DataflowPipelineTranslatorTest.java b/runners/google-cloud-dataflow-java/src/test/java/org/apache/beam/runners/dataflow/DataflowPipelineTranslatorTest.java index 953f7a638ede..63c019042c53 100644 --- a/runners/google-cloud-dataflow-java/src/test/java/org/apache/beam/runners/dataflow/DataflowPipelineTranslatorTest.java +++ b/runners/google-cloud-dataflow-java/src/test/java/org/apache/beam/runners/dataflow/DataflowPipelineTranslatorTest.java @@ -89,6 +89,7 @@ import org.apache.beam.sdk.io.gcp.bigquery.BigQueryIO; import org.apache.beam.sdk.io.gcp.bigquery.BigQueryIO.Write.CreateDisposition; import org.apache.beam.sdk.io.range.OffsetRange; +import org.apache.beam.sdk.options.ExperimentalOptions; import org.apache.beam.sdk.options.PipelineOptions; import org.apache.beam.sdk.options.PipelineOptionsFactory; import org.apache.beam.sdk.options.StreamingOptions; @@ -1902,21 +1903,22 @@ public SimpleAutoValue apply(byte[] input) { })) .apply(Window.into(FixedWindows.of(Duration.standardMinutes(1)))); { + // Without the experiment, SchemaCoder is treated as unknown. SdkComponents sdkComponents = createSdkComponents(options); RunnerApi.Pipeline pipelineProto = PipelineTranslation.toProto(pipeline, sdkComponents, true); Map coders = pipelineProto.getComponents().getCodersMap(); assertTrue(coders.containsKey("SchemaCoder")); - assertEquals("beam:coder:schema:v1", coders.get("SchemaCoder").getSpec().getUrn()); + assertEquals("beam:coders:javasdk:0.1", coders.get("SchemaCoder").getSpec().getUrn()); } - // Prior to version 2.74, SchemaCoders are translated as custom java coders. { - options.as(StreamingOptions.class).setUpdateCompatibilityVersion("2.73"); + // Add the experiment to get the known coder urn instead. + ExperimentalOptions.addExperiment(options, "use_known_schema_coder"); SdkComponents sdkComponents = createSdkComponents(options); RunnerApi.Pipeline pipelineProto = PipelineTranslation.toProto(pipeline, sdkComponents, true); Map coders = pipelineProto.getComponents().getCodersMap(); assertTrue(coders.containsKey("SchemaCoder")); - assertEquals("beam:coders:javasdk:0.1", coders.get("SchemaCoder").getSpec().getUrn()); + assertEquals("beam:coder:schema:v1", coders.get("SchemaCoder").getSpec().getUrn()); } } } diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/ModelCoderRegistrar.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/ModelCoderRegistrar.java index 1f9f1eaafbed..39c81c6decb5 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/ModelCoderRegistrar.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/ModelCoderRegistrar.java @@ -34,8 +34,8 @@ 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.ExperimentalOptions; import org.apache.beam.sdk.options.PipelineOptions; -import org.apache.beam.sdk.options.StreamingOptions; import org.apache.beam.sdk.schemas.SchemaCoder; import org.apache.beam.sdk.transforms.windowing.GlobalWindow; import org.apache.beam.sdk.transforms.windowing.IntervalWindow.IntervalWindowCoder; @@ -141,7 +141,7 @@ public Map, CoderTranslator> getCoderTra @Override public boolean isKnownCoder(Coder coder, PipelineOptions options) { if (coder.getClass() == SchemaCoder.class - && StreamingOptions.updateCompatibilityVersionLessThan(options, "2.74")) { + && !ExperimentalOptions.hasExperiment(options, "use_known_schema_coder")) { return false; } return BEAM_MODEL_CODER_URNS.containsKey(coder.getClass()); From 35e0ec8cba4651a356b2686d69dce4612b094e99 Mon Sep 17 00:00:00 2001 From: Andrew Crites Date: Mon, 3 Aug 2026 09:39:02 -0700 Subject: [PATCH 5/5] Moves CHANGES to correct release version and gets rid of extraneous changes to this PR. --- .../beam_PostCommit_Java_ValidatesRunner_ULR.json | 3 +-- CHANGES.md | 6 +++--- runners/portability/java/build.gradle | 5 ----- .../apache/beam/sdk/util/construction/CoderTranslator.java | 1 - .../apache/beam/sdk/expansion/service/ExpansionService.java | 2 +- 5 files changed, 5 insertions(+), 12 deletions(-) diff --git a/.github/trigger_files/beam_PostCommit_Java_ValidatesRunner_ULR.json b/.github/trigger_files/beam_PostCommit_Java_ValidatesRunner_ULR.json index fbd81891f93b..6e2f429dd24e 100644 --- a/.github/trigger_files/beam_PostCommit_Java_ValidatesRunner_ULR.json +++ b/.github/trigger_files/beam_PostCommit_Java_ValidatesRunner_ULR.json @@ -2,6 +2,5 @@ "https://github.com/apache/beam/pull/34902": "Introducing OutputBuilder", "comment": "Modify this file in a trivial way to cause this test suite to run", "https://github.com/apache/beam/pull/31156": "noting that PR #31156 should run this test", - "https://github.com/apache/beam/pull/35159": "moving WindowedValue and making an interface", - "https://github.com/apache/beam/pull/38497": "sickbay two failed tests" + "https://github.com/apache/beam/pull/35159": "moving WindowedValue and making an interface" } diff --git a/CHANGES.md b/CHANGES.md index e4db58147260..6baf3b470e55 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -89,6 +89,9 @@ * Fixed unbounded checkpoint state growth for splittable DoFns that self-checkpoint on the portable Flink runner (Java) ([#27648](https://github.com/apache/beam/issues/27648)). * Improved Java pipeline performance by avoiding repeated `DoFn` type descriptor resolution when creating cached invokers ([#39309](https://github.com/apache/beam/issues/39309)). * (Python) Fixed a memory leak in Python SDK caused by storing exceptions with potentially large stack frames in a cache ([#39406](https://github.com/apache/beam/issues/39406)). +* Portable Java SDK now encodes SchemaCoders in a portable way ([#34672](https://github.com/apache/beam/issues/34672)). + - The new coder translation is protected by the "use_known_schema_coder" PipelineOptions experiment. + - Fixes ([#36496](https://github.com/apache/beam/issues/36496)), ([#30276](https://github.com/apache/beam/issues/30276)), ([#29245](https://github.com/apache/beam/issues/29245)). ## Security Fixes @@ -173,9 +176,6 @@ ## Breaking Changes -* Portable Java SDK now encodes SchemaCoders in a portable way ([#34672](https://github.com/apache/beam/issues/34672)). - - Original custom Java coder encoding can still be obtained using [StreamingOptions.setUpdateCompatibilityVersion("2.73")](https://github.com/apache/beam/blob/2cf0930e7ae1aa389c26ce6639b584877a3e31d9/sdks/java/core/src/main/java/org/apache/beam/sdk/options/StreamingOptions.java#L47) ([#34672](https://github.com/apache/beam/issues/34672)). - - Fixes ([#36496](https://github.com/apache/beam/issues/36496)), ([#30276](https://github.com/apache/beam/issues/30276)), ([#29245](https://github.com/apache/beam/issues/29245)). * (Python) Made Beartype the default fallback type checking tool. This can be disabled with the `--disable_beartype` pipeline option. ([#38275](https://github.com/apache/beam/issues/38275)) ## Deprecations diff --git a/runners/portability/java/build.gradle b/runners/portability/java/build.gradle index aa147e8426c4..6e3b431e802b 100644 --- a/runners/portability/java/build.gradle +++ b/runners/portability/java/build.gradle @@ -214,11 +214,6 @@ def createUlrValidatesRunnerTask = { name, environmentType, dockerImageTask = "" // TODO(https://github.com/apache/beam/issues/31231) excludeTestsMatching 'org.apache.beam.sdk.transforms.RedistributeTest.testRedistributePreservesMetadata' - // TODO(https://github.com/apache/beam/issues/33859): Failed with "KeyError: 'beam:coder:schema:v1'". - // New schema coder urn is not yet supported in runners other than dataflow - excludeTestsMatching 'org.apache.beam.sdk.transforms.PerKeyOrderingTest.testMultipleStatefulOrderingWithShuffle' - excludeTestsMatching 'org.apache.beam.sdk.transforms.PerKeyOrderingTest.testMultipleStatefulOrderingWithoutShuffle' - for (String test : sickbayTests) { excludeTestsMatching test } diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/CoderTranslator.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/CoderTranslator.java index 78f5b61c0f0e..3d89c4c7ff4a 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/CoderTranslator.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/CoderTranslator.java @@ -28,7 +28,6 @@ * additional payload, which is not currently supported. This exists as a temporary measure. */ public interface CoderTranslator> { - /** Extract all component {@link Coder coders} within a coder. */ List> getComponents(T from); diff --git a/sdks/java/expansion-service/src/main/java/org/apache/beam/sdk/expansion/service/ExpansionService.java b/sdks/java/expansion-service/src/main/java/org/apache/beam/sdk/expansion/service/ExpansionService.java index c93de2014798..6ecb029c5d97 100644 --- a/sdks/java/expansion-service/src/main/java/org/apache/beam/sdk/expansion/service/ExpansionService.java +++ b/sdks/java/expansion-service/src/main/java/org/apache/beam/sdk/expansion/service/ExpansionService.java @@ -600,7 +600,7 @@ private Map loadRegisteredTransforms() { pipeline.getOptions().as(ExperimentalOptions.class), "use_sdf_read"); } else { LOG.warn( - "Using use_deprecated_read in portable runners is runner-dependent. The " + "Using use_depreacted_read in portable runners is runner-dependent. The " + "ExpansionService will respect that, but if your runner does not have support for " + "native Read transform, your Pipeline will fail during Pipeline submission."); }