Skip to content

Commit a6d2b7d

Browse files
committed
add draining to output builder, encode draining
1 parent 85853a3 commit a6d2b7d

3 files changed

Lines changed: 121 additions & 36 deletions

File tree

sdks/java/core/src/main/java/org/apache/beam/sdk/values/OutputBuilder.java

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -48,5 +48,7 @@ public interface OutputBuilder<T> extends WindowedValue<T> {
4848

4949
OutputBuilder<T> setRecordOffset(@Nullable Long recordOffset);
5050

51+
OutputBuilder<T> setDraining(@Nullable Boolean drain);
52+
5153
void output();
5254
}

sdks/java/core/src/main/java/org/apache/beam/sdk/values/WindowedValue.java

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -52,6 +52,9 @@ public interface WindowedValue<T> {
5252
@Nullable
5353
Long getRecordOffset();
5454

55+
@Nullable
56+
Boolean isDraining();
57+
5558
/**
5659
* A representation of each of the actual values represented by this compressed {@link
5760
* WindowedValue}, one per window.

0 commit comments

Comments
 (0)