From 6dd67e6e58117898d9ca47ab7cad0277caa906b0 Mon Sep 17 00:00:00 2001 From: Radek Stankiewicz Date: Fri, 6 Jun 2025 13:05:01 +0200 Subject: [PATCH 1/4] element metadata capability --- .../runners/dataflow/BatchViewOverrides.java | 19 +- .../worker/StreamingDataflowWorker.java | 5 + .../worker/util/ValueInEmptyWindows.java | 17 ++ .../sdk/transforms/windowing/PaneInfo.java | 64 ++++-- .../apache/beam/sdk/util/ElementMetadata.java | 28 +++ .../sdk/util/construction/Environments.java | 1 + .../apache/beam/sdk/values/WindowedValue.java | 7 + .../beam/sdk/values/WindowedValues.java | 203 +++++++++++++----- .../transforms/windowing/PaneInfoTest.java | 20 ++ .../beam/sdk/util/WindowedValueTest.java | 28 +++ .../org/apache/beam/fn/harness/FnHarness.java | 9 + 11 files changed, 337 insertions(+), 64 deletions(-) create mode 100644 sdks/java/core/src/main/java/org/apache/beam/sdk/util/ElementMetadata.java diff --git a/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/BatchViewOverrides.java b/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/BatchViewOverrides.java index 8cb82c2ca42d..2605628d5e69 100644 --- a/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/BatchViewOverrides.java +++ b/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/BatchViewOverrides.java @@ -63,6 +63,7 @@ import org.apache.beam.sdk.transforms.windowing.PaneInfo; import org.apache.beam.sdk.transforms.windowing.Window; import org.apache.beam.sdk.util.CoderUtils; +import org.apache.beam.sdk.util.ElementMetadata; import org.apache.beam.sdk.util.SystemDoFnInternal; import org.apache.beam.sdk.values.KV; import org.apache.beam.sdk.values.PCollection; @@ -1362,12 +1363,19 @@ private static WindowedValue valueInEmptyWindows(T value) { return new ValueInEmptyWindows<>(value); } - private static class ValueInEmptyWindows implements WindowedValue { + private static class ValueInEmptyWindows extends WindowedValue { private final T value; + private final @Nullable ElementMetadata elementMetadata; private ValueInEmptyWindows(T value) { this.value = value; + this.elementMetadata = null; + } + + private ValueInEmptyWindows(T value, @Nullable ElementMetadata elementMetadata) { + this.value = value; + this.elementMetadata = elementMetadata; } @Override @@ -1375,6 +1383,11 @@ public WindowedValue withValue(NewT value) { return new ValueInEmptyWindows<>(value); } + @Override + public WindowedValue withElementMetadata(@Nullable ElementMetadata elementMetadata) { + return new ValueInEmptyWindows<>(this.getValue(), elementMetadata); + } + @Override public T getValue() { return value; @@ -1396,8 +1409,8 @@ public PaneInfo getPane() { } @Override - public Iterable> explodeWindows() { - return Collections.emptyList(); + public @Nullable ElementMetadata getElementMetadata() { + return elementMetadata; } @Override 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..15b78934cc44 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 @@ -109,6 +109,8 @@ import org.apache.beam.sdk.io.FileSystems; import org.apache.beam.sdk.io.gcp.bigquery.BigQuerySinkMetrics; import org.apache.beam.sdk.metrics.MetricsEnvironment; +import org.apache.beam.sdk.transforms.windowing.PaneInfo; +import org.apache.beam.sdk.util.WindowedValue; import org.apache.beam.sdk.util.construction.CoderTranslation; 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; @@ -816,6 +818,9 @@ public static void main(String[] args) throws Exception { CoderTranslation.verifyModelCodersRegistered(); + WindowedValue.FullWindowedValueCoder.setMetadataSupported(); + PaneInfo.PaneInfoCoder.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/util/ValueInEmptyWindows.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/ValueInEmptyWindows.java index ca0e2279fb03..ff328e679c4b 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/ValueInEmptyWindows.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/ValueInEmptyWindows.java @@ -22,6 +22,7 @@ import java.util.Objects; import org.apache.beam.sdk.transforms.windowing.BoundedWindow; import org.apache.beam.sdk.transforms.windowing.PaneInfo; +import org.apache.beam.sdk.util.ElementMetadata; import org.apache.beam.sdk.values.WindowedValue; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.MoreObjects; import org.checkerframework.checker.nullness.qual.Nullable; @@ -39,9 +40,15 @@ */ public class ValueInEmptyWindows implements WindowedValue { private final T value; + private final @Nullable ElementMetadata elementMetadata; public ValueInEmptyWindows(T value) { + this(value, null); + } + + public ValueInEmptyWindows(T value, @Nullable ElementMetadata elementMetadata) { this.value = value; + this.elementMetadata = elementMetadata; } @Override @@ -54,6 +61,11 @@ public Iterable> explodeWindows() { return Collections.emptyList(); } + @Override + public @Nullable ElementMetadata getElementMetadata() { + return elementMetadata; + } + @Override public T getValue() { return value; @@ -64,6 +76,11 @@ public WindowedValue withValue(NewT newValue) { return new ValueInEmptyWindows<>(newValue); } + @Override + public WindowedValue withElementMetadata(@Nullable ElementMetadata elementMetadata) { + return new ValueInEmptyWindows<>(this.getValue(), elementMetadata); + } + @Override public Instant getTimestamp() { return BoundedWindow.TIMESTAMP_MIN_VALUE; 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..6b59821b5c55 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,13 @@ public String toString() { /** A Coder for encoding PaneInfo instances. */ public static class PaneInfoCoder extends AtomicCoder { + private static final byte ELEMENT_METADATA_MASK = (byte) 0x80; + protected static boolean metadataSupported = false; + + public static void setMetadataSupported() { + metadataSupported = true; + } + private enum Encoding { FIRST, ONE_INDEX, @@ -337,16 +372,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 +396,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 +411,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/util/ElementMetadata.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/ElementMetadata.java new file mode 100644 index 000000000000..fa34b433c8ef --- /dev/null +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/ElementMetadata.java @@ -0,0 +1,28 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.beam.sdk.util; + +import com.google.auto.value.AutoValue; + +@AutoValue +public abstract class ElementMetadata { + + public static ElementMetadata create() { + return new AutoValue_ElementMetadata(); + } +} diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/Environments.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/Environments.java index 05ecb21fd956..fb37e7507aa2 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/Environments.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/Environments.java @@ -495,6 +495,7 @@ public static Set getJavaCapabilities() { capabilities.add(BeamUrns.getUrn(StandardProtocols.Enum.DATA_SAMPLING)); capabilities.add(BeamUrns.getUrn(StandardProtocols.Enum.SDK_CONSUMING_RECEIVED_DATA)); capabilities.add(BeamUrns.getUrn(StandardProtocols.Enum.ORDERED_LIST_STATE)); + capabilities.add(BeamUrns.getUrn(StandardProtocols.Enum.ELEMENT_METADATA)); return capabilities.build(); } diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/values/WindowedValue.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/values/WindowedValue.java index 2a5236f0147f..90f6b5e5d467 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/values/WindowedValue.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/values/WindowedValue.java @@ -20,6 +20,8 @@ import java.util.Collection; import org.apache.beam.sdk.transforms.windowing.BoundedWindow; import org.apache.beam.sdk.transforms.windowing.PaneInfo; +import org.apache.beam.sdk.util.ElementMetadata; +import org.checkerframework.checker.nullness.qual.Nullable; import org.joda.time.Instant; /** @@ -34,6 +36,9 @@ public interface WindowedValue { /** The timestamp of this value in event time. */ Instant getTimestamp(); + @Nullable + ElementMetadata getElementMetadata(); + /** Returns the windows of this {@code WindowedValue}. */ Collection getWindows(); @@ -57,4 +62,6 @@ default PaneInfo getPaneInfo() { * value. */ WindowedValue withValue(OtherT value); + + WindowedValue withElementMetadata(@Nullable ElementMetadata elementMetadata); } 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 3c044990de37..f5b4640aa1bc 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 @@ -34,6 +34,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; @@ -45,6 +46,7 @@ import org.apache.beam.sdk.transforms.windowing.GlobalWindow; import org.apache.beam.sdk.transforms.windowing.PaneInfo; import org.apache.beam.sdk.transforms.windowing.PaneInfo.PaneInfoCoder; +import org.apache.beam.sdk.util.ElementMetadata; import org.apache.beam.sdk.util.common.ElementByteSizeObserver; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.MoreObjects; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList; @@ -63,40 +65,64 @@ private WindowedValues() {} // non-instantiable utility class /** Returns a {@code WindowedValue} with the given value, timestamp, and windows. */ public static WindowedValue of( - T value, Instant timestamp, Collection windows, PaneInfo pane) { + T value, + Instant timestamp, + Collection windows, + PaneInfo pane, + @Nullable ElementMetadata elementMetadata) { checkArgument(pane != null, "WindowedValue requires PaneInfo, but it was null"); checkArgument(windows.size() > 0, "WindowedValue requires windows, but there were none"); - if (windows.size() == 1) { - return of(value, timestamp, windows.iterator().next(), pane); + return of(value, timestamp, windows.iterator().next(), pane, elementMetadata); } else { - return new TimestampedValueInMultipleWindows<>(value, timestamp, windows, pane); + return new TimestampedValueInMultipleWindows<>( + value, timestamp, windows, pane, elementMetadata); } } + public static WindowedValue of( + T value, Instant timestamp, Collection windows, PaneInfo pane) { + return of(value, timestamp, windows, pane, null); + } + /** @deprecated for use only in compatibility with old broken code */ @Deprecated static WindowedValue createWithoutValidation( - T value, Instant timestamp, Collection windows, PaneInfo pane) { + T value, + Instant timestamp, + Collection windows, + PaneInfo pane, + @Nullable ElementMetadata elementMetadata) { if (windows.size() == 1) { - return of(value, timestamp, windows.iterator().next(), pane); + return of(value, timestamp, windows.iterator().next(), pane, elementMetadata); } else { - return new TimestampedValueInMultipleWindows<>(value, timestamp, windows, pane); + return new TimestampedValueInMultipleWindows<>( + value, timestamp, windows, pane, elementMetadata); } } /** Returns a {@code WindowedValue} with the given value, timestamp, and window. */ public static WindowedValue of( T value, Instant timestamp, BoundedWindow window, PaneInfo pane) { + return of(value, timestamp, window, pane, null); + } + + public static WindowedValue of( + T value, + Instant timestamp, + BoundedWindow window, + PaneInfo pane, + @Nullable ElementMetadata elementMetadata) { + checkArgument(pane != null, "WindowedValue requires PaneInfo, but it was null"); boolean isGlobal = GlobalWindow.INSTANCE.equals(window); if (isGlobal && BoundedWindow.TIMESTAMP_MIN_VALUE.equals(timestamp)) { return valueInGlobalWindow(value, pane); } else if (isGlobal) { - return new TimestampedValueInGlobalWindow<>(value, timestamp, pane); + return new TimestampedValueInGlobalWindow<>(value, timestamp, pane, elementMetadata); } else { - return new TimestampedValueInSingleWindow<>(value, timestamp, window, pane); + return new TimestampedValueInSingleWindow<>(value, timestamp, window, pane, elementMetadata); } } @@ -141,17 +167,9 @@ public static WindowedValue timestampedValueInGlobalWindow( } } - /** - * Returns a new {@code WindowedValue} that is a copy of this one, but with a different value, - * which may have a new type {@code NewT}. - */ - public static WindowedValue withValue( - WindowedValue windowedValue, NewT newValue) { - return WindowedValues.of( - newValue, - windowedValue.getTimestamp(), - windowedValue.getWindows(), - windowedValue.getPaneInfo()); + /** Returns {@code true} if this WindowedValue has exactly one window. */ + public boolean isSingleWindowedValue() { + return false; } public static boolean equals( @@ -200,10 +218,13 @@ private abstract static class SimpleWindowedValue implements WindowedValue private final T value; private final PaneInfo pane; + private final @Nullable ElementMetadata elementMetadata; - protected SimpleWindowedValue(T value, PaneInfo pane) { + protected SimpleWindowedValue( + T value, PaneInfo pane, @Nullable ElementMetadata elementMetadata) { this.value = value; this.pane = checkNotNull(pane); + this.elementMetadata = elementMetadata; } @Override @@ -211,6 +232,11 @@ public PaneInfo getPane() { return pane; } + @Override + public @Nullable ElementMetadata getElementMetadata() { + return elementMetadata; + } + @Override public T getValue() { return value; @@ -232,8 +258,9 @@ public Iterable> explodeWindows() { /** The abstract superclass of WindowedValue representations where timestamp == MIN. */ private abstract static class MinTimestampWindowedValue extends SimpleWindowedValue { - public MinTimestampWindowedValue(T value, PaneInfo pane) { - super(value, pane); + public MinTimestampWindowedValue( + T value, PaneInfo pane, @Nullable ElementMetadata elementMetadata) { + super(value, pane, elementMetadata); } @Override @@ -246,8 +273,22 @@ public Instant getTimestamp() { private static class ValueInGlobalWindow extends MinTimestampWindowedValue implements SingleWindowedValue { + public ValueInGlobalWindow(T value, PaneInfo pane, @Nullable ElementMetadata elementMetadata) { + super(value, pane, elementMetadata); + } + public ValueInGlobalWindow(T value, PaneInfo pane) { - super(value, pane); + this(value, pane, null); + } + + @Override + public WindowedValue withValue(NewT newValue) { + return new ValueInGlobalWindow<>(newValue, getPane(), getElementMetadata()); + } + + @Override + public WindowedValue withElementMetadata(@Nullable ElementMetadata elementMetadata) { + return new ValueInGlobalWindow<>(getValue(), getPane(), elementMetadata); } @Override @@ -260,11 +301,6 @@ public BoundedWindow getWindow() { return GlobalWindow.INSTANCE; } - @Override - public WindowedValue withValue(NewT newValue) { - return new ValueInGlobalWindow<>(newValue, getPane()); - } - @Override public boolean equals(@Nullable Object o) { if (o instanceof ValueInGlobalWindow) { @@ -294,8 +330,9 @@ public String toString() { private abstract static class TimestampedWindowedValue extends SimpleWindowedValue { private final Instant timestamp; - public TimestampedWindowedValue(T value, Instant timestamp, PaneInfo pane) { - super(value, pane); + public TimestampedWindowedValue( + T value, Instant timestamp, PaneInfo pane, @Nullable ElementMetadata elementMetadata) { + super(value, pane, elementMetadata); this.timestamp = checkNotNull(timestamp); } @@ -312,8 +349,25 @@ public Instant getTimestamp() { private static class TimestampedValueInGlobalWindow extends TimestampedWindowedValue implements SingleWindowedValue { + public TimestampedValueInGlobalWindow( + T value, Instant timestamp, PaneInfo pane, @Nullable ElementMetadata elementMetadata) { + super(value, timestamp, pane, elementMetadata); + } + public TimestampedValueInGlobalWindow(T value, Instant timestamp, PaneInfo pane) { - super(value, timestamp, pane); + this(value, timestamp, pane, null); + } + + @Override + public WindowedValue withValue(NewT newValue) { + return new TimestampedValueInGlobalWindow<>( + newValue, getTimestamp(), getPane(), getElementMetadata()); + } + + @Override + public WindowedValue withElementMetadata(@Nullable ElementMetadata elementMetadata) { + return new TimestampedValueInGlobalWindow<>( + getValue(), getTimestamp(), getPane(), elementMetadata); } @Override @@ -326,11 +380,6 @@ public BoundedWindow getWindow() { return GlobalWindow.INSTANCE; } - @Override - public WindowedValue withValue(NewT newValue) { - return new TimestampedValueInGlobalWindow<>(newValue, getTimestamp(), getPane()); - } - @Override public boolean equals(@Nullable Object o) { if (o instanceof TimestampedValueInGlobalWindow) { @@ -372,14 +421,25 @@ private static class TimestampedValueInSingleWindow extends TimestampedWindow private final BoundedWindow window; public TimestampedValueInSingleWindow( - T value, Instant timestamp, BoundedWindow window, PaneInfo pane) { - super(value, timestamp, pane); + T value, + Instant timestamp, + BoundedWindow window, + PaneInfo pane, + @Nullable ElementMetadata elementMetadata) { + super(value, timestamp, pane, elementMetadata); this.window = checkNotNull(window); } @Override public WindowedValue withValue(NewT newValue) { - return new TimestampedValueInSingleWindow<>(newValue, getTimestamp(), window, getPane()); + return new TimestampedValueInSingleWindow<>( + newValue, getTimestamp(), window, getPane(), getElementMetadata()); + } + + @Override + public WindowedValue withElementMetadata(@Nullable ElementMetadata elementMetadata) { + return new TimestampedValueInSingleWindow<>( + getValue(), getTimestamp(), getWindow(), getPane(), elementMetadata); } @Override @@ -425,25 +485,36 @@ public String toString() { } } + /** The representation of a WindowedValue, excluding the special cases captured above. */ /** The representation of a WindowedValue, excluding the special cases captured above. */ private static class TimestampedValueInMultipleWindows extends TimestampedWindowedValue { private Collection windows; public TimestampedValueInMultipleWindows( - T value, Instant timestamp, Collection windows, PaneInfo pane) { - super(value, timestamp, pane); + T value, + Instant timestamp, + Collection windows, + PaneInfo pane, + @Nullable ElementMetadata elementMetadata) { + super(value, timestamp, pane, elementMetadata); this.windows = checkNotNull(windows); } @Override - public Collection getWindows() { - return windows; + public WindowedValue withValue(NewT newValue) { + return new TimestampedValueInMultipleWindows<>( + newValue, getTimestamp(), windows, getPane(), getElementMetadata()); } @Override - public WindowedValue withValue(NewT newValue) { + public WindowedValue withElementMetadata(@Nullable ElementMetadata elementMetadata) { return new TimestampedValueInMultipleWindows<>( - newValue, getTimestamp(), getWindows(), getPane()); + getValue(), getTimestamp(), getWindows(), getPane(), elementMetadata); + } + + @Override + public Collection getWindows() { + return windows; } @Override @@ -515,6 +586,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); @@ -538,6 +618,12 @@ public static class FullWindowedValueCoder extends WindowedValueCoder { // Precompute and cache the coder for a list of windows. private final Coder> windowsCoder; + private static boolean metadataSupported = false; + + public static void setMetadataSupported() { + metadataSupported = true; + } + public static FullWindowedValueCoder of( Coder valueCoder, Coder windowCoder) { return new FullWindowedValueCoder<>(valueCoder, windowCoder); @@ -582,7 +668,16 @@ public void encode(WindowedValue windowedElem, OutputStream outStream, Contex InstantCoder.of().encode(windowedElem.getTimestamp(), outStream); windowsCoder.encode(windowedElem.getWindows(), outStream); PaneInfoCoder.INSTANCE.encode(windowedElem.getPane(), outStream); + if (isMetadataSupported()) { + // propagate metadata only if runner and other environments support this capability + // ElementMetadata elementMetadata = windowedElem.getElementMetadata(); + BeamFnApi.Elements.ElementMetadata.Builder builder = + BeamFnApi.Elements.ElementMetadata.newBuilder(); + BeamFnApi.Elements.ElementMetadata em = builder.build(); + em.writeDelimitedTo(outStream); + } valueCoder.encode(windowedElem.getValue(), outStream, context); + // todo add support for metadata } @Override @@ -596,11 +691,23 @@ public WindowedValue decode(InputStream inStream, Context context) Instant timestamp = InstantCoder.of().decode(inStream); Collection windows = windowsCoder.decode(inStream); PaneInfo pane = PaneInfoCoder.INSTANCE.decode(inStream); + ElementMetadata elementMetadata = ElementMetadata.create(); + if (isMetadataSupported() && pane.isElementMetadata()) { + // read metadata only if runner and other environments support this capability + // read metadata only if pane has provided information about additional metadata + BeamFnApi.Elements.ElementMetadata metadata = + BeamFnApi.Elements.ElementMetadata.parseDelimitedFrom(inStream); + if (metadata != null) { + elementMetadata = ElementMetadata.create(); + } + } T value = valueCoder.decode(inStream, context); + // todo add support for metadata // Because there are some remaining (incorrect) uses of WindowedValue with no windows, // we call this deprecated no-validation path when decoding - return WindowedValues.createWithoutValidation(value, timestamp, windows, pane); + return WindowedValues.createWithoutValidation( + value, timestamp, windows, pane, elementMetadata); } @Override @@ -794,7 +901,7 @@ public WindowedValue decode(InputStream inStream) throws CoderException, IOEx @Override public WindowedValue decode(InputStream inStream, Context context) throws CoderException, IOException { - return WindowedValues.withValue(windowedValuePrototype, valueCoder.decode(inStream, context)); + return windowedValuePrototype.withValue(valueCoder.decode(inStream, context)); } @Override 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 db1579333e57..73cb4e9e57e6 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,34 @@ public void testWindowedValueCoder() throws CoderException { Assert.assertArrayEquals(value.getWindows().toArray(), decodedValue.getWindows().toArray()); } + @Test + public void testWindowedValueWithElementMetadataCoder() throws CoderException { + 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, + ElementMetadata.create()); + + 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()); + ElementMetadata elementMetadata = decodedValue.getElementMetadata(); + Assert.assertNotNull(elementMetadata); + } + @Test public void testFullWindowedValueCoderIsSerializableWithWellKnownCoderType() { CoderProperties.coderSerializable( diff --git a/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/FnHarness.java b/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/FnHarness.java index 831337072b06..67e3bcdef14a 100644 --- a/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/FnHarness.java +++ b/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/FnHarness.java @@ -49,6 +49,7 @@ import org.apache.beam.model.fnexecution.v1.BeamFnApi.ProcessBundleDescriptor; import org.apache.beam.model.fnexecution.v1.BeamFnControlGrpc; import org.apache.beam.model.pipeline.v1.Endpoints; +import org.apache.beam.model.pipeline.v1.RunnerApi.StandardProtocols; import org.apache.beam.runners.core.metrics.MetricsContainerImpl; import org.apache.beam.runners.core.metrics.ShortIdMap; import org.apache.beam.sdk.fn.IdGenerator; @@ -64,6 +65,9 @@ import org.apache.beam.sdk.options.ExperimentalOptions; import org.apache.beam.sdk.options.PipelineOptions; import org.apache.beam.sdk.options.SdkHarnessOptions; +import org.apache.beam.sdk.transforms.windowing.PaneInfo; +import org.apache.beam.sdk.util.WindowedValue; +import org.apache.beam.sdk.util.construction.BeamUrns; import org.apache.beam.sdk.util.construction.CoderTranslation; import org.apache.beam.sdk.util.construction.PipelineOptionsTranslation; import org.apache.beam.vendor.grpc.v1p69p0.com.google.protobuf.TextFormat; @@ -195,6 +199,11 @@ public static void main(Function environmentVarGetter) throws Ex ? Collections.emptySet() : ImmutableSet.copyOf(runnerCapabilitesOrNull.split("\\s+")); + if (runnerCapabilites.contains(BeamUrns.getUrn(StandardProtocols.Enum.ELEMENT_METADATA))) { + WindowedValue.FullWindowedValueCoder.setMetadataSupported(); + PaneInfo.PaneInfoCoder.setMetadataSupported(); + } + main( id, options, From 9a38d8d39697f62ab27ff3f01335223f227a253e Mon Sep 17 00:00:00 2001 From: Radek Stankiewicz Date: Fri, 6 Jun 2025 13:13:44 +0200 Subject: [PATCH 2/4] element metadata capability --- .../java/org/apache/beam/sdk/values/WindowedValues.java | 6 ------ 1 file changed, 6 deletions(-) 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 f5b4640aa1bc..a98ae7feef04 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 @@ -618,12 +618,6 @@ public static class FullWindowedValueCoder extends WindowedValueCoder { // Precompute and cache the coder for a list of windows. private final Coder> windowsCoder; - private static boolean metadataSupported = false; - - public static void setMetadataSupported() { - metadataSupported = true; - } - public static FullWindowedValueCoder of( Coder valueCoder, Coder windowCoder) { return new FullWindowedValueCoder<>(valueCoder, windowCoder); From 6778135136ddcc072024ffe9165165f1c8affb60 Mon Sep 17 00:00:00 2001 From: Radek Stankiewicz Date: Fri, 6 Jun 2025 13:27:13 +0200 Subject: [PATCH 3/4] element metadata capability --- .../beam/runners/dataflow/worker/StreamingDataflowWorker.java | 4 ++-- .../src/main/java/org/apache/beam/fn/harness/FnHarness.java | 4 ++-- 2 files changed, 4 insertions(+), 4 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 15b78934cc44..cece4391547b 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,8 +110,8 @@ import org.apache.beam.sdk.io.gcp.bigquery.BigQuerySinkMetrics; import org.apache.beam.sdk.metrics.MetricsEnvironment; import org.apache.beam.sdk.transforms.windowing.PaneInfo; -import org.apache.beam.sdk.util.WindowedValue; 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; @@ -818,7 +818,7 @@ public static void main(String[] args) throws Exception { CoderTranslation.verifyModelCodersRegistered(); - WindowedValue.FullWindowedValueCoder.setMetadataSupported(); + WindowedValues.FullWindowedValueCoder.setMetadataSupported(); PaneInfo.PaneInfoCoder.setMetadataSupported(); LOG.debug("Creating StreamingDataflowWorker from options: {}", options); diff --git a/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/FnHarness.java b/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/FnHarness.java index 67e3bcdef14a..10200ace03cd 100644 --- a/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/FnHarness.java +++ b/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/FnHarness.java @@ -66,10 +66,10 @@ import org.apache.beam.sdk.options.PipelineOptions; import org.apache.beam.sdk.options.SdkHarnessOptions; import org.apache.beam.sdk.transforms.windowing.PaneInfo; -import org.apache.beam.sdk.util.WindowedValue; import org.apache.beam.sdk.util.construction.BeamUrns; import org.apache.beam.sdk.util.construction.CoderTranslation; import org.apache.beam.sdk.util.construction.PipelineOptionsTranslation; +import org.apache.beam.sdk.values.WindowedValues; import org.apache.beam.vendor.grpc.v1p69p0.com.google.protobuf.TextFormat; import org.apache.beam.vendor.grpc.v1p69p0.io.grpc.ManagedChannel; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting; @@ -200,7 +200,7 @@ public static void main(Function environmentVarGetter) throws Ex : ImmutableSet.copyOf(runnerCapabilitesOrNull.split("\\s+")); if (runnerCapabilites.contains(BeamUrns.getUrn(StandardProtocols.Enum.ELEMENT_METADATA))) { - WindowedValue.FullWindowedValueCoder.setMetadataSupported(); + WindowedValues.FullWindowedValueCoder.setMetadataSupported(); PaneInfo.PaneInfoCoder.setMetadataSupported(); } From f470b0e77a84bf58aa6be36ff9cf1f662b317706 Mon Sep 17 00:00:00 2001 From: Radek Stankiewicz Date: Fri, 6 Jun 2025 13:55:04 +0200 Subject: [PATCH 4/4] fix BatchViewOverrides --- .../apache/beam/runners/dataflow/BatchViewOverrides.java | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/BatchViewOverrides.java b/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/BatchViewOverrides.java index 2605628d5e69..7fc000d9853d 100644 --- a/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/BatchViewOverrides.java +++ b/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/BatchViewOverrides.java @@ -1363,7 +1363,7 @@ private static WindowedValue valueInEmptyWindows(T value) { return new ValueInEmptyWindows<>(value); } - private static class ValueInEmptyWindows extends WindowedValue { + private static class ValueInEmptyWindows implements WindowedValue { private final T value; private final @Nullable ElementMetadata elementMetadata; @@ -1408,6 +1408,11 @@ public PaneInfo getPane() { return PaneInfo.NO_FIRING; } + @Override + public Iterable> explodeWindows() { + return Collections.emptyList(); + } + @Override public @Nullable ElementMetadata getElementMetadata() { return elementMetadata;