From 85853a3edf0f792c3cb8a8f70ae0f27b712c7fd5 Mon Sep 17 00:00:00 2001 From: Radek Stankiewicz Date: Thu, 9 Oct 2025 13:48:11 +0200 Subject: [PATCH 01/10] proto change --- .../beam/model/fn_execution/v1/beam_fn_api.proto | 10 +++++++++- 1 file changed, 9 insertions(+), 1 deletion(-) diff --git a/model/fn-execution/src/main/proto/org/apache/beam/model/fn_execution/v1/beam_fn_api.proto b/model/fn-execution/src/main/proto/org/apache/beam/model/fn_execution/v1/beam_fn_api.proto index 9360522ab409..9b32048b4995 100644 --- a/model/fn-execution/src/main/proto/org/apache/beam/model/fn_execution/v1/beam_fn_api.proto +++ b/model/fn-execution/src/main/proto/org/apache/beam/model/fn_execution/v1/beam_fn_api.proto @@ -740,10 +740,18 @@ message Elements { bool is_last = 4; } + message DrainMode { + enum Enum { + UNSPECIFIED = 0; + NOT_DRAINING = 1; + DRAINING = 2; + } + } + // Element metadata passed as part of WindowedValue to make WindowedValue // extensible and backward compatible message ElementMetadata { - // empty message - add drain, kind, tracing metadata in the future + optional DrainMode.Enum drain = 1; } // Represent the encoded user timer for a given instruction, transform and From a6d2b7dabf5802b3af6da1b186e31f0e879fcb82 Mon Sep 17 00:00:00 2001 From: Radek Stankiewicz Date: Thu, 9 Oct 2025 13:50:14 +0200 Subject: [PATCH 02/10] add draining to output builder, encode draining --- .../apache/beam/sdk/values/OutputBuilder.java | 2 + .../apache/beam/sdk/values/WindowedValue.java | 3 + .../beam/sdk/values/WindowedValues.java | 152 +++++++++++++----- 3 files changed, 121 insertions(+), 36 deletions(-) diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/values/OutputBuilder.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/values/OutputBuilder.java index a7f8bc8e03b1..5762d32ae832 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/values/OutputBuilder.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/values/OutputBuilder.java @@ -48,5 +48,7 @@ public interface OutputBuilder extends WindowedValue { OutputBuilder setRecordOffset(@Nullable Long recordOffset); + OutputBuilder setDraining(@Nullable Boolean drain); + void output(); } 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 ea6be129ecb4..762a602bc3f9 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 @@ -52,6 +52,9 @@ public interface WindowedValue { @Nullable Long getRecordOffset(); + @Nullable + Boolean isDraining(); + /** * A representation of each of the actual values represented by this compressed {@link * WindowedValue}, one per window. 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 99e9d5e83a64..1f458a4a7729 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 @@ -99,6 +99,7 @@ public static class Builder implements OutputBuilder { private @MonotonicNonNull Collection windows; private @Nullable String recordId; private @Nullable Long recordOffset; + private @Nullable Boolean draining; @Override public Builder setValue(T value) { @@ -142,6 +143,12 @@ public Builder setRecordOffset(@Nullable Long recordOffset) { return this; } + @Override + public Builder setDraining(@Nullable Boolean draining) { + this.draining = draining; + return this; + } + public Builder setReceiver(WindowedValueReceiver receiver) { this.receiver = receiver; return this; @@ -190,6 +197,11 @@ public PaneInfo getPaneInfo() { return recordOffset; } + @Override + public @Nullable Boolean isDraining() { + return draining; + } + @Override public Collection> explodeWindows() { throw new UnsupportedOperationException( @@ -218,7 +230,8 @@ public void output() { } public WindowedValue build() { - return WindowedValues.of(getValue(), getTimestamp(), getWindows(), getPaneInfo()); + return WindowedValues.of( + getValue(), getTimestamp(), getWindows(), getPaneInfo(), null, null, isDraining()); } @Override @@ -228,6 +241,7 @@ public String toString() { .add("timestamp", getTimestamp()) .add("windows", getWindows()) .add("paneInfo", getPaneInfo()) + .add("draining", isDraining()) .add("receiver", receiver) .toString(); } @@ -235,7 +249,7 @@ public String toString() { public static WindowedValue of( T value, Instant timestamp, Collection windows, PaneInfo paneInfo) { - return of(value, timestamp, windows, paneInfo, null, null); + return of(value, timestamp, windows, paneInfo, null, null, null); } /** Returns a {@code WindowedValue} with the given value, timestamp, and windows. */ @@ -245,27 +259,32 @@ public static WindowedValue of( Collection windows, PaneInfo paneInfo, @Nullable String currentRecordId, - @Nullable Long currentRecordOffset) { + @Nullable Long currentRecordOffset, + @Nullable Boolean draining) { checkArgument(paneInfo != 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(), paneInfo); + return of(value, timestamp, windows.iterator().next(), paneInfo, draining); } else { return new TimestampedValueInMultipleWindows<>( - value, timestamp, windows, paneInfo, currentRecordId, currentRecordOffset); + value, timestamp, windows, paneInfo, currentRecordId, currentRecordOffset, draining); } } /** @deprecated for use only in compatibility with old broken code */ @Deprecated static WindowedValue createWithoutValidation( - T value, Instant timestamp, Collection windows, PaneInfo paneInfo) { + T value, + Instant timestamp, + Collection windows, + PaneInfo paneInfo, + @Nullable Boolean draining) { if (windows.size() == 1) { - return of(value, timestamp, windows.iterator().next(), paneInfo); + return of(value, timestamp, windows.iterator().next(), paneInfo, draining); } else { return new TimestampedValueInMultipleWindows<>( - value, timestamp, windows, paneInfo, null, null); + value, timestamp, windows, paneInfo, null, null, draining); } } @@ -274,13 +293,26 @@ public static WindowedValue of( T value, Instant timestamp, BoundedWindow window, PaneInfo paneInfo) { checkArgument(paneInfo != null, "WindowedValue requires PaneInfo, but it was null"); + return of(value, timestamp, window, paneInfo, null); + } + + /** Returns a {@code WindowedValue} with the given value, timestamp, and window. */ + public static WindowedValue of( + T value, + Instant timestamp, + BoundedWindow window, + PaneInfo paneInfo, + @Nullable Boolean draining) { + checkArgument(paneInfo != 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, paneInfo); } else if (isGlobal) { - return new TimestampedValueInGlobalWindow<>(value, timestamp, paneInfo, null, null); + return new TimestampedValueInGlobalWindow<>(value, timestamp, paneInfo, null, null, draining); } else { - return new TimestampedValueInSingleWindow<>(value, timestamp, window, paneInfo, null, null); + return new TimestampedValueInSingleWindow<>( + value, timestamp, window, paneInfo, null, null, draining); } } @@ -289,7 +321,7 @@ public static WindowedValue of( * default timestamp and pane. */ public static WindowedValue valueInGlobalWindow(T value) { - return new ValueInGlobalWindow<>(value, PaneInfo.NO_FIRING, null, null); + return new ValueInGlobalWindow<>(value, PaneInfo.NO_FIRING, null, null, null); } /** @@ -297,7 +329,7 @@ public static WindowedValue valueInGlobalWindow(T value) { * default timestamp and the specified pane. */ public static WindowedValue valueInGlobalWindow(T value, PaneInfo paneInfo) { - return new ValueInGlobalWindow<>(value, paneInfo, null, null); + return new ValueInGlobalWindow<>(value, paneInfo, null, null, null); } /** @@ -308,7 +340,8 @@ public static WindowedValue timestampedValueInGlobalWindow(T value, Insta if (BoundedWindow.TIMESTAMP_MIN_VALUE.equals(timestamp)) { return valueInGlobalWindow(value); } else { - return new TimestampedValueInGlobalWindow<>(value, timestamp, PaneInfo.NO_FIRING, null, null); + return new TimestampedValueInGlobalWindow<>( + value, timestamp, PaneInfo.NO_FIRING, null, null, null); } } @@ -321,7 +354,7 @@ public static WindowedValue timestampedValueInGlobalWindow( if (paneInfo.equals(PaneInfo.NO_FIRING)) { return timestampedValueInGlobalWindow(value, timestamp); } else { - return new TimestampedValueInGlobalWindow<>(value, timestamp, paneInfo, null, null); + return new TimestampedValueInGlobalWindow<>(value, timestamp, paneInfo, null, null, null); } } @@ -337,7 +370,8 @@ public static WindowedValue withValue( windowedValue.getWindows(), windowedValue.getPaneInfo(), windowedValue.getRecordId(), - windowedValue.getRecordOffset()); + windowedValue.getRecordOffset(), + windowedValue.isDraining()); } public static boolean equals( @@ -388,6 +422,7 @@ private abstract static class SimpleWindowedValue implements WindowedValue private final PaneInfo paneInfo; private final @Nullable String currentRecordId; private final @Nullable Long currentRecordOffset; + private final @Nullable Boolean draining; @Override public @Nullable String getRecordId() { @@ -399,15 +434,22 @@ private abstract static class SimpleWindowedValue implements WindowedValue return currentRecordOffset; } + @Override + public @Nullable Boolean isDraining() { + return draining; + } + protected SimpleWindowedValue( T value, PaneInfo paneInfo, @Nullable String currentRecordId, - @Nullable Long currentRecordOffset) { + @Nullable Long currentRecordOffset, + @Nullable Boolean draining) { this.value = value; this.paneInfo = checkNotNull(paneInfo); this.currentRecordId = currentRecordId; this.currentRecordOffset = currentRecordOffset; + this.draining = draining; } @Override @@ -455,8 +497,9 @@ public MinTimestampWindowedValue( T value, PaneInfo pane, @Nullable String currentRecordId, - @Nullable Long currentRecordOffset) { - super(value, pane, currentRecordId, currentRecordOffset); + @Nullable Long currentRecordOffset, + @Nullable Boolean draining) { + super(value, pane, currentRecordId, currentRecordOffset, draining); } @Override @@ -473,8 +516,9 @@ public ValueInGlobalWindow( T value, PaneInfo paneInfo, @Nullable String currentRecordId, - @Nullable Long currentRecordOffset) { - super(value, paneInfo, currentRecordId, currentRecordOffset); + @Nullable Long currentRecordOffset, + @Nullable Boolean draining) { + super(value, paneInfo, currentRecordId, currentRecordOffset, draining); } @Override @@ -489,7 +533,8 @@ public BoundedWindow getWindow() { @Override public WindowedValue withValue(NewT newValue) { - return new ValueInGlobalWindow<>(newValue, getPaneInfo(), getRecordId(), getRecordOffset()); + return new ValueInGlobalWindow<>( + newValue, getPaneInfo(), getRecordId(), getRecordOffset(), isDraining()); } @Override @@ -513,6 +558,7 @@ public String toString() { return MoreObjects.toStringHelper(getClass()) .add("value", getValue()) .add("paneInfo", getPaneInfo()) + .add("draining", isDraining()) .toString(); } } @@ -526,8 +572,9 @@ public TimestampedWindowedValue( Instant timestamp, PaneInfo paneInfo, @Nullable String currentRecordId, - @Nullable Long currentRecordOffset) { - super(value, paneInfo, currentRecordId, currentRecordOffset); + @Nullable Long currentRecordOffset, + @Nullable Boolean draining) { + super(value, paneInfo, currentRecordId, currentRecordOffset, draining); this.timestamp = checkNotNull(timestamp); } @@ -549,8 +596,9 @@ public TimestampedValueInGlobalWindow( Instant timestamp, PaneInfo paneInfo, @Nullable String currentRecordId, - @Nullable Long currentRecordOffset) { - super(value, timestamp, paneInfo, currentRecordId, currentRecordOffset); + @Nullable Long currentRecordOffset, + @Nullable Boolean draining) { + super(value, timestamp, paneInfo, currentRecordId, currentRecordOffset, draining); } @Override @@ -566,7 +614,7 @@ public BoundedWindow getWindow() { @Override public WindowedValue withValue(NewT newValue) { return new TimestampedValueInGlobalWindow<>( - newValue, getTimestamp(), getPaneInfo(), getRecordId(), getRecordOffset()); + newValue, getTimestamp(), getPaneInfo(), getRecordId(), getRecordOffset(), isDraining()); } @Override @@ -596,6 +644,7 @@ public String toString() { .add("value", getValue()) .add("timestamp", getTimestamp()) .add("paneInfo", getPaneInfo()) + .add("draining", isDraining()) .toString(); } } @@ -615,15 +664,22 @@ public TimestampedValueInSingleWindow( BoundedWindow window, PaneInfo paneInfo, @Nullable String currentRecordId, - @Nullable Long currentRecordOffset) { - super(value, timestamp, paneInfo, currentRecordId, currentRecordOffset); + @Nullable Long currentRecordOffset, + @Nullable Boolean draining) { + super(value, timestamp, paneInfo, currentRecordId, currentRecordOffset, draining); this.window = checkNotNull(window); } @Override public WindowedValue withValue(NewT newValue) { return new TimestampedValueInSingleWindow<>( - newValue, getTimestamp(), window, getPaneInfo(), getRecordId(), getRecordOffset()); + newValue, + getTimestamp(), + window, + getPaneInfo(), + getRecordId(), + getRecordOffset(), + isDraining()); } @Override @@ -665,6 +721,7 @@ public String toString() { .add("timestamp", getTimestamp()) .add("window", window) .add("paneInfo", getPaneInfo()) + .add("draining", isDraining()) .toString(); } } @@ -679,8 +736,9 @@ public TimestampedValueInMultipleWindows( Collection windows, PaneInfo paneInfo, @Nullable String currentRecordId, - @Nullable Long currentRecordOffset) { - super(value, timestamp, paneInfo, currentRecordId, currentRecordOffset); + @Nullable Long currentRecordOffset, + @Nullable Boolean draining) { + super(value, timestamp, paneInfo, currentRecordId, currentRecordOffset, draining); this.windows = checkNotNull(windows); } @@ -692,7 +750,13 @@ public Collection getWindows() { @Override public WindowedValue withValue(NewT newValue) { return new TimestampedValueInMultipleWindows<>( - newValue, getTimestamp(), getWindows(), getPaneInfo(), getRecordId(), getRecordOffset()); + newValue, + getTimestamp(), + getWindows(), + getPaneInfo(), + getRecordId(), + getRecordOffset(), + isDraining()); } @Override @@ -730,6 +794,7 @@ public String toString() { .add("timestamp", getTimestamp()) .add("windows", windows) .add("paneInfo", getPaneInfo()) + .add("draining", isDraining()) .toString(); } @@ -845,7 +910,16 @@ public void encode(WindowedValue windowedElem, OutputStream outStream, Contex if (metadataSupported) { BeamFnApi.Elements.ElementMetadata.Builder builder = BeamFnApi.Elements.ElementMetadata.newBuilder(); - BeamFnApi.Elements.ElementMetadata em = builder.build(); + BeamFnApi.Elements.ElementMetadata em = + builder + .setDrain( + windowedElem.isDraining() != null + ? (Boolean.TRUE.equals(windowedElem.isDraining()) + ? BeamFnApi.Elements.DrainMode.Enum.DRAINING + : BeamFnApi.Elements.DrainMode.Enum.NOT_DRAINING) + : BeamFnApi.Elements.DrainMode.Enum.UNSPECIFIED) + .build(); + ByteArrayCoder.of().encode(em.toByteArray(), outStream); } valueCoder.encode(windowedElem.getValue(), outStream, context); @@ -857,20 +931,26 @@ 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); Collection windows = windowsCoder.decode(inStream); PaneInfo paneInfo = PaneInfoCoder.INSTANCE.decode(inStream); + Boolean draining = null; if (isMetadataSupported() && paneInfo.isElementMetadata()) { - BeamFnApi.Elements.ElementMetadata.parseFrom(ByteArrayCoder.of().decode(inStream)); + BeamFnApi.Elements.ElementMetadata elementMetadata = + BeamFnApi.Elements.ElementMetadata.parseFrom(ByteArrayCoder.of().decode(inStream)); + boolean b = elementMetadata.hasDrain(); + draining = + b + ? elementMetadata.getDrain().equals(BeamFnApi.Elements.DrainMode.Enum.DRAINING) + : null; } T value = valueCoder.decode(inStream, context); // 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, paneInfo); + return WindowedValues.createWithoutValidation(value, timestamp, windows, paneInfo, draining); } @Override From 951943e9bba2eb63129e6269cda0f16556aa0dcf Mon Sep 17 00:00:00 2001 From: Radek Stankiewicz Date: Wed, 15 Oct 2025 12:02:26 +0200 Subject: [PATCH 03/10] add draining to output builder --- .../java/org/apache/beam/sdk/util/WindowedValueTest.java | 6 +++++- 1 file changed, 5 insertions(+), 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 3e3973e3720b..50e2f8f506fc 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 @@ -89,7 +89,10 @@ public void testWindowedValueWithElementMetadataCoder() throws CoderException { new IntervalWindow(timestamp, timestamp.plus(Duration.millis(1000))), new IntervalWindow( timestamp.plus(Duration.millis(1000)), timestamp.plus(Duration.millis(2000)))), - PaneInfo.NO_FIRING); + PaneInfo.NO_FIRING, + null, + null, + true); Coder> windowedValueCoder = WindowedValues.getFullCoder(StringUtf8Coder.of(), IntervalWindow.getCoder()); @@ -101,6 +104,7 @@ public void testWindowedValueWithElementMetadataCoder() throws CoderException { Assert.assertEquals(value.getValue(), decodedValue.getValue()); Assert.assertEquals(value.getTimestamp(), decodedValue.getTimestamp()); Assert.assertArrayEquals(value.getWindows().toArray(), decodedValue.getWindows().toArray()); + Assert.assertTrue(value.isDraining()); } @Test From 7ea3109a0e4256cdfe91413a1eb665dd3728162f Mon Sep 17 00:00:00 2001 From: Radek Stankiewicz Date: Wed, 15 Oct 2025 15:43:53 +0200 Subject: [PATCH 04/10] default impls --- .../org/apache/beam/runners/dataflow/BatchViewOverrides.java | 5 +++++ .../runners/dataflow/worker/util/ValueInEmptyWindows.java | 5 +++++ .../java/org/apache/beam/runners/spark/util/TimerUtils.java | 5 +++++ 3 files changed, 15 insertions(+) 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 10b41bb5b5ba..7026396a42b9 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 @@ -1378,6 +1378,11 @@ public T getValue() { return value; } + @Override + public @Nullable Boolean isDraining() { + return null; + } + @Override public Instant getTimestamp() { return BoundedWindow.TIMESTAMP_MIN_VALUE; 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 a51c9ed419e1..90b31b974e2f 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 @@ -59,6 +59,11 @@ public PaneInfo getPaneInfo() { return null; } + @Override + public @Nullable Boolean isDraining() { + return null; + } + @Override public Iterable> explodeWindows() { return Collections.emptyList(); diff --git a/runners/spark/src/main/java/org/apache/beam/runners/spark/util/TimerUtils.java b/runners/spark/src/main/java/org/apache/beam/runners/spark/util/TimerUtils.java index 03735355de51..0760c1aeb649 100644 --- a/runners/spark/src/main/java/org/apache/beam/runners/spark/util/TimerUtils.java +++ b/runners/spark/src/main/java/org/apache/beam/runners/spark/util/TimerUtils.java @@ -115,6 +115,11 @@ public PaneInfo getPaneInfo() { return null; } + @Override + public @Nullable Boolean isDraining() { + return null; + } + @Override public @Nullable Long getRecordOffset() { return null; From 3d9e403b3ed8614567702da3af43c6aacea3b545 Mon Sep 17 00:00:00 2001 From: Radek Stankiewicz Date: Wed, 15 Oct 2025 18:58:04 +0200 Subject: [PATCH 05/10] comment --- .../test/java/org/apache/beam/sdk/util/WindowedValueTest.java | 2 +- 1 file changed, 1 insertion(+), 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 50e2f8f506fc..be37f54f35bb 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 @@ -92,7 +92,7 @@ public void testWindowedValueWithElementMetadataCoder() throws CoderException { PaneInfo.NO_FIRING, null, null, - true); + true); // drain is persisted as part of metadata Coder> windowedValueCoder = WindowedValues.getFullCoder(StringUtf8Coder.of(), IntervalWindow.getCoder()); From 499039fe37def138bb4c3df5c19607b9658b4e06 Mon Sep 17 00:00:00 2001 From: Radek Stankiewicz Date: Wed, 15 Oct 2025 19:35:00 +0200 Subject: [PATCH 06/10] remove nullable --- .../runners/dataflow/BatchViewOverrides.java | 4 ++-- .../worker/util/ValueInEmptyWindows.java | 4 ++-- .../beam/runners/spark/util/TimerUtils.java | 4 ++-- .../apache/beam/sdk/values/OutputBuilder.java | 2 +- .../apache/beam/sdk/values/WindowedValue.java | 3 +-- .../beam/sdk/values/WindowedValues.java | 20 +++++++++---------- 6 files changed, 17 insertions(+), 20 deletions(-) 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 7026396a42b9..3fd46eb9b0de 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 @@ -1379,8 +1379,8 @@ public T getValue() { } @Override - public @Nullable Boolean isDraining() { - return null; + public boolean isDraining() { + return false; } @Override 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 90b31b974e2f..cbc673b15c0f 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 @@ -60,8 +60,8 @@ public PaneInfo getPaneInfo() { } @Override - public @Nullable Boolean isDraining() { - return null; + public boolean isDraining() { + return false; } @Override diff --git a/runners/spark/src/main/java/org/apache/beam/runners/spark/util/TimerUtils.java b/runners/spark/src/main/java/org/apache/beam/runners/spark/util/TimerUtils.java index 0760c1aeb649..162144ca283f 100644 --- a/runners/spark/src/main/java/org/apache/beam/runners/spark/util/TimerUtils.java +++ b/runners/spark/src/main/java/org/apache/beam/runners/spark/util/TimerUtils.java @@ -116,8 +116,8 @@ public PaneInfo getPaneInfo() { } @Override - public @Nullable Boolean isDraining() { - return null; + public boolean isDraining() { + return false; } @Override diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/values/OutputBuilder.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/values/OutputBuilder.java index 5762d32ae832..05b72d52264b 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/values/OutputBuilder.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/values/OutputBuilder.java @@ -48,7 +48,7 @@ public interface OutputBuilder extends WindowedValue { OutputBuilder setRecordOffset(@Nullable Long recordOffset); - OutputBuilder setDraining(@Nullable Boolean drain); + OutputBuilder setDraining(boolean drain); void output(); } 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 762a602bc3f9..3097c8e33a92 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 @@ -52,8 +52,7 @@ public interface WindowedValue { @Nullable Long getRecordOffset(); - @Nullable - Boolean isDraining(); + boolean isDraining(); /** * A representation of each of the actual values represented by this compressed {@link 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 1f458a4a7729..f5a243a0b27b 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 @@ -99,7 +99,7 @@ public static class Builder implements OutputBuilder { private @MonotonicNonNull Collection windows; private @Nullable String recordId; private @Nullable Long recordOffset; - private @Nullable Boolean draining; + private boolean draining; @Override public Builder setValue(T value) { @@ -144,7 +144,7 @@ public Builder setRecordOffset(@Nullable Long recordOffset) { } @Override - public Builder setDraining(@Nullable Boolean draining) { + public Builder setDraining(boolean draining) { this.draining = draining; return this; } @@ -198,7 +198,7 @@ public PaneInfo getPaneInfo() { } @Override - public @Nullable Boolean isDraining() { + public boolean isDraining() { return draining; } @@ -435,7 +435,7 @@ private abstract static class SimpleWindowedValue implements WindowedValue } @Override - public @Nullable Boolean isDraining() { + public boolean isDraining() { return draining; } @@ -913,11 +913,9 @@ public void encode(WindowedValue windowedElem, OutputStream outStream, Contex BeamFnApi.Elements.ElementMetadata em = builder .setDrain( - windowedElem.isDraining() != null - ? (Boolean.TRUE.equals(windowedElem.isDraining()) - ? BeamFnApi.Elements.DrainMode.Enum.DRAINING - : BeamFnApi.Elements.DrainMode.Enum.NOT_DRAINING) - : BeamFnApi.Elements.DrainMode.Enum.UNSPECIFIED) + Boolean.TRUE.equals(windowedElem.isDraining()) + ? BeamFnApi.Elements.DrainMode.Enum.DRAINING + : BeamFnApi.Elements.DrainMode.Enum.NOT_DRAINING) .build(); ByteArrayCoder.of().encode(em.toByteArray(), outStream); @@ -936,7 +934,7 @@ public WindowedValue decode(InputStream inStream, Context context) Instant timestamp = InstantCoder.of().decode(inStream); Collection windows = windowsCoder.decode(inStream); PaneInfo paneInfo = PaneInfoCoder.INSTANCE.decode(inStream); - Boolean draining = null; + boolean draining = false; if (isMetadataSupported() && paneInfo.isElementMetadata()) { BeamFnApi.Elements.ElementMetadata elementMetadata = BeamFnApi.Elements.ElementMetadata.parseFrom(ByteArrayCoder.of().decode(inStream)); @@ -944,7 +942,7 @@ public WindowedValue decode(InputStream inStream, Context context) draining = b ? elementMetadata.getDrain().equals(BeamFnApi.Elements.DrainMode.Enum.DRAINING) - : null; + : false; } T value = valueCoder.decode(inStream, context); From cca50bff5aa89b1c6b402e21cb2aa917ccd9f760 Mon Sep 17 00:00:00 2001 From: Radek Stankiewicz Date: Wed, 15 Oct 2025 19:39:24 +0200 Subject: [PATCH 07/10] remove nullable --- .../beam/sdk/values/WindowedValues.java | 34 ++++++++----------- 1 file changed, 15 insertions(+), 19 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 f5a243a0b27b..91fd49ee92ff 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 @@ -293,16 +293,12 @@ public static WindowedValue of( T value, Instant timestamp, BoundedWindow window, PaneInfo paneInfo) { checkArgument(paneInfo != null, "WindowedValue requires PaneInfo, but it was null"); - return of(value, timestamp, window, paneInfo, null); + return of(value, timestamp, window, paneInfo, false); } /** Returns a {@code WindowedValue} with the given value, timestamp, and window. */ public static WindowedValue of( - T value, - Instant timestamp, - BoundedWindow window, - PaneInfo paneInfo, - @Nullable Boolean draining) { + T value, Instant timestamp, BoundedWindow window, PaneInfo paneInfo, boolean draining) { checkArgument(paneInfo != null, "WindowedValue requires PaneInfo, but it was null"); boolean isGlobal = GlobalWindow.INSTANCE.equals(window); @@ -321,7 +317,7 @@ public static WindowedValue of( * default timestamp and pane. */ public static WindowedValue valueInGlobalWindow(T value) { - return new ValueInGlobalWindow<>(value, PaneInfo.NO_FIRING, null, null, null); + return new ValueInGlobalWindow<>(value, PaneInfo.NO_FIRING, null, null, false); } /** @@ -329,7 +325,7 @@ public static WindowedValue valueInGlobalWindow(T value) { * default timestamp and the specified pane. */ public static WindowedValue valueInGlobalWindow(T value, PaneInfo paneInfo) { - return new ValueInGlobalWindow<>(value, paneInfo, null, null, null); + return new ValueInGlobalWindow<>(value, paneInfo, null, null, false); } /** @@ -341,7 +337,7 @@ public static WindowedValue timestampedValueInGlobalWindow(T value, Insta return valueInGlobalWindow(value); } else { return new TimestampedValueInGlobalWindow<>( - value, timestamp, PaneInfo.NO_FIRING, null, null, null); + value, timestamp, PaneInfo.NO_FIRING, null, null, false); } } @@ -354,7 +350,7 @@ public static WindowedValue timestampedValueInGlobalWindow( if (paneInfo.equals(PaneInfo.NO_FIRING)) { return timestampedValueInGlobalWindow(value, timestamp); } else { - return new TimestampedValueInGlobalWindow<>(value, timestamp, paneInfo, null, null, null); + return new TimestampedValueInGlobalWindow<>(value, timestamp, paneInfo, null, null, false); } } @@ -422,7 +418,7 @@ private abstract static class SimpleWindowedValue implements WindowedValue private final PaneInfo paneInfo; private final @Nullable String currentRecordId; private final @Nullable Long currentRecordOffset; - private final @Nullable Boolean draining; + private final boolean draining; @Override public @Nullable String getRecordId() { @@ -444,7 +440,7 @@ protected SimpleWindowedValue( PaneInfo paneInfo, @Nullable String currentRecordId, @Nullable Long currentRecordOffset, - @Nullable Boolean draining) { + boolean draining) { this.value = value; this.paneInfo = checkNotNull(paneInfo); this.currentRecordId = currentRecordId; @@ -498,7 +494,7 @@ public MinTimestampWindowedValue( PaneInfo pane, @Nullable String currentRecordId, @Nullable Long currentRecordOffset, - @Nullable Boolean draining) { + boolean draining) { super(value, pane, currentRecordId, currentRecordOffset, draining); } @@ -517,7 +513,7 @@ public ValueInGlobalWindow( PaneInfo paneInfo, @Nullable String currentRecordId, @Nullable Long currentRecordOffset, - @Nullable Boolean draining) { + boolean draining) { super(value, paneInfo, currentRecordId, currentRecordOffset, draining); } @@ -573,7 +569,7 @@ public TimestampedWindowedValue( PaneInfo paneInfo, @Nullable String currentRecordId, @Nullable Long currentRecordOffset, - @Nullable Boolean draining) { + boolean draining) { super(value, paneInfo, currentRecordId, currentRecordOffset, draining); this.timestamp = checkNotNull(timestamp); } @@ -597,7 +593,7 @@ public TimestampedValueInGlobalWindow( PaneInfo paneInfo, @Nullable String currentRecordId, @Nullable Long currentRecordOffset, - @Nullable Boolean draining) { + boolean draining) { super(value, timestamp, paneInfo, currentRecordId, currentRecordOffset, draining); } @@ -665,7 +661,7 @@ public TimestampedValueInSingleWindow( PaneInfo paneInfo, @Nullable String currentRecordId, @Nullable Long currentRecordOffset, - @Nullable Boolean draining) { + boolean draining) { super(value, timestamp, paneInfo, currentRecordId, currentRecordOffset, draining); this.window = checkNotNull(window); } @@ -737,7 +733,7 @@ public TimestampedValueInMultipleWindows( PaneInfo paneInfo, @Nullable String currentRecordId, @Nullable Long currentRecordOffset, - @Nullable Boolean draining) { + boolean draining) { super(value, timestamp, paneInfo, currentRecordId, currentRecordOffset, draining); this.windows = checkNotNull(windows); } @@ -913,7 +909,7 @@ public void encode(WindowedValue windowedElem, OutputStream outStream, Contex BeamFnApi.Elements.ElementMetadata em = builder .setDrain( - Boolean.TRUE.equals(windowedElem.isDraining()) + windowedElem.isDraining() ? BeamFnApi.Elements.DrainMode.Enum.DRAINING : BeamFnApi.Elements.DrainMode.Enum.NOT_DRAINING) .build(); From f8037f0f5bf0abca803edbb2d9d57c528146a546 Mon Sep 17 00:00:00 2001 From: Radek Stankiewicz Date: Wed, 15 Oct 2025 19:51:55 +0200 Subject: [PATCH 08/10] remove nullable --- .../java/org/apache/beam/sdk/values/WindowedValues.java | 6 +++--- 1 file changed, 3 insertions(+), 3 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 91fd49ee92ff..639102c28ba0 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 @@ -249,7 +249,7 @@ public String toString() { public static WindowedValue of( T value, Instant timestamp, Collection windows, PaneInfo paneInfo) { - return of(value, timestamp, windows, paneInfo, null, null, null); + return of(value, timestamp, windows, paneInfo, null, null, false); } /** Returns a {@code WindowedValue} with the given value, timestamp, and windows. */ @@ -260,7 +260,7 @@ public static WindowedValue of( PaneInfo paneInfo, @Nullable String currentRecordId, @Nullable Long currentRecordOffset, - @Nullable Boolean draining) { + boolean draining) { checkArgument(paneInfo != null, "WindowedValue requires PaneInfo, but it was null"); checkArgument(windows.size() > 0, "WindowedValue requires windows, but there were none"); @@ -279,7 +279,7 @@ static WindowedValue createWithoutValidation( Instant timestamp, Collection windows, PaneInfo paneInfo, - @Nullable Boolean draining) { + boolean draining) { if (windows.size() == 1) { return of(value, timestamp, windows.iterator().next(), paneInfo, draining); } else { From ac99abf3d794924ad55b7c5c9dceabd6ceaec4bb Mon Sep 17 00:00:00 2001 From: Radek Stankiewicz Date: Fri, 24 Oct 2025 14:32:54 +0200 Subject: [PATCH 09/10] rename field --- .../beam/runners/spark/util/TimerUtils.java | 2 +- .../apache/beam/sdk/values/OutputBuilder.java | 2 +- .../apache/beam/sdk/values/WindowedValue.java | 2 +- .../beam/sdk/values/WindowedValues.java | 99 ++++++++++--------- .../beam/sdk/util/WindowedValueTest.java | 2 +- 5 files changed, 57 insertions(+), 50 deletions(-) diff --git a/runners/spark/src/main/java/org/apache/beam/runners/spark/util/TimerUtils.java b/runners/spark/src/main/java/org/apache/beam/runners/spark/util/TimerUtils.java index 162144ca283f..0be36d67388c 100644 --- a/runners/spark/src/main/java/org/apache/beam/runners/spark/util/TimerUtils.java +++ b/runners/spark/src/main/java/org/apache/beam/runners/spark/util/TimerUtils.java @@ -116,7 +116,7 @@ public PaneInfo getPaneInfo() { } @Override - public boolean isDraining() { + public boolean causedByDrain() { return false; } diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/values/OutputBuilder.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/values/OutputBuilder.java index 05b72d52264b..03e3088e5256 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/values/OutputBuilder.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/values/OutputBuilder.java @@ -48,7 +48,7 @@ public interface OutputBuilder extends WindowedValue { OutputBuilder setRecordOffset(@Nullable Long recordOffset); - OutputBuilder setDraining(boolean drain); + OutputBuilder setCausedByDrain(boolean causedByDrain); void output(); } 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 3097c8e33a92..bcd58b903171 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 @@ -52,7 +52,7 @@ public interface WindowedValue { @Nullable Long getRecordOffset(); - boolean isDraining(); + boolean causedByDrain(); /** * A representation of each of the actual values represented by this compressed {@link 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 639102c28ba0..518b9a62647e 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 @@ -99,7 +99,7 @@ public static class Builder implements OutputBuilder { private @MonotonicNonNull Collection windows; private @Nullable String recordId; private @Nullable Long recordOffset; - private boolean draining; + private boolean causedByDrain; @Override public Builder setValue(T value) { @@ -144,8 +144,8 @@ public Builder setRecordOffset(@Nullable Long recordOffset) { } @Override - public Builder setDraining(boolean draining) { - this.draining = draining; + public Builder setCausedByDrain(boolean causedByDrain) { + this.causedByDrain = causedByDrain; return this; } @@ -198,8 +198,8 @@ public PaneInfo getPaneInfo() { } @Override - public boolean isDraining() { - return draining; + public boolean causedByDrain() { + return causedByDrain; } @Override @@ -231,7 +231,7 @@ public void output() { public WindowedValue build() { return WindowedValues.of( - getValue(), getTimestamp(), getWindows(), getPaneInfo(), null, null, isDraining()); + getValue(), getTimestamp(), getWindows(), getPaneInfo(), null, null, causedByDrain()); } @Override @@ -241,7 +241,7 @@ public String toString() { .add("timestamp", getTimestamp()) .add("windows", getWindows()) .add("paneInfo", getPaneInfo()) - .add("draining", isDraining()) + .add("causedByDrain", causedByDrain()) .add("receiver", receiver) .toString(); } @@ -260,15 +260,15 @@ public static WindowedValue of( PaneInfo paneInfo, @Nullable String currentRecordId, @Nullable Long currentRecordOffset, - boolean draining) { + boolean causedByDrain) { checkArgument(paneInfo != 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(), paneInfo, draining); + return of(value, timestamp, windows.iterator().next(), paneInfo, causedByDrain); } else { return new TimestampedValueInMultipleWindows<>( - value, timestamp, windows, paneInfo, currentRecordId, currentRecordOffset, draining); + value, timestamp, windows, paneInfo, currentRecordId, currentRecordOffset, causedByDrain); } } @@ -279,12 +279,12 @@ static WindowedValue createWithoutValidation( Instant timestamp, Collection windows, PaneInfo paneInfo, - boolean draining) { + boolean causedByDrain) { if (windows.size() == 1) { - return of(value, timestamp, windows.iterator().next(), paneInfo, draining); + return of(value, timestamp, windows.iterator().next(), paneInfo, causedByDrain); } else { return new TimestampedValueInMultipleWindows<>( - value, timestamp, windows, paneInfo, null, null, draining); + value, timestamp, windows, paneInfo, null, null, causedByDrain); } } @@ -298,17 +298,18 @@ public static WindowedValue of( /** Returns a {@code WindowedValue} with the given value, timestamp, and window. */ public static WindowedValue of( - T value, Instant timestamp, BoundedWindow window, PaneInfo paneInfo, boolean draining) { + T value, Instant timestamp, BoundedWindow window, PaneInfo paneInfo, boolean causedByDrain) { checkArgument(paneInfo != 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, paneInfo); } else if (isGlobal) { - return new TimestampedValueInGlobalWindow<>(value, timestamp, paneInfo, null, null, draining); + return new TimestampedValueInGlobalWindow<>( + value, timestamp, paneInfo, null, null, causedByDrain); } else { return new TimestampedValueInSingleWindow<>( - value, timestamp, window, paneInfo, null, null, draining); + value, timestamp, window, paneInfo, null, null, causedByDrain); } } @@ -367,7 +368,7 @@ public static WindowedValue withValue( windowedValue.getPaneInfo(), windowedValue.getRecordId(), windowedValue.getRecordOffset(), - windowedValue.isDraining()); + windowedValue.causedByDrain()); } public static boolean equals( @@ -418,7 +419,7 @@ private abstract static class SimpleWindowedValue implements WindowedValue private final PaneInfo paneInfo; private final @Nullable String currentRecordId; private final @Nullable Long currentRecordOffset; - private final boolean draining; + private final boolean causedByDrain; @Override public @Nullable String getRecordId() { @@ -431,8 +432,8 @@ private abstract static class SimpleWindowedValue implements WindowedValue } @Override - public boolean isDraining() { - return draining; + public boolean causedByDrain() { + return causedByDrain; } protected SimpleWindowedValue( @@ -440,12 +441,12 @@ protected SimpleWindowedValue( PaneInfo paneInfo, @Nullable String currentRecordId, @Nullable Long currentRecordOffset, - boolean draining) { + boolean causedByDrain) { this.value = value; this.paneInfo = checkNotNull(paneInfo); this.currentRecordId = currentRecordId; this.currentRecordOffset = currentRecordOffset; - this.draining = draining; + this.causedByDrain = causedByDrain; } @Override @@ -494,8 +495,8 @@ public MinTimestampWindowedValue( PaneInfo pane, @Nullable String currentRecordId, @Nullable Long currentRecordOffset, - boolean draining) { - super(value, pane, currentRecordId, currentRecordOffset, draining); + boolean causedByDrain) { + super(value, pane, currentRecordId, currentRecordOffset, causedByDrain); } @Override @@ -513,8 +514,8 @@ public ValueInGlobalWindow( PaneInfo paneInfo, @Nullable String currentRecordId, @Nullable Long currentRecordOffset, - boolean draining) { - super(value, paneInfo, currentRecordId, currentRecordOffset, draining); + boolean causedByDrain) { + super(value, paneInfo, currentRecordId, currentRecordOffset, causedByDrain); } @Override @@ -530,7 +531,7 @@ public BoundedWindow getWindow() { @Override public WindowedValue withValue(NewT newValue) { return new ValueInGlobalWindow<>( - newValue, getPaneInfo(), getRecordId(), getRecordOffset(), isDraining()); + newValue, getPaneInfo(), getRecordId(), getRecordOffset(), causedByDrain()); } @Override @@ -554,7 +555,7 @@ public String toString() { return MoreObjects.toStringHelper(getClass()) .add("value", getValue()) .add("paneInfo", getPaneInfo()) - .add("draining", isDraining()) + .add("causedByDrain", causedByDrain()) .toString(); } } @@ -569,8 +570,8 @@ public TimestampedWindowedValue( PaneInfo paneInfo, @Nullable String currentRecordId, @Nullable Long currentRecordOffset, - boolean draining) { - super(value, paneInfo, currentRecordId, currentRecordOffset, draining); + boolean causedByDrain) { + super(value, paneInfo, currentRecordId, currentRecordOffset, causedByDrain); this.timestamp = checkNotNull(timestamp); } @@ -593,8 +594,8 @@ public TimestampedValueInGlobalWindow( PaneInfo paneInfo, @Nullable String currentRecordId, @Nullable Long currentRecordOffset, - boolean draining) { - super(value, timestamp, paneInfo, currentRecordId, currentRecordOffset, draining); + boolean causedByDrain) { + super(value, timestamp, paneInfo, currentRecordId, currentRecordOffset, causedByDrain); } @Override @@ -610,7 +611,12 @@ public BoundedWindow getWindow() { @Override public WindowedValue withValue(NewT newValue) { return new TimestampedValueInGlobalWindow<>( - newValue, getTimestamp(), getPaneInfo(), getRecordId(), getRecordOffset(), isDraining()); + newValue, + getTimestamp(), + getPaneInfo(), + getRecordId(), + getRecordOffset(), + causedByDrain()); } @Override @@ -640,7 +646,7 @@ public String toString() { .add("value", getValue()) .add("timestamp", getTimestamp()) .add("paneInfo", getPaneInfo()) - .add("draining", isDraining()) + .add("causedByDrain", causedByDrain()) .toString(); } } @@ -661,8 +667,8 @@ public TimestampedValueInSingleWindow( PaneInfo paneInfo, @Nullable String currentRecordId, @Nullable Long currentRecordOffset, - boolean draining) { - super(value, timestamp, paneInfo, currentRecordId, currentRecordOffset, draining); + boolean causedByDrain) { + super(value, timestamp, paneInfo, currentRecordId, currentRecordOffset, causedByDrain); this.window = checkNotNull(window); } @@ -675,7 +681,7 @@ public WindowedValue withValue(NewT newValue) { getPaneInfo(), getRecordId(), getRecordOffset(), - isDraining()); + causedByDrain()); } @Override @@ -717,7 +723,7 @@ public String toString() { .add("timestamp", getTimestamp()) .add("window", window) .add("paneInfo", getPaneInfo()) - .add("draining", isDraining()) + .add("causedByDrain", causedByDrain()) .toString(); } } @@ -733,8 +739,8 @@ public TimestampedValueInMultipleWindows( PaneInfo paneInfo, @Nullable String currentRecordId, @Nullable Long currentRecordOffset, - boolean draining) { - super(value, timestamp, paneInfo, currentRecordId, currentRecordOffset, draining); + boolean causedByDrain) { + super(value, timestamp, paneInfo, currentRecordId, currentRecordOffset, causedByDrain); this.windows = checkNotNull(windows); } @@ -752,7 +758,7 @@ public WindowedValue withValue(NewT newValue) { getPaneInfo(), getRecordId(), getRecordOffset(), - isDraining()); + causedByDrain()); } @Override @@ -790,7 +796,7 @@ public String toString() { .add("timestamp", getTimestamp()) .add("windows", windows) .add("paneInfo", getPaneInfo()) - .add("draining", isDraining()) + .add("causedByDrain", causedByDrain()) .toString(); } @@ -909,7 +915,7 @@ public void encode(WindowedValue windowedElem, OutputStream outStream, Contex BeamFnApi.Elements.ElementMetadata em = builder .setDrain( - windowedElem.isDraining() + windowedElem.causedByDrain() ? BeamFnApi.Elements.DrainMode.Enum.DRAINING : BeamFnApi.Elements.DrainMode.Enum.NOT_DRAINING) .build(); @@ -930,12 +936,12 @@ public WindowedValue decode(InputStream inStream, Context context) Instant timestamp = InstantCoder.of().decode(inStream); Collection windows = windowsCoder.decode(inStream); PaneInfo paneInfo = PaneInfoCoder.INSTANCE.decode(inStream); - boolean draining = false; + boolean causedByDrain = false; if (isMetadataSupported() && paneInfo.isElementMetadata()) { BeamFnApi.Elements.ElementMetadata elementMetadata = BeamFnApi.Elements.ElementMetadata.parseFrom(ByteArrayCoder.of().decode(inStream)); boolean b = elementMetadata.hasDrain(); - draining = + causedByDrain = b ? elementMetadata.getDrain().equals(BeamFnApi.Elements.DrainMode.Enum.DRAINING) : false; @@ -944,7 +950,8 @@ public WindowedValue decode(InputStream inStream, Context context) // 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, paneInfo, draining); + return WindowedValues.createWithoutValidation( + value, timestamp, windows, paneInfo, causedByDrain); } @Override 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 be37f54f35bb..915399311859 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 @@ -104,7 +104,7 @@ public void testWindowedValueWithElementMetadataCoder() throws CoderException { Assert.assertEquals(value.getValue(), decodedValue.getValue()); Assert.assertEquals(value.getTimestamp(), decodedValue.getTimestamp()); Assert.assertArrayEquals(value.getWindows().toArray(), decodedValue.getWindows().toArray()); - Assert.assertTrue(value.isDraining()); + Assert.assertTrue(value.causedByDrain()); } @Test From fc4b29c42561916d687e8a68d06ce8b05d7a1c62 Mon Sep 17 00:00:00 2001 From: Radek Stankiewicz Date: Fri, 24 Oct 2025 15:02:52 +0200 Subject: [PATCH 10/10] rename field --- .../org/apache/beam/runners/dataflow/BatchViewOverrides.java | 2 +- .../beam/runners/dataflow/worker/util/ValueInEmptyWindows.java | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) 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 3fd46eb9b0de..e7bb4dc9c0ac 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 @@ -1379,7 +1379,7 @@ public T getValue() { } @Override - public boolean isDraining() { + public boolean causedByDrain() { return false; } 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 cbc673b15c0f..00bb282c6845 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 @@ -60,7 +60,7 @@ public PaneInfo getPaneInfo() { } @Override - public boolean isDraining() { + public boolean causedByDrain() { return false; }