From 4d87d3c803f23dd016946524fb30c42bfdf30dd0 Mon Sep 17 00:00:00 2001 From: Radek Stankiewicz Date: Tue, 9 Sep 2025 15:17:05 +0200 Subject: [PATCH 1/9] encode empty metadata within WindowedValue --- .../worker/StreamingDataflowWorker.java | 5 ++ .../worker/UngroupedWindmillReader.java | 1 + .../worker/WindmillKeyedWorkItem.java | 2 +- .../runners/dataflow/worker/WindmillSink.java | 21 ++++++- .../StreamingGroupAlsoByWindowFnsTest.java | 8 ++- ...ngGroupAlsoByWindowsReshuffleDoFnTest.java | 8 ++- .../worker/WindmillKeyedWorkItemTest.java | 9 ++- .../sdk/transforms/windowing/PaneInfo.java | 59 +++++++++++++++---- .../beam/sdk/values/WindowedValues.java | 23 +++++++- .../transforms/windowing/PaneInfoTest.java | 20 +++++++ .../beam/sdk/util/WindowedValueTest.java | 27 +++++++++ 11 files changed, 162 insertions(+), 21 deletions(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java index 2a4b111af225..59269f8ea55e 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java @@ -110,6 +110,7 @@ import org.apache.beam.sdk.io.gcp.bigquery.BigQuerySinkMetrics; import org.apache.beam.sdk.metrics.MetricsEnvironment; import org.apache.beam.sdk.util.construction.CoderTranslation; +import org.apache.beam.sdk.values.WindowedValues; 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.base.Preconditions; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.cache.CacheStats; @@ -163,6 +164,7 @@ public final class StreamingDataflowWorker { private static final Random CLIENT_ID_GENERATOR = new Random(); private static final String CHANNELZ_PATH = "/channelz"; private static final String BEAM_FN_API_EXPERIMENT = "beam_fn_api"; + private static final String ELEMENT_METADATA_SUPPORTED_EXPERIMENT = "element_metadata_supported"; private static final String STREAMING_ENGINE_USE_JOB_SETTINGS_FOR_HEARTBEAT_POOL_EXPERIMENT = "streaming_engine_use_job_settings_for_heartbeat_pool"; // Experiment make the monitor within BoundedQueueExecutor fair @@ -815,6 +817,9 @@ public static void main(String[] args) throws Exception { validateWorkerOptions(options); CoderTranslation.verifyModelCodersRegistered(); + if (DataflowRunner.hasExperiment(options, ELEMENT_METADATA_SUPPORTED_EXPERIMENT)) { + WindowedValues.FullWindowedValueCoder.setMetadataSupported(); + } LOG.debug("Creating StreamingDataflowWorker from options: {}", options); StreamingDataflowWorker worker = StreamingDataflowWorker.fromOptions(options); diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/UngroupedWindmillReader.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/UngroupedWindmillReader.java index e031d1bb50eb..5df15e8d2109 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/UngroupedWindmillReader.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/UngroupedWindmillReader.java @@ -117,6 +117,7 @@ protected WindowedValue decodeMessage(Windmill.Message message) throws IOExce Collection windows = WindmillSink.decodeMetadataWindows(windowsCoder, message.getMetadata()); PaneInfo paneInfo = WindmillSink.decodeMetadataPane(message.getMetadata()); + WindmillSink.decodeAdditionalMetadata(windowsCoder, message.getMetadata()); if (valueCoder instanceof KvCoder) { KvCoder kvCoder = (KvCoder) valueCoder; InputStream key = context.getSerializedKey().newInput(); diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItem.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItem.java index cee4894e3d68..bae41a156967 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItem.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItem.java @@ -108,7 +108,7 @@ public Iterable> elementsIterable() { Collection windows = WindmillSink.decodeMetadataWindows(windowsCoder, message.getMetadata()); PaneInfo paneInfo = WindmillSink.decodeMetadataPane(message.getMetadata()); - + WindmillSink.decodeAdditionalMetadata(windowsCoder, message.getMetadata()); InputStream inputStream = message.getData().newInput(); ElemT value = valueCoder.decode(inputStream, Coder.Context.OUTER); return WindowedValues.of(value, timestamp, windows, paneInfo); diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java index 7cb6f2223472..4765eb74e7f5 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java @@ -26,6 +26,7 @@ import java.util.Collection; import java.util.HashMap; import java.util.Map; +import org.apache.beam.model.fnexecution.v1.BeamFnApi; import org.apache.beam.runners.dataflow.util.CloudObject; import org.apache.beam.runners.dataflow.worker.util.common.worker.Sink; import org.apache.beam.runners.dataflow.worker.windmill.Windmill; @@ -40,6 +41,7 @@ import org.apache.beam.sdk.values.ValueWithRecordId; import org.apache.beam.sdk.values.ValueWithRecordId.ValueWithRecordIdCoder; import org.apache.beam.sdk.values.WindowedValue; +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.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap; @@ -75,7 +77,8 @@ private static ByteString encodeMetadata( ByteStringOutputStream stream, Coder> windowsCoder, Collection windows, - PaneInfo paneInfo) + PaneInfo paneInfo, + BeamFnApi.Elements.ElementMetadata metadata) throws IOException { try { PaneInfoCoder.INSTANCE.encode(paneInfo, stream); @@ -101,12 +104,26 @@ public static PaneInfo decodeMetadataPane(ByteString metadata) throws IOExceptio return PaneInfoCoder.INSTANCE.decode(inStream); } + public static BeamFnApi.Elements.ElementMetadata decodeAdditionalMetadata( + Coder> windowsCoder, ByteString metadata) + throws IOException { + InputStream inStream = metadata.newInput(); + PaneInfo paneInfo = PaneInfoCoder.INSTANCE.decode(inStream); + windowsCoder.decode(inStream); + if (paneInfo.isElementMetadata() && WindowedValues.WindowedValueCoder.isMetadataSupported()) { + return BeamFnApi.Elements.ElementMetadata.parseDelimitedFrom(inStream); + } else { + // empty + return BeamFnApi.Elements.ElementMetadata.newBuilder().build(); + } + } + public static Collection decodeMetadataWindows( Coder> windowsCoder, ByteString metadata) throws IOException { InputStream inStream = metadata.newInput(); PaneInfoCoder.INSTANCE.decode(inStream); - return windowsCoder.decode(inStream, Coder.Context.OUTER); + return windowsCoder.decode(inStream); } /** A {@link SinkFactory.Registrar} for windmill sinks. */ diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingGroupAlsoByWindowFnsTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingGroupAlsoByWindowFnsTest.java index c89a031b3728..1182a2c0b9e9 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingGroupAlsoByWindowFnsTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingGroupAlsoByWindowFnsTest.java @@ -29,6 +29,7 @@ import java.util.Arrays; import java.util.Collection; import java.util.List; +import org.apache.beam.model.fnexecution.v1.BeamFnApi; import org.apache.beam.runners.core.DoFnRunner; import org.apache.beam.runners.core.DoFnRunners; import org.apache.beam.runners.core.InMemoryStateInternals; @@ -176,7 +177,12 @@ private void addElement( valueCoder.encode(value, dataOutput, Context.OUTER); messageBundle .addMessagesBuilder() - .setMetadata(WindmillSink.encodeMetadata(windowsCoder, windows, PaneInfo.NO_FIRING)) + .setMetadata( + WindmillSink.encodeMetadata( + windowsCoder, + windows, + PaneInfo.NO_FIRING, + BeamFnApi.Elements.ElementMetadata.newBuilder().build())) .setData(dataOutput.toByteString()) .setTimestamp(WindmillTimeUtils.harnessToWindmillTimestamp(timestamp)); } diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingGroupAlsoByWindowsReshuffleDoFnTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingGroupAlsoByWindowsReshuffleDoFnTest.java index c169c9b46a57..a348c0f00214 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingGroupAlsoByWindowsReshuffleDoFnTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingGroupAlsoByWindowsReshuffleDoFnTest.java @@ -24,6 +24,7 @@ import java.util.Arrays; import java.util.Collection; import java.util.List; +import org.apache.beam.model.fnexecution.v1.BeamFnApi; import org.apache.beam.runners.core.DoFnRunner; import org.apache.beam.runners.core.KeyedWorkItem; import org.apache.beam.runners.core.NullSideInputReader; @@ -114,7 +115,12 @@ private void addElement( valueCoder.encode(value, dataOutput, Context.OUTER); messageBundle .addMessagesBuilder() - .setMetadata(WindmillSink.encodeMetadata(windowsCoder, windows, PaneInfo.NO_FIRING)) + .setMetadata( + WindmillSink.encodeMetadata( + windowsCoder, + windows, + PaneInfo.NO_FIRING, + BeamFnApi.Elements.ElementMetadata.newBuilder().build())) .setData(dataOutput.toByteString()) .setTimestamp(WindmillTimeUtils.harnessToWindmillTimestamp(timestamp)); } diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItemTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItemTest.java index ffe71176367a..53a36722e41c 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItemTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItemTest.java @@ -22,6 +22,7 @@ import java.io.IOException; import java.util.Collection; import java.util.Collections; +import org.apache.beam.model.fnexecution.v1.BeamFnApi; import org.apache.beam.runners.core.KeyedWorkItem; import org.apache.beam.runners.core.StateNamespace; import org.apache.beam.runners.core.StateNamespaces; @@ -107,10 +108,14 @@ private void addElement( long timestamp, String value, IntervalWindow window, - PaneInfo paneInfo) + PaneInfo pane) throws IOException { ByteString encodedMetadata = - WindmillSink.encodeMetadata(WINDOWS_CODER, Collections.singletonList(window), paneInfo); + WindmillSink.encodeMetadata( + WINDOWS_CODER, + Collections.singletonList(window), + pane, + BeamFnApi.Elements.ElementMetadata.newBuilder().build()); chunk .addMessagesBuilder() .setTimestamp(WindmillTimeUtils.harnessToWindmillTimestamp(new Instant(timestamp))) diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/windowing/PaneInfo.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/windowing/PaneInfo.java index 6e4c694d48e3..80f931b8e6a6 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/windowing/PaneInfo.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/windowing/PaneInfo.java @@ -146,10 +146,10 @@ private static byte encodedByte(boolean isFirst, boolean isLast, Timing timing) ImmutableMap.Builder decodingBuilder = ImmutableMap.builder(); for (Timing timing : Timing.values()) { long onTimeIndex = timing == Timing.EARLY ? -1 : 0; - register(decodingBuilder, new PaneInfo(true, true, timing, 0, onTimeIndex)); - register(decodingBuilder, new PaneInfo(true, false, timing, 0, onTimeIndex)); - register(decodingBuilder, new PaneInfo(false, true, timing, -1, onTimeIndex)); - register(decodingBuilder, new PaneInfo(false, false, timing, -1, onTimeIndex)); + register(decodingBuilder, new PaneInfo(true, true, timing, 0, onTimeIndex, false)); + register(decodingBuilder, new PaneInfo(true, false, timing, 0, onTimeIndex, false)); + register(decodingBuilder, new PaneInfo(false, true, timing, -1, onTimeIndex, false)); + register(decodingBuilder, new PaneInfo(false, false, timing, -1, onTimeIndex, false)); } BYTE_TO_PANE_INFO = decodingBuilder.build(); } @@ -159,7 +159,7 @@ private static void register(ImmutableMap.Builder builder, PaneI } private final byte encodedByte; - + private final boolean containsElementMetadata; private final boolean isFirst; private final boolean isLast; private final Timing timing; @@ -177,13 +177,20 @@ private static void register(ImmutableMap.Builder builder, PaneI public static final PaneInfo ON_TIME_AND_ONLY_FIRING = PaneInfo.createPane(true, true, Timing.ON_TIME, 0, 0); - private PaneInfo(boolean isFirst, boolean isLast, Timing timing, long index, long onTimeIndex) { + private PaneInfo( + boolean isFirst, + boolean isLast, + Timing timing, + long index, + long onTimeIndex, + boolean containsElementMetadata) { this.encodedByte = encodedByte(isFirst, isLast, timing); this.isFirst = isFirst; this.isLast = isLast; this.timing = timing; this.index = index; this.nonSpeculativeIndex = onTimeIndex; + this.containsElementMetadata = containsElementMetadata; } public static PaneInfo createPane(boolean isFirst, boolean isLast, Timing timing) { @@ -191,13 +198,25 @@ public static PaneInfo createPane(boolean isFirst, boolean isLast, Timing timing return createPane(isFirst, isLast, timing, 0, timing == Timing.EARLY ? -1 : 0); } + /** Factory method to create a {@link PaneInfo} with the specified parameters. */ /** Factory method to create a {@link PaneInfo} with the specified parameters. */ public static PaneInfo createPane( boolean isFirst, boolean isLast, Timing timing, long index, long onTimeIndex) { + return createPane(isFirst, isLast, timing, index, onTimeIndex, false); + } + + /** Factory method to create a {@link PaneInfo} with the specified parameters. */ + public static PaneInfo createPane( + boolean isFirst, + boolean isLast, + Timing timing, + long index, + long onTimeIndex, + boolean containsElementMetadata) { if (isFirst || timing == Timing.UNKNOWN) { return checkNotNull(BYTE_TO_PANE_INFO.get(encodedByte(isFirst, isLast, timing))); } else { - return new PaneInfo(isFirst, isLast, timing, index, onTimeIndex); + return new PaneInfo(isFirst, isLast, timing, index, onTimeIndex, containsElementMetadata); } } @@ -219,6 +238,15 @@ public boolean isFirst() { return isFirst; } + public boolean isElementMetadata() { + return containsElementMetadata; + } + + public PaneInfo withElementMetadata(boolean elementMetadata) { + return new PaneInfo( + this.isFirst, this.isLast, this.timing, index, nonSpeculativeIndex, elementMetadata); + } + /** Return true if this is the last pane that will be produced in the associated window. */ public boolean isLast() { return isLast; @@ -295,6 +323,8 @@ public String toString() { /** A Coder for encoding PaneInfo instances. */ public static class PaneInfoCoder extends AtomicCoder { + private static final byte ELEMENT_METADATA_MASK = (byte) 0x80; + private enum Encoding { FIRST, ONE_INDEX, @@ -337,16 +367,17 @@ private PaneInfoCoder() {} public void encode(PaneInfo value, final OutputStream outStream) throws CoderException, IOException { Encoding encoding = chooseEncoding(value); + byte elementMetadata = value.containsElementMetadata ? ELEMENT_METADATA_MASK : 0x00; switch (chooseEncoding(value)) { case FIRST: - outStream.write(value.encodedByte); + outStream.write(value.encodedByte | elementMetadata); break; case ONE_INDEX: - outStream.write(value.encodedByte | encoding.tag); + outStream.write(value.encodedByte | encoding.tag | elementMetadata); VarInt.encode(value.index, outStream); break; case TWO_INDICES: - outStream.write(value.encodedByte | encoding.tag); + outStream.write(value.encodedByte | encoding.tag | elementMetadata); VarInt.encode(value.index, outStream); VarInt.encode(value.nonSpeculativeIndex, outStream); break; @@ -360,9 +391,10 @@ public PaneInfo decode(final InputStream inStream) throws CoderException, IOExce byte keyAndTag = (byte) inStream.read(); PaneInfo base = Preconditions.checkNotNull(BYTE_TO_PANE_INFO.get((byte) (keyAndTag & 0x0F))); long index, onTimeIndex; - switch (Encoding.fromTag(keyAndTag)) { + boolean elementMetadata = (keyAndTag & ELEMENT_METADATA_MASK) != 0; + switch (Encoding.fromTag((byte) (keyAndTag & ~ELEMENT_METADATA_MASK))) { case FIRST: - return base; + return base.withElementMetadata(elementMetadata); case ONE_INDEX: index = VarInt.decodeLong(inStream); onTimeIndex = base.timing == Timing.EARLY ? -1 : index; @@ -374,7 +406,8 @@ public PaneInfo decode(final InputStream inStream) throws CoderException, IOExce default: throw new CoderException("Unknown encoding " + (keyAndTag & 0xF0)); } - return new PaneInfo(base.isFirst, base.isLast, base.timing, index, onTimeIndex); + return new PaneInfo( + base.isFirst, base.isLast, base.timing, index, onTimeIndex, elementMetadata); } @Override diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/values/WindowedValues.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/values/WindowedValues.java index 9b079b8699b9..26ba4948bc0a 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/values/WindowedValues.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/values/WindowedValues.java @@ -36,6 +36,7 @@ import java.util.List; import java.util.Objects; import java.util.Set; +import org.apache.beam.model.fnexecution.v1.BeamFnApi; import org.apache.beam.sdk.annotations.Internal; import org.apache.beam.sdk.coders.ByteArrayCoder; import org.apache.beam.sdk.coders.Coder; @@ -763,6 +764,15 @@ public static ParamWindowedValueCoder getParamWindowedValueCoder(Coder /** Abstract class for {@code WindowedValue} coder. */ public abstract static class WindowedValueCoder extends StructuredCoder> { final Coder valueCoder; + private static boolean metadataSupported = false; + + public static void setMetadataSupported() { + metadataSupported = true; + } + + public static boolean isMetadataSupported() { + return metadataSupported; + } WindowedValueCoder(Coder valueCoder) { this.valueCoder = checkNotNull(valueCoder); @@ -829,7 +839,15 @@ public void encode(WindowedValue windowedElem, OutputStream outStream, Contex throws CoderException, IOException { InstantCoder.of().encode(windowedElem.getTimestamp(), outStream); windowsCoder.encode(windowedElem.getWindows(), outStream); - PaneInfoCoder.INSTANCE.encode(windowedElem.getPaneInfo(), outStream); + boolean metadataSupported = isMetadataSupported(); + PaneInfoCoder.INSTANCE.encode( + windowedElem.getPaneInfo().withElementMetadata(metadataSupported), outStream); + if (metadataSupported) { + BeamFnApi.Elements.ElementMetadata.Builder builder = + BeamFnApi.Elements.ElementMetadata.newBuilder(); + BeamFnApi.Elements.ElementMetadata em = builder.build(); + em.writeDelimitedTo(outStream); + } valueCoder.encode(windowedElem.getValue(), outStream, context); } @@ -844,6 +862,9 @@ public WindowedValue decode(InputStream inStream, Context context) Instant timestamp = InstantCoder.of().decode(inStream); Collection windows = windowsCoder.decode(inStream); PaneInfo paneInfo = PaneInfoCoder.INSTANCE.decode(inStream); + if (isMetadataSupported() && paneInfo.isElementMetadata()) { + BeamFnApi.Elements.ElementMetadata.parseDelimitedFrom(inStream); + } T value = valueCoder.decode(inStream, context); // Because there are some remaining (incorrect) uses of WindowedValue with no windows, diff --git a/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/windowing/PaneInfoTest.java b/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/windowing/PaneInfoTest.java index 946deba036db..cda8ee1ea55c 100644 --- a/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/windowing/PaneInfoTest.java +++ b/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/windowing/PaneInfoTest.java @@ -52,6 +52,22 @@ public void testEncodingRoundTrip() throws Exception { } } + @Test + public void testEncodingRoundTripWithElementMetadata() throws Exception { + Coder coder = PaneInfo.PaneInfoCoder.INSTANCE; + for (Timing timing : Timing.values()) { + long onTimeIndex = timing == Timing.EARLY ? -1 : 37; + CoderProperties.coderDecodeEncodeEqual( + coder, PaneInfo.createPane(false, false, timing, 389, onTimeIndex, true)); + CoderProperties.coderDecodeEncodeEqual( + coder, PaneInfo.createPane(false, true, timing, 5077, onTimeIndex, true)); + CoderProperties.coderDecodeEncodeEqual( + coder, PaneInfo.createPane(true, false, timing, 0, 0, true)); + CoderProperties.coderDecodeEncodeEqual( + coder, PaneInfo.createPane(true, true, timing, 0, 0, true)); + } + } + @Test public void testEncodings() { assertEquals( @@ -82,5 +98,9 @@ public void testEncodings() { "PaneInfo encoding should remain the same.", 0xF, PaneInfo.createPane(true, true, Timing.UNKNOWN).getEncodedByte()); + assertEquals( + "PaneInfo encoding should remain the same.", + 0x1, + PaneInfo.createPane(true, false, Timing.EARLY, 1, -1, true).getEncodedByte()); } } diff --git a/sdks/java/core/src/test/java/org/apache/beam/sdk/util/WindowedValueTest.java b/sdks/java/core/src/test/java/org/apache/beam/sdk/util/WindowedValueTest.java index 18660c5e6c36..957a839eae94 100644 --- a/sdks/java/core/src/test/java/org/apache/beam/sdk/util/WindowedValueTest.java +++ b/sdks/java/core/src/test/java/org/apache/beam/sdk/util/WindowedValueTest.java @@ -77,6 +77,33 @@ public void testWindowedValueCoder() throws CoderException { Assert.assertArrayEquals(value.getWindows().toArray(), decodedValue.getWindows().toArray()); } + @Test + public void testWindowedValueWithElementMetadataCoder() throws CoderException { + WindowedValues.WindowedValueCoder.setMetadataSupported(); + PaneInfo.PaneInfoCoder.setMetadataSupported(); + Instant timestamp = new Instant(1234); + WindowedValue value = + WindowedValues.of( + "abc", + new Instant(1234), + Arrays.asList( + new IntervalWindow(timestamp, timestamp.plus(Duration.millis(1000))), + new IntervalWindow( + timestamp.plus(Duration.millis(1000)), timestamp.plus(Duration.millis(2000)))), + PaneInfo.NO_FIRING); + + Coder> windowedValueCoder = + WindowedValues.getFullCoder(StringUtf8Coder.of(), IntervalWindow.getCoder()); + + byte[] encodedValue = CoderUtils.encodeToByteArray(windowedValueCoder, value); + WindowedValue decodedValue = + CoderUtils.decodeFromByteArray(windowedValueCoder, encodedValue); + + Assert.assertEquals(value.getValue(), decodedValue.getValue()); + Assert.assertEquals(value.getTimestamp(), decodedValue.getTimestamp()); + Assert.assertArrayEquals(value.getWindows().toArray(), decodedValue.getWindows().toArray()); + } + @Test public void testFullWindowedValueCoderIsSerializableWithWellKnownCoderType() { CoderProperties.coderSerializable( From 7cfd6f8c0539c2ca7fac4f87498686b67e949964 Mon Sep 17 00:00:00 2001 From: Radek Stankiewicz Date: Fri, 3 Oct 2025 11:39:38 +0200 Subject: [PATCH 2/9] add todos for future otel metadata --- .../worker/UngroupedWindmillReader.java | 6 ++++- .../worker/WindmillKeyedWorkItem.java | 5 +++- .../runners/dataflow/worker/WindmillSink.java | 25 ++++++++++++++----- .../sdk/transforms/windowing/PaneInfo.java | 1 - 4 files changed, 28 insertions(+), 9 deletions(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/UngroupedWindmillReader.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/UngroupedWindmillReader.java index 5df15e8d2109..d2e3e84ee085 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/UngroupedWindmillReader.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/UngroupedWindmillReader.java @@ -117,7 +117,9 @@ protected WindowedValue decodeMessage(Windmill.Message message) throws IOExce Collection windows = WindmillSink.decodeMetadataWindows(windowsCoder, message.getMetadata()); PaneInfo paneInfo = WindmillSink.decodeMetadataPane(message.getMetadata()); - WindmillSink.decodeAdditionalMetadata(windowsCoder, message.getMetadata()); + if(WindowedValues.WindowedValueCoder.isMetadataSupported()) { + WindmillSink.decodeAdditionalMetadata(windowsCoder, message.getMetadata()); + } if (valueCoder instanceof KvCoder) { KvCoder kvCoder = (KvCoder) valueCoder; InputStream key = context.getSerializedKey().newInput(); @@ -126,9 +128,11 @@ protected WindowedValue decodeMessage(Windmill.Message message) throws IOExce @SuppressWarnings("unchecked") T result = (T) KV.of(decode(kvCoder.getKeyCoder(), key), decode(kvCoder.getValueCoder(), data)); + // todo #33176 propagate metadata to windowed value return WindowedValues.of(result, timestampMillis, windows, paneInfo); } else { notifyElementRead(data.available() + metadata.available()); + // todo #33176 propagate metadata to windowed value return WindowedValues.of(decode(valueCoder, data), timestampMillis, windows, paneInfo); } } diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItem.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItem.java index bae41a156967..1b5a52ad71cf 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItem.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItem.java @@ -108,9 +108,12 @@ public Iterable> elementsIterable() { Collection windows = WindmillSink.decodeMetadataWindows(windowsCoder, message.getMetadata()); PaneInfo paneInfo = WindmillSink.decodeMetadataPane(message.getMetadata()); - WindmillSink.decodeAdditionalMetadata(windowsCoder, message.getMetadata()); + if(WindowedValues.WindowedValueCoder.isMetadataSupported()) { + WindmillSink.decodeAdditionalMetadata(windowsCoder, message.getMetadata()); + } InputStream inputStream = message.getData().newInput(); ElemT value = valueCoder.decode(inputStream, Coder.Context.OUTER); + // todo #33176 specify additional metadata in the future return WindowedValues.of(value, timestamp, windows, paneInfo); } catch (IOException e) { throw new RuntimeException(e); diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java index 4765eb74e7f5..b9f3e235379b 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java @@ -30,6 +30,7 @@ import org.apache.beam.runners.dataflow.util.CloudObject; import org.apache.beam.runners.dataflow.worker.util.common.worker.Sink; import org.apache.beam.runners.dataflow.worker.windmill.Windmill; +import org.apache.beam.sdk.coders.ByteArrayCoder; import org.apache.beam.sdk.coders.Coder; import org.apache.beam.sdk.coders.KvCoder; import org.apache.beam.sdk.options.PipelineOptions; @@ -81,8 +82,16 @@ private static ByteString encodeMetadata( BeamFnApi.Elements.ElementMetadata metadata) throws IOException { try { - PaneInfoCoder.INSTANCE.encode(paneInfo, stream); - windowsCoder.encode(windows, stream, Coder.Context.OUTER); + // element metadata is behind the experiment + boolean elementMetadata = WindowedValues.WindowedValueCoder.isMetadataSupported(); + if (elementMetadata) { + PaneInfoCoder.INSTANCE.encode(paneInfo, stream); + windowsCoder.encode(windows, stream); + ByteArrayCoder.of().encode(metadata.toByteArray(), stream, Coder.Context.OUTER); + }else { + PaneInfoCoder.INSTANCE.encode(paneInfo, stream); + windowsCoder.encode(windows, stream, Coder.Context.OUTER); + } return stream.toByteStringAndReset(); } catch (Exception e) { stream.reset(); @@ -93,10 +102,11 @@ private static ByteString encodeMetadata( public static ByteString encodeMetadata( Coder> windowsCoder, Collection windows, - PaneInfo paneInfo) + PaneInfo paneInfo, + BeamFnApi.Elements.ElementMetadata metadata) throws IOException { ByteStringOutputStream stream = new ByteStringOutputStream(); - return encodeMetadata(stream, windowsCoder, windows, paneInfo); + return encodeMetadata(stream, windowsCoder, windows, paneInfo, metadata); } public static PaneInfo decodeMetadataPane(ByteString metadata) throws IOException { @@ -110,7 +120,7 @@ public static BeamFnApi.Elements.ElementMetadata decodeAdditionalMetadata( InputStream inStream = metadata.newInput(); PaneInfo paneInfo = PaneInfoCoder.INSTANCE.decode(inStream); windowsCoder.decode(inStream); - if (paneInfo.isElementMetadata() && WindowedValues.WindowedValueCoder.isMetadataSupported()) { + if (paneInfo.isElementMetadata()) { return BeamFnApi.Elements.ElementMetadata.parseDelimitedFrom(inStream); } else { // empty @@ -201,8 +211,11 @@ private ByteString encode(Coder coder, EncodeT object) throws public long add(WindowedValue data) throws IOException { ByteString key, value; ByteString id = ByteString.EMPTY; + // todo - add here all windowedValue metadata + BeamFnApi.Elements.ElementMetadata additionalMetadata = + BeamFnApi.Elements.ElementMetadata.newBuilder().build(); ByteString metadata = - encodeMetadata(stream, windowsCoder, data.getWindows(), data.getPaneInfo()); + encodeMetadata(stream, windowsCoder, data.getWindows(), data.getPaneInfo(), additionalMetadata); if (valueCoder instanceof KvCoder) { KvCoder kvCoder = (KvCoder) valueCoder; KV kv = (KV) data.getValue(); diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/windowing/PaneInfo.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/windowing/PaneInfo.java index 80f931b8e6a6..bc83687bae4e 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/windowing/PaneInfo.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/windowing/PaneInfo.java @@ -198,7 +198,6 @@ public static PaneInfo createPane(boolean isFirst, boolean isLast, Timing timing return createPane(isFirst, isLast, timing, 0, timing == Timing.EARLY ? -1 : 0); } - /** Factory method to create a {@link PaneInfo} with the specified parameters. */ /** Factory method to create a {@link PaneInfo} with the specified parameters. */ public static PaneInfo createPane( boolean isFirst, boolean isLast, Timing timing, long index, long onTimeIndex) { From 5b7e0fc21da9904fb3c52888a7c6d47d78bb1a49 Mon Sep 17 00:00:00 2001 From: Radek Stankiewicz Date: Fri, 3 Oct 2025 14:03:23 +0200 Subject: [PATCH 3/9] add todos for future otel metadata --- .../RedistributeByKeyOverrideFactory.java | 1 + .../runners/dataflow/worker/WindmillSink.java | 4 ++-- .../beam/sdk/transforms/Redistribute.java | 1 + .../org/apache/beam/sdk/transforms/Reify.java | 1 + .../apache/beam/sdk/transforms/Reshuffle.java | 1 + .../beam/sdk/values/ValueInSingleWindow.java | 19 ++++++++++++++++++- .../beam/sdk/values/WindowedValues.java | 5 +++-- 7 files changed, 27 insertions(+), 5 deletions(-) diff --git a/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/RedistributeByKeyOverrideFactory.java b/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/RedistributeByKeyOverrideFactory.java index 4375cc5adcfe..47ff5b764910 100644 --- a/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/RedistributeByKeyOverrideFactory.java +++ b/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/RedistributeByKeyOverrideFactory.java @@ -135,6 +135,7 @@ public Duration getAllowedTimestampSkew() { public void processElement( @Element KV> kv, OutputReceiver> outputReceiver) { + // todo #33176 specify additional metadata in the future outputReceiver .builder(KV.of(kv.getKey(), kv.getValue().getValue())) .setTimestamp(kv.getValue().getTimestamp()) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java index b9f3e235379b..3fe6af954cb5 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java @@ -121,7 +121,7 @@ public static BeamFnApi.Elements.ElementMetadata decodeAdditionalMetadata( PaneInfo paneInfo = PaneInfoCoder.INSTANCE.decode(inStream); windowsCoder.decode(inStream); if (paneInfo.isElementMetadata()) { - return BeamFnApi.Elements.ElementMetadata.parseDelimitedFrom(inStream); + return BeamFnApi.Elements.ElementMetadata.parseFrom(ByteArrayCoder.of().decode(inStream)); } else { // empty return BeamFnApi.Elements.ElementMetadata.newBuilder().build(); @@ -211,7 +211,7 @@ private ByteString encode(Coder coder, EncodeT object) throws public long add(WindowedValue data) throws IOException { ByteString key, value; ByteString id = ByteString.EMPTY; - // todo - add here all windowedValue metadata + // todo #33176 specify additional metadata in the future BeamFnApi.Elements.ElementMetadata additionalMetadata = BeamFnApi.Elements.ElementMetadata.newBuilder().build(); ByteString metadata = diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Redistribute.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Redistribute.java index a01b5f570a57..7382f29310c5 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Redistribute.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Redistribute.java @@ -179,6 +179,7 @@ public Duration getAllowedTimestampSkew() { public void processElement( @Element KV> kv, OutputReceiver> outputReceiver) { + // todo #33176 specify additional metadata in the future outputReceiver .builder(KV.of(kv.getKey(), kv.getValue().getValue())) .setTimestamp(kv.getValue().getTimestamp()) diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Reify.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Reify.java index 797af9538c53..905e94ccc4d0 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Reify.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Reify.java @@ -136,6 +136,7 @@ public PCollection>> expand(PCollection> i KvCoder coder = (KvCoder) input.getCoder(); return input .apply( + // todo #33176 specify additional metadata in the future ParDo.of( new DoFn, KV>>() { @ProcessElement diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Reshuffle.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Reshuffle.java index b2de48342d7c..0a8d058107b8 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Reshuffle.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Reshuffle.java @@ -184,6 +184,7 @@ public Duration getAllowedTimestampSkew() { public void processElement( @Element KV> kv, OutputReceiver> outputReceiver) { + // todo #33176 specify additional metadata in the future outputReceiver .builder(KV.of(kv.getKey(), kv.getValue().getValue())) .setTimestamp(kv.getValue().getTimestamp()) diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/values/ValueInSingleWindow.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/values/ValueInSingleWindow.java index 7dc5fef52ecb..56cb58ed6d15 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/values/ValueInSingleWindow.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/values/ValueInSingleWindow.java @@ -22,7 +22,9 @@ import java.io.InputStream; import java.io.OutputStream; import java.util.List; +import org.apache.beam.model.fnexecution.v1.BeamFnApi; import org.apache.beam.sdk.annotations.Internal; +import org.apache.beam.sdk.coders.ByteArrayCoder; import org.apache.beam.sdk.coders.InstantCoder; import org.apache.beam.sdk.coders.StructuredCoder; import org.apache.beam.sdk.transforms.windowing.BoundedWindow; @@ -64,6 +66,7 @@ public T getValue() { public abstract @Nullable Long getCurrentRecordOffset(); + // todo #33176 specify additional metadata in the future public static ValueInSingleWindow of( T value, Instant timestamp, @@ -110,7 +113,16 @@ public void encode(ValueInSingleWindow windowedElem, OutputStream outStream, throws IOException { InstantCoder.of().encode(windowedElem.getTimestamp(), outStream); windowCoder.encode(windowedElem.getWindow(), outStream); - PaneInfo.PaneInfoCoder.INSTANCE.encode(windowedElem.getPaneInfo(), outStream); + boolean metadataSupported = WindowedValues.WindowedValueCoder.isMetadataSupported(); + PaneInfo.PaneInfoCoder.INSTANCE.encode(windowedElem.getPaneInfo().withElementMetadata(metadataSupported), outStream); + if(metadataSupported){ + BeamFnApi.Elements.ElementMetadata.Builder builder = + BeamFnApi.Elements.ElementMetadata.newBuilder(); + // todo #33176 specify additional metadata in the future + BeamFnApi.Elements.ElementMetadata metadata = builder.build(); + ByteArrayCoder.of().encode(metadata.toByteArray(), outStream); + } + valueCoder.encode(windowedElem.getValue(), outStream, context); } @@ -124,7 +136,12 @@ public ValueInSingleWindow decode(InputStream inStream, Context context) thro Instant timestamp = InstantCoder.of().decode(inStream); BoundedWindow window = windowCoder.decode(inStream); PaneInfo paneInfo = PaneInfo.PaneInfoCoder.INSTANCE.decode(inStream); + if(WindowedValues.WindowedValueCoder.isMetadataSupported() && paneInfo.isElementMetadata()) { + BeamFnApi.Elements.ElementMetadata.parseFrom(ByteArrayCoder.of().decode(inStream)); + } + T value = valueCoder.decode(inStream, context); + // todo #33176 specify additional metadata in the future return new AutoValue_ValueInSingleWindow<>(value, timestamp, window, paneInfo, null, null); } diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/values/WindowedValues.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/values/WindowedValues.java index 26ba4948bc0a..91763002d294 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/values/WindowedValues.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/values/WindowedValues.java @@ -846,7 +846,8 @@ public void encode(WindowedValue windowedElem, OutputStream outStream, Contex BeamFnApi.Elements.ElementMetadata.Builder builder = BeamFnApi.Elements.ElementMetadata.newBuilder(); BeamFnApi.Elements.ElementMetadata em = builder.build(); - em.writeDelimitedTo(outStream); + ByteArrayCoder.of().encode(em.toByteArray(), outStream); + } valueCoder.encode(windowedElem.getValue(), outStream, context); } @@ -863,7 +864,7 @@ public WindowedValue decode(InputStream inStream, Context context) Collection windows = windowsCoder.decode(inStream); PaneInfo paneInfo = PaneInfoCoder.INSTANCE.decode(inStream); if (isMetadataSupported() && paneInfo.isElementMetadata()) { - BeamFnApi.Elements.ElementMetadata.parseDelimitedFrom(inStream); + BeamFnApi.Elements.ElementMetadata.parseFrom(ByteArrayCoder.of().decode(inStream)); } T value = valueCoder.decode(inStream, context); From 7654478a5d2495382cdaa4e6ab57f65e36492e8a Mon Sep 17 00:00:00 2001 From: Radek Stankiewicz Date: Fri, 3 Oct 2025 14:27:29 +0200 Subject: [PATCH 4/9] spotless --- .../runners/dataflow/worker/UngroupedWindmillReader.java | 2 +- .../runners/dataflow/worker/WindmillKeyedWorkItem.java | 2 +- .../beam/runners/dataflow/worker/WindmillSink.java | 7 ++++--- .../main/java/org/apache/beam/sdk/transforms/Reify.java | 2 +- .../org/apache/beam/sdk/values/ValueInSingleWindow.java | 9 +++++---- .../java/org/apache/beam/sdk/values/WindowedValues.java | 1 - 6 files changed, 12 insertions(+), 11 deletions(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/UngroupedWindmillReader.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/UngroupedWindmillReader.java index d2e3e84ee085..a9a033c89ad7 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/UngroupedWindmillReader.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/UngroupedWindmillReader.java @@ -117,7 +117,7 @@ protected WindowedValue decodeMessage(Windmill.Message message) throws IOExce Collection windows = WindmillSink.decodeMetadataWindows(windowsCoder, message.getMetadata()); PaneInfo paneInfo = WindmillSink.decodeMetadataPane(message.getMetadata()); - if(WindowedValues.WindowedValueCoder.isMetadataSupported()) { + if (WindowedValues.WindowedValueCoder.isMetadataSupported()) { WindmillSink.decodeAdditionalMetadata(windowsCoder, message.getMetadata()); } if (valueCoder instanceof KvCoder) { diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItem.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItem.java index 1b5a52ad71cf..6690377d3de6 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItem.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItem.java @@ -108,7 +108,7 @@ public Iterable> elementsIterable() { Collection windows = WindmillSink.decodeMetadataWindows(windowsCoder, message.getMetadata()); PaneInfo paneInfo = WindmillSink.decodeMetadataPane(message.getMetadata()); - if(WindowedValues.WindowedValueCoder.isMetadataSupported()) { + if (WindowedValues.WindowedValueCoder.isMetadataSupported()) { WindmillSink.decodeAdditionalMetadata(windowsCoder, message.getMetadata()); } InputStream inputStream = message.getData().newInput(); diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java index 3fe6af954cb5..4c5841f2de97 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java @@ -88,7 +88,7 @@ private static ByteString encodeMetadata( PaneInfoCoder.INSTANCE.encode(paneInfo, stream); windowsCoder.encode(windows, stream); ByteArrayCoder.of().encode(metadata.toByteArray(), stream, Coder.Context.OUTER); - }else { + } else { PaneInfoCoder.INSTANCE.encode(paneInfo, stream); windowsCoder.encode(windows, stream, Coder.Context.OUTER); } @@ -213,9 +213,10 @@ public long add(WindowedValue data) throws IOException { ByteString id = ByteString.EMPTY; // todo #33176 specify additional metadata in the future BeamFnApi.Elements.ElementMetadata additionalMetadata = - BeamFnApi.Elements.ElementMetadata.newBuilder().build(); + BeamFnApi.Elements.ElementMetadata.newBuilder().build(); ByteString metadata = - encodeMetadata(stream, windowsCoder, data.getWindows(), data.getPaneInfo(), additionalMetadata); + encodeMetadata( + stream, windowsCoder, data.getWindows(), data.getPaneInfo(), additionalMetadata); if (valueCoder instanceof KvCoder) { KvCoder kvCoder = (KvCoder) valueCoder; KV kv = (KV) data.getValue(); diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Reify.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Reify.java index 905e94ccc4d0..af125d9e63e8 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Reify.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Reify.java @@ -136,7 +136,7 @@ public PCollection>> expand(PCollection> i KvCoder coder = (KvCoder) input.getCoder(); return input .apply( - // todo #33176 specify additional metadata in the future + // todo #33176 specify additional metadata in the future ParDo.of( new DoFn, KV>>() { @ProcessElement diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/values/ValueInSingleWindow.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/values/ValueInSingleWindow.java index 56cb58ed6d15..2bbc28cac9fe 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/values/ValueInSingleWindow.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/values/ValueInSingleWindow.java @@ -114,10 +114,11 @@ public void encode(ValueInSingleWindow windowedElem, OutputStream outStream, InstantCoder.of().encode(windowedElem.getTimestamp(), outStream); windowCoder.encode(windowedElem.getWindow(), outStream); boolean metadataSupported = WindowedValues.WindowedValueCoder.isMetadataSupported(); - PaneInfo.PaneInfoCoder.INSTANCE.encode(windowedElem.getPaneInfo().withElementMetadata(metadataSupported), outStream); - if(metadataSupported){ + PaneInfo.PaneInfoCoder.INSTANCE.encode( + windowedElem.getPaneInfo().withElementMetadata(metadataSupported), outStream); + if (metadataSupported) { BeamFnApi.Elements.ElementMetadata.Builder builder = - BeamFnApi.Elements.ElementMetadata.newBuilder(); + BeamFnApi.Elements.ElementMetadata.newBuilder(); // todo #33176 specify additional metadata in the future BeamFnApi.Elements.ElementMetadata metadata = builder.build(); ByteArrayCoder.of().encode(metadata.toByteArray(), outStream); @@ -136,7 +137,7 @@ public ValueInSingleWindow decode(InputStream inStream, Context context) thro Instant timestamp = InstantCoder.of().decode(inStream); BoundedWindow window = windowCoder.decode(inStream); PaneInfo paneInfo = PaneInfo.PaneInfoCoder.INSTANCE.decode(inStream); - if(WindowedValues.WindowedValueCoder.isMetadataSupported() && paneInfo.isElementMetadata()) { + if (WindowedValues.WindowedValueCoder.isMetadataSupported() && paneInfo.isElementMetadata()) { BeamFnApi.Elements.ElementMetadata.parseFrom(ByteArrayCoder.of().decode(inStream)); } diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/values/WindowedValues.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/values/WindowedValues.java index 91763002d294..d90637da9b1b 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/values/WindowedValues.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/values/WindowedValues.java @@ -847,7 +847,6 @@ public void encode(WindowedValue windowedElem, OutputStream outStream, Contex BeamFnApi.Elements.ElementMetadata.newBuilder(); BeamFnApi.Elements.ElementMetadata em = builder.build(); ByteArrayCoder.of().encode(em.toByteArray(), outStream); - } valueCoder.encode(windowedElem.getValue(), outStream, context); } From 0274977bacf313b561559d1c31241bfdc608bc16 Mon Sep 17 00:00:00 2001 From: Radek Stankiewicz Date: Fri, 3 Oct 2025 14:56:47 +0200 Subject: [PATCH 5/9] fix tests --- .../test/java/org/apache/beam/sdk/util/WindowedValueTest.java | 1 - 1 file changed, 1 deletion(-) diff --git a/sdks/java/core/src/test/java/org/apache/beam/sdk/util/WindowedValueTest.java b/sdks/java/core/src/test/java/org/apache/beam/sdk/util/WindowedValueTest.java index 957a839eae94..3e3973e3720b 100644 --- a/sdks/java/core/src/test/java/org/apache/beam/sdk/util/WindowedValueTest.java +++ b/sdks/java/core/src/test/java/org/apache/beam/sdk/util/WindowedValueTest.java @@ -80,7 +80,6 @@ public void testWindowedValueCoder() throws CoderException { @Test public void testWindowedValueWithElementMetadataCoder() throws CoderException { WindowedValues.WindowedValueCoder.setMetadataSupported(); - PaneInfo.PaneInfoCoder.setMetadataSupported(); Instant timestamp = new Instant(1234); WindowedValue value = WindowedValues.of( From bbb947471153f17d14417d22639aae1ad1fb27e3 Mon Sep 17 00:00:00 2001 From: Radek Stankiewicz Date: Fri, 3 Oct 2025 19:33:18 +0200 Subject: [PATCH 6/9] fix --- .../org/apache/beam/runners/dataflow/worker/WindmillSink.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java index 4c5841f2de97..18c023d62aa8 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java @@ -85,7 +85,7 @@ private static ByteString encodeMetadata( // element metadata is behind the experiment boolean elementMetadata = WindowedValues.WindowedValueCoder.isMetadataSupported(); if (elementMetadata) { - PaneInfoCoder.INSTANCE.encode(paneInfo, stream); + PaneInfoCoder.INSTANCE.encode(paneInfo.withElementMetadata(elementMetadata), stream); windowsCoder.encode(windows, stream); ByteArrayCoder.of().encode(metadata.toByteArray(), stream, Coder.Context.OUTER); } else { From 99edc7650e653ccf899be99f49e532f8479b5a23 Mon Sep 17 00:00:00 2001 From: Radek Stankiewicz Date: Fri, 3 Oct 2025 19:33:41 +0200 Subject: [PATCH 7/9] fix --- .../org/apache/beam/runners/dataflow/worker/WindmillSink.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java index 18c023d62aa8..6d3484f6642d 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java @@ -85,7 +85,7 @@ private static ByteString encodeMetadata( // element metadata is behind the experiment boolean elementMetadata = WindowedValues.WindowedValueCoder.isMetadataSupported(); if (elementMetadata) { - PaneInfoCoder.INSTANCE.encode(paneInfo.withElementMetadata(elementMetadata), stream); + PaneInfoCoder.INSTANCE.encode(paneInfo.withElementMetadata(true), stream); windowsCoder.encode(windows, stream); ByteArrayCoder.of().encode(metadata.toByteArray(), stream, Coder.Context.OUTER); } else { From 1e38bec43d57cb13675eb6ef7a28b3dc00652008 Mon Sep 17 00:00:00 2001 From: Radek Stankiewicz Date: Thu, 9 Oct 2025 13:50:14 +0200 Subject: [PATCH 8/9] fix bugs --- .../org/apache/beam/runners/dataflow/worker/WindmillSink.java | 3 ++- .../java/org/apache/beam/sdk/values/ValueInSingleWindow.java | 1 + 2 files changed, 3 insertions(+), 1 deletion(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java index 6d3484f6642d..d54d94f47d7c 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java @@ -121,7 +121,8 @@ public static BeamFnApi.Elements.ElementMetadata decodeAdditionalMetadata( PaneInfo paneInfo = PaneInfoCoder.INSTANCE.decode(inStream); windowsCoder.decode(inStream); if (paneInfo.isElementMetadata()) { - return BeamFnApi.Elements.ElementMetadata.parseFrom(ByteArrayCoder.of().decode(inStream)); + return BeamFnApi.Elements.ElementMetadata.parseFrom( + ByteArrayCoder.of().decode(inStream, Coder.Context.OUTER)); } else { // empty return BeamFnApi.Elements.ElementMetadata.newBuilder().build(); diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/values/ValueInSingleWindow.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/values/ValueInSingleWindow.java index 2bbc28cac9fe..21df11119831 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/values/ValueInSingleWindow.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/values/ValueInSingleWindow.java @@ -133,6 +133,7 @@ public ValueInSingleWindow decode(InputStream inStream) throws IOException { } @Override + @SuppressWarnings("IgnoredPureGetter") public ValueInSingleWindow decode(InputStream inStream, Context context) throws IOException { Instant timestamp = InstantCoder.of().decode(inStream); BoundedWindow window = windowCoder.decode(inStream); From a4054cda70fb2f03e6fd85313e17eaea75c6b883 Mon Sep 17 00:00:00 2001 From: Radek Stankiewicz Date: Mon, 13 Oct 2025 15:11:12 +0200 Subject: [PATCH 9/9] fix bugs --- .../src/main/java/org/apache/beam/sdk/values/WindowedValues.java | 1 + 1 file changed, 1 insertion(+) diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/values/WindowedValues.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/values/WindowedValues.java index d90637da9b1b..99e9d5e83a64 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/values/WindowedValues.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/values/WindowedValues.java @@ -857,6 +857,7 @@ public WindowedValue decode(InputStream inStream) throws CoderException, IOEx } @Override + @SuppressWarnings("IgnoredPureGetter") public WindowedValue decode(InputStream inStream, Context context) throws CoderException, IOException { Instant timestamp = InstantCoder.of().decode(inStream);