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 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..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 @@ -1378,6 +1378,11 @@ public T getValue() { return value; } + @Override + public boolean causedByDrain() { + return false; + } + @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..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 @@ -59,6 +59,11 @@ public PaneInfo getPaneInfo() { return null; } + @Override + public boolean causedByDrain() { + return false; + } + @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..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 @@ -115,6 +115,11 @@ public PaneInfo getPaneInfo() { return null; } + @Override + public boolean causedByDrain() { + return false; + } + @Override public @Nullable Long getRecordOffset() { return null; 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..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,5 +48,7 @@ public interface OutputBuilder extends WindowedValue { OutputBuilder setRecordOffset(@Nullable Long recordOffset); + 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 ea6be129ecb4..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,6 +52,8 @@ public interface WindowedValue { @Nullable Long getRecordOffset(); + boolean causedByDrain(); + /** * 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..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,6 +99,7 @@ public static class Builder implements OutputBuilder { private @MonotonicNonNull Collection windows; private @Nullable String recordId; private @Nullable Long recordOffset; + private boolean causedByDrain; @Override public Builder setValue(T value) { @@ -142,6 +143,12 @@ public Builder setRecordOffset(@Nullable Long recordOffset) { return this; } + @Override + public Builder setCausedByDrain(boolean causedByDrain) { + this.causedByDrain = causedByDrain; + return this; + } + public Builder setReceiver(WindowedValueReceiver receiver) { this.receiver = receiver; return this; @@ -190,6 +197,11 @@ public PaneInfo getPaneInfo() { return recordOffset; } + @Override + public boolean causedByDrain() { + return causedByDrain; + } + @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, causedByDrain()); } @Override @@ -228,6 +241,7 @@ public String toString() { .add("timestamp", getTimestamp()) .add("windows", getWindows()) .add("paneInfo", getPaneInfo()) + .add("causedByDrain", causedByDrain()) .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, false); } /** 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, + 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); + return of(value, timestamp, windows.iterator().next(), paneInfo, causedByDrain); } else { return new TimestampedValueInMultipleWindows<>( - value, timestamp, windows, paneInfo, currentRecordId, currentRecordOffset); + value, timestamp, windows, paneInfo, currentRecordId, currentRecordOffset, causedByDrain); } } /** @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, + boolean causedByDrain) { if (windows.size() == 1) { - return of(value, timestamp, windows.iterator().next(), paneInfo); + return of(value, timestamp, windows.iterator().next(), paneInfo, causedByDrain); } else { return new TimestampedValueInMultipleWindows<>( - value, timestamp, windows, paneInfo, null, null); + value, timestamp, windows, paneInfo, null, null, causedByDrain); } } @@ -274,13 +293,23 @@ 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, false); + } + + /** Returns a {@code WindowedValue} with the given value, timestamp, and window. */ + public static WindowedValue of( + 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); + return new TimestampedValueInGlobalWindow<>( + value, timestamp, paneInfo, null, null, causedByDrain); } else { - return new TimestampedValueInSingleWindow<>(value, timestamp, window, paneInfo, null, null); + return new TimestampedValueInSingleWindow<>( + value, timestamp, window, paneInfo, null, null, causedByDrain); } } @@ -289,7 +318,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, false); } /** @@ -297,7 +326,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, false); } /** @@ -308,7 +337,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, false); } } @@ -321,7 +351,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, false); } } @@ -337,7 +367,8 @@ public static WindowedValue withValue( windowedValue.getWindows(), windowedValue.getPaneInfo(), windowedValue.getRecordId(), - windowedValue.getRecordOffset()); + windowedValue.getRecordOffset(), + windowedValue.causedByDrain()); } public static boolean equals( @@ -388,6 +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 causedByDrain; @Override public @Nullable String getRecordId() { @@ -399,15 +431,22 @@ private abstract static class SimpleWindowedValue implements WindowedValue return currentRecordOffset; } + @Override + public boolean causedByDrain() { + return causedByDrain; + } + protected SimpleWindowedValue( T value, PaneInfo paneInfo, @Nullable String currentRecordId, - @Nullable Long currentRecordOffset) { + @Nullable Long currentRecordOffset, + boolean causedByDrain) { this.value = value; this.paneInfo = checkNotNull(paneInfo); this.currentRecordId = currentRecordId; this.currentRecordOffset = currentRecordOffset; + this.causedByDrain = causedByDrain; } @Override @@ -455,8 +494,9 @@ public MinTimestampWindowedValue( T value, PaneInfo pane, @Nullable String currentRecordId, - @Nullable Long currentRecordOffset) { - super(value, pane, currentRecordId, currentRecordOffset); + @Nullable Long currentRecordOffset, + boolean causedByDrain) { + super(value, pane, currentRecordId, currentRecordOffset, causedByDrain); } @Override @@ -473,8 +513,9 @@ public ValueInGlobalWindow( T value, PaneInfo paneInfo, @Nullable String currentRecordId, - @Nullable Long currentRecordOffset) { - super(value, paneInfo, currentRecordId, currentRecordOffset); + @Nullable Long currentRecordOffset, + boolean causedByDrain) { + super(value, paneInfo, currentRecordId, currentRecordOffset, causedByDrain); } @Override @@ -489,7 +530,8 @@ public BoundedWindow getWindow() { @Override public WindowedValue withValue(NewT newValue) { - return new ValueInGlobalWindow<>(newValue, getPaneInfo(), getRecordId(), getRecordOffset()); + return new ValueInGlobalWindow<>( + newValue, getPaneInfo(), getRecordId(), getRecordOffset(), causedByDrain()); } @Override @@ -513,6 +555,7 @@ public String toString() { return MoreObjects.toStringHelper(getClass()) .add("value", getValue()) .add("paneInfo", getPaneInfo()) + .add("causedByDrain", causedByDrain()) .toString(); } } @@ -526,8 +569,9 @@ public TimestampedWindowedValue( Instant timestamp, PaneInfo paneInfo, @Nullable String currentRecordId, - @Nullable Long currentRecordOffset) { - super(value, paneInfo, currentRecordId, currentRecordOffset); + @Nullable Long currentRecordOffset, + boolean causedByDrain) { + super(value, paneInfo, currentRecordId, currentRecordOffset, causedByDrain); this.timestamp = checkNotNull(timestamp); } @@ -549,8 +593,9 @@ public TimestampedValueInGlobalWindow( Instant timestamp, PaneInfo paneInfo, @Nullable String currentRecordId, - @Nullable Long currentRecordOffset) { - super(value, timestamp, paneInfo, currentRecordId, currentRecordOffset); + @Nullable Long currentRecordOffset, + boolean causedByDrain) { + super(value, timestamp, paneInfo, currentRecordId, currentRecordOffset, causedByDrain); } @Override @@ -566,7 +611,12 @@ public BoundedWindow getWindow() { @Override public WindowedValue withValue(NewT newValue) { return new TimestampedValueInGlobalWindow<>( - newValue, getTimestamp(), getPaneInfo(), getRecordId(), getRecordOffset()); + newValue, + getTimestamp(), + getPaneInfo(), + getRecordId(), + getRecordOffset(), + causedByDrain()); } @Override @@ -596,6 +646,7 @@ public String toString() { .add("value", getValue()) .add("timestamp", getTimestamp()) .add("paneInfo", getPaneInfo()) + .add("causedByDrain", causedByDrain()) .toString(); } } @@ -615,15 +666,22 @@ public TimestampedValueInSingleWindow( BoundedWindow window, PaneInfo paneInfo, @Nullable String currentRecordId, - @Nullable Long currentRecordOffset) { - super(value, timestamp, paneInfo, currentRecordId, currentRecordOffset); + @Nullable Long currentRecordOffset, + boolean causedByDrain) { + super(value, timestamp, paneInfo, currentRecordId, currentRecordOffset, causedByDrain); 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(), + causedByDrain()); } @Override @@ -665,6 +723,7 @@ public String toString() { .add("timestamp", getTimestamp()) .add("window", window) .add("paneInfo", getPaneInfo()) + .add("causedByDrain", causedByDrain()) .toString(); } } @@ -679,8 +738,9 @@ public TimestampedValueInMultipleWindows( Collection windows, PaneInfo paneInfo, @Nullable String currentRecordId, - @Nullable Long currentRecordOffset) { - super(value, timestamp, paneInfo, currentRecordId, currentRecordOffset); + @Nullable Long currentRecordOffset, + boolean causedByDrain) { + super(value, timestamp, paneInfo, currentRecordId, currentRecordOffset, causedByDrain); this.windows = checkNotNull(windows); } @@ -692,7 +752,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(), + causedByDrain()); } @Override @@ -730,6 +796,7 @@ public String toString() { .add("timestamp", getTimestamp()) .add("windows", windows) .add("paneInfo", getPaneInfo()) + .add("causedByDrain", causedByDrain()) .toString(); } @@ -845,7 +912,14 @@ 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.causedByDrain() + ? BeamFnApi.Elements.DrainMode.Enum.DRAINING + : BeamFnApi.Elements.DrainMode.Enum.NOT_DRAINING) + .build(); + ByteArrayCoder.of().encode(em.toByteArray(), outStream); } valueCoder.encode(windowedElem.getValue(), outStream, context); @@ -857,20 +931,27 @@ 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 causedByDrain = false; 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(); + causedByDrain = + b + ? elementMetadata.getDrain().equals(BeamFnApi.Elements.DrainMode.Enum.DRAINING) + : false; } 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, 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 3e3973e3720b..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 @@ -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); // drain is persisted as part of metadata 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.causedByDrain()); } @Test