Skip to content

Commit a805235

Browse files
committed
WIP: add OutputBuilder
1 parent 5b862dd commit a805235

8 files changed

Lines changed: 320 additions & 134 deletions

File tree

buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1182,6 +1182,7 @@ class BeamModulePlugin implements Plugin<Project> {
11821182

11831183
List<String> skipDefRegexes = []
11841184
skipDefRegexes << "AutoValue_.*"
1185+
skipDefRegexes << "AutoBuilder_.*"
11851186
skipDefRegexes << "AutoOneOf_.*"
11861187
skipDefRegexes << ".*\\.jmh_generated\\..*"
11871188
skipDefRegexes += configuration.generatedClassPatterns
@@ -1275,7 +1276,8 @@ class BeamModulePlugin implements Plugin<Project> {
12751276
'**/org/apache/beam/gradle/**',
12761277
'**/org/apache/beam/model/**',
12771278
'**/org/apache/beam/runners/dataflow/worker/windmill/**',
1278-
'**/AutoValue_*'
1279+
'**/AutoValue_*',
1280+
'**/AutoBuilder_*',
12791281
]
12801282

12811283
def jacocoEnabled = project.hasProperty('enableJacocoReport')

sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/DoFn.java

Lines changed: 16 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -45,6 +45,7 @@
4545
import org.apache.beam.sdk.transforms.windowing.BoundedWindow;
4646
import org.apache.beam.sdk.transforms.windowing.PaneInfo;
4747
import org.apache.beam.sdk.transforms.windowing.Window;
48+
import org.apache.beam.sdk.values.OutputBuilder;
4849
import org.apache.beam.sdk.values.PCollection;
4950
import org.apache.beam.sdk.values.PCollectionView;
5051
import org.apache.beam.sdk.values.Row;
@@ -391,17 +392,27 @@ public TypeDescriptor<OutputT> getOutputTypeDescriptor() {
391392

392393
/** Receives values of the given type. */
393394
public interface OutputReceiver<T> {
394-
void output(T output);
395+
OutputBuilder<T> builder();
395396

396-
void outputWithTimestamp(T output, Instant timestamp);
397+
default void output(T value) {
398+
builder().setValue(value).output();
399+
}
400+
401+
default void outputWithTimestamp(T value, Instant timestamp) {
402+
builder().setValue(value).setTimestamp(timestamp).output();
403+
}
397404

398405
default void outputWindowedValue(
399-
T output,
406+
T value,
400407
Instant timestamp,
401408
Collection<? extends BoundedWindow> windows,
402409
PaneInfo paneInfo) {
403-
throw new UnsupportedOperationException(
404-
String.format("Not implemented: %s.outputWindowedValue", this.getClass().getName()));
410+
builder()
411+
.setValue(value)
412+
.setTimestamp(timestamp)
413+
.setWindows(windows)
414+
.setPaneInfo(paneInfo)
415+
.output();
405416
}
406417
}
407418

sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/DoFnOutputReceivers.java

Lines changed: 125 additions & 43 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@
1717
*/
1818
package org.apache.beam.sdk.transforms;
1919

20+
import static org.apache.beam.sdk.util.Preconditions.checkStateNotNull;
2021
import static org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkNotNull;
2122
import static org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkState;
2223

@@ -29,6 +30,8 @@
2930
import org.apache.beam.sdk.transforms.DoFn.OutputReceiver;
3031
import org.apache.beam.sdk.transforms.windowing.BoundedWindow;
3132
import org.apache.beam.sdk.transforms.windowing.PaneInfo;
33+
import org.apache.beam.sdk.util.WindowedValue;
34+
import org.apache.beam.sdk.values.OutputBuilder;
3235
import org.apache.beam.sdk.values.Row;
3336
import org.apache.beam.sdk.values.TupleTag;
3437
import org.checkerframework.checker.nullness.qual.Nullable;
@@ -40,40 +43,127 @@
4043
"nullness" // TODO(https://github.com/apache/beam/issues/20497)
4144
})
4245
public class DoFnOutputReceivers {
43-
private static class RowOutputReceiver<T> implements OutputReceiver<Row> {
44-
WindowedContextOutputReceiver<T> outputReceiver;
45-
SchemaCoder<T> schemaCoder;
4646

47-
public RowOutputReceiver(
48-
DoFn<?, ?>.WindowedContext context,
49-
@Nullable TupleTag<T> outputTag,
50-
SchemaCoder<T> schemaCoder) {
51-
outputReceiver = new WindowedContextOutputReceiver<>(context, outputTag);
52-
this.schemaCoder = checkNotNull(schemaCoder);
47+
/**
48+
* OutputBuilder implementations may extend this class to simplify wrapping of another
49+
* OutputBuilder.
50+
*
51+
* <p>Because a common cause of delegation is to adjust the type of output received, extenders
52+
* <i>must</i> implement {#link setValue} and {#link getValue}.
53+
*/
54+
public abstract static class DelegatingOutputBuilder<InnerT, OuterT>
55+
implements OutputBuilder<OuterT> {
56+
57+
private final OutputBuilder<InnerT> delegate;
58+
59+
public DelegatingOutputBuilder(OutputBuilder<InnerT> delegate) {
60+
this.delegate = delegate;
5361
}
5462

63+
protected OutputBuilder<InnerT> getDelegate() {
64+
return delegate;
65+
}
66+
67+
@Override
68+
public abstract OutputBuilder<OuterT> setValue(OuterT value);
69+
5570
@Override
56-
public void output(Row output) {
57-
outputReceiver.output(schemaCoder.getFromRowFunction().apply(output));
71+
public OutputBuilder<OuterT> setTimestamp(Instant timestamp) {
72+
delegate.setTimestamp(timestamp);
73+
return this;
5874
}
5975

6076
@Override
61-
public void outputWithTimestamp(Row output, Instant timestamp) {
62-
outputReceiver.outputWithTimestamp(schemaCoder.getFromRowFunction().apply(output), timestamp);
77+
public OutputBuilder<OuterT> setWindow(BoundedWindow window) {
78+
delegate.setWindow(window);
79+
return this;
6380
}
6481

6582
@Override
66-
public void outputWindowedValue(
67-
Row output,
68-
Instant timestamp,
69-
Collection<? extends BoundedWindow> windows,
70-
PaneInfo paneInfo) {
71-
outputReceiver.outputWindowedValue(
72-
schemaCoder.getFromRowFunction().apply(output), timestamp, windows, paneInfo);
83+
public OutputBuilder<OuterT> setWindows(Collection<? extends BoundedWindow> windows) {
84+
delegate.setWindows(windows);
85+
return this;
86+
}
87+
88+
@Override
89+
public OutputBuilder<OuterT> setPaneInfo(PaneInfo paneInfo) {
90+
delegate.setPaneInfo(paneInfo);
91+
return this;
92+
}
93+
94+
@Override
95+
public abstract OuterT getValue();
96+
97+
@Override
98+
public Instant getTimestamp() {
99+
return delegate.getTimestamp();
100+
}
101+
102+
@Override
103+
public Collection<? extends BoundedWindow> getWindows() {
104+
return delegate.getWindows();
105+
}
106+
107+
@Override
108+
public PaneInfo getPaneInfo() {
109+
return delegate.getPaneInfo();
110+
}
111+
112+
@Override
113+
public void output() {
114+
delegate.output();
73115
}
74116
}
75117

76-
private static class WindowedContextOutputReceiver<T> implements OutputReceiver<T> {
118+
private static class RowOutputBuilder<T> extends DelegatingOutputBuilder<T, Row>
119+
implements OutputBuilder<Row> {
120+
121+
private @Nullable Row currentRow;
122+
private final SchemaCoder<T> schemaCoder;
123+
124+
RowOutputBuilder(OutputBuilder<T> delegate, SchemaCoder<T> schemaCoder) {
125+
super(delegate);
126+
this.schemaCoder = schemaCoder;
127+
}
128+
129+
@Override
130+
public OutputBuilder<Row> setValue(Row value) {
131+
currentRow = value;
132+
getDelegate().setValue(schemaCoder.getFromRowFunction().apply(value));
133+
return (OutputBuilder<Row>) this;
134+
}
135+
136+
@Override
137+
public Row getValue() {
138+
checkStateNotNull(currentRow, "getValue() called before value set");
139+
return currentRow;
140+
}
141+
}
142+
143+
private static class RowOutputReceiver<T> implements OutputReceiver<Row> {
144+
WindowedContextOutputReceiver<T> outputReceiver;
145+
SchemaCoder<T> schemaCoder;
146+
147+
@Override
148+
public OutputBuilder<Row> builder() {
149+
return new RowOutputBuilder<>(outputReceiver.builder(), schemaCoder);
150+
}
151+
152+
public RowOutputReceiver(
153+
DoFn<?, ?>.WindowedContext context,
154+
@Nullable TupleTag<T> outputTag,
155+
SchemaCoder<T> schemaCoder) {
156+
outputReceiver = new WindowedContextOutputReceiver<>(context, outputTag);
157+
this.schemaCoder = checkNotNull(schemaCoder);
158+
}
159+
}
160+
161+
/**
162+
* OutputReceiver that delegates all its core functionality to DoFn.WindowedContext which predates
163+
* OutputReceiver and has most of the same methods.
164+
*/
165+
private static class WindowedContextOutputReceiver<T>
166+
implements OutputReceiver<T>, WindowedValue.Outputter<T> {
77167
DoFn<?, ?>.WindowedContext context;
78168
@Nullable TupleTag<T> outputTag;
79169

@@ -84,34 +174,26 @@ public WindowedContextOutputReceiver(
84174
}
85175

86176
@Override
87-
public void output(T output) {
88-
if (outputTag != null) {
89-
context.output(outputTag, output);
90-
} else {
91-
((DoFn<?, T>.WindowedContext) context).output(output);
92-
}
93-
}
94-
95-
@Override
96-
public void outputWithTimestamp(T output, Instant timestamp) {
97-
if (outputTag != null) {
98-
context.outputWithTimestamp(outputTag, output, timestamp);
99-
} else {
100-
((DoFn<?, T>.WindowedContext) context).outputWithTimestamp(output, timestamp);
101-
}
177+
public OutputBuilder<T> builder() {
178+
return WindowedValue.builder(this);
102179
}
103180

104181
@Override
105-
public void outputWindowedValue(
106-
T output,
107-
Instant timestamp,
108-
Collection<? extends BoundedWindow> windows,
109-
PaneInfo paneInfo) {
182+
public void output(OutputBuilder<T> outputBuilder) {
110183
if (outputTag != null) {
111-
context.outputWindowedValue(outputTag, output, timestamp, windows, paneInfo);
184+
context.outputWindowedValue(
185+
outputTag,
186+
outputBuilder.getValue(),
187+
outputBuilder.getTimestamp(),
188+
outputBuilder.getWindows(),
189+
outputBuilder.getPaneInfo());
112190
} else {
113191
((DoFn<?, T>.WindowedContext) context)
114-
.outputWindowedValue(output, timestamp, windows, paneInfo);
192+
.outputWindowedValue(
193+
outputBuilder.getValue(),
194+
outputBuilder.getTimestamp(),
195+
outputBuilder.getWindows(),
196+
outputBuilder.getPaneInfo());
115197
}
116198
}
117199
}

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

Lines changed: 84 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -17,9 +17,11 @@
1717
*/
1818
package org.apache.beam.sdk.util;
1919

20+
import static org.apache.beam.sdk.util.Preconditions.checkStateNotNull;
2021
import static org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkArgument;
2122
import static org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkNotNull;
2223

24+
import com.google.auto.value.AutoBuilder;
2325
import java.io.ByteArrayInputStream;
2426
import java.io.ByteArrayOutputStream;
2527
import java.io.IOException;
@@ -46,8 +48,11 @@
4648
import org.apache.beam.sdk.transforms.windowing.PaneInfo;
4749
import org.apache.beam.sdk.transforms.windowing.PaneInfo.PaneInfoCoder;
4850
import org.apache.beam.sdk.util.common.ElementByteSizeObserver;
51+
import org.apache.beam.sdk.values.OutputBuilder;
52+
import org.apache.beam.vendor.grpc.v1p69p0.com.google.common.collect.Iterables;
4953
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.MoreObjects;
5054
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList;
55+
import org.checkerframework.checker.nullness.qual.MonotonicNonNull;
5156
import org.checkerframework.checker.nullness.qual.Nullable;
5257
import org.joda.time.Instant;
5358

@@ -59,16 +64,91 @@
5964
@Internal
6065
public abstract class WindowedValue<T> {
6166

67+
public interface Outputter<T> {
68+
void output(OutputBuilder<T> outputBuilder);
69+
}
70+
71+
/** Create a typical {@link OutputBuilder} that */
72+
public static <T> OutputBuilder<T> builder(Outputter<T> outputter) {
73+
return new AutoBuilder_WindowedValue_Builder<T>().setOutputter(outputter);
74+
}
75+
76+
@AutoBuilder(callMethod = "of", ofClass = WindowedValue.class)
77+
abstract static class Builder<T> implements OutputBuilder<T> {
78+
79+
private @MonotonicNonNull Outputter<T> outputter;
80+
81+
@Override
82+
public abstract OutputBuilder<T> setValue(T value);
83+
84+
@Override
85+
public abstract OutputBuilder<T> setTimestamp(Instant timestamp);
86+
87+
@Override
88+
public abstract OutputBuilder<T> setWindows(Collection<? extends BoundedWindow> windows);
89+
90+
@Override
91+
public abstract OutputBuilder<T> setPaneInfo(PaneInfo pane);
92+
93+
@Override
94+
public OutputBuilder<T> setWindow(BoundedWindow window) {
95+
return setWindows(Collections.singleton(window));
96+
}
97+
98+
public Builder<T> setOutputter(Outputter<T> outputter) {
99+
this.outputter = outputter;
100+
return this;
101+
}
102+
103+
abstract WindowedValue<T> autoBuild();
104+
105+
@Override
106+
public void output() {
107+
checkStateNotNull(outputter, "Cannot output until Outputter set").output(this);
108+
}
109+
110+
@Override
111+
public abstract T getValue();
112+
113+
@Override
114+
public abstract Instant getTimestamp();
115+
116+
@Override
117+
public abstract Collection<? extends BoundedWindow> getWindows();
118+
119+
@Override
120+
public abstract PaneInfo getPaneInfo();
121+
122+
public WindowedValue<T> build() {
123+
if (getWindows().size() > 1) {
124+
return autoBuild();
125+
}
126+
127+
BoundedWindow window = Iterables.getOnlyElement(getWindows());
128+
129+
if (!GlobalWindow.INSTANCE.equals(window)) {
130+
return new TimestampedValueInSingleWindow<>(
131+
getValue(), getTimestamp(), window, getPaneInfo());
132+
}
133+
134+
if (BoundedWindow.TIMESTAMP_MIN_VALUE.equals(getTimestamp())) {
135+
return new ValueInGlobalWindow<>(getValue(), getPaneInfo());
136+
}
137+
138+
return new TimestampedValueInGlobalWindow<>(getValue(), getTimestamp(), getPaneInfo());
139+
};
140+
}
141+
62142
/** Returns a {@code WindowedValue} with the given value, timestamp, and windows. */
63143
public static <T> WindowedValue<T> of(
64-
T value, Instant timestamp, Collection<? extends BoundedWindow> windows, PaneInfo pane) {
65-
checkArgument(pane != null, "WindowedValue requires PaneInfo, but it was null");
144+
T value, Instant timestamp, Collection<? extends BoundedWindow> windows, PaneInfo paneInfo) {
145+
checkArgument(paneInfo != null, "WindowedValue requires PaneInfo, but it was null");
66146
checkArgument(windows.size() > 0, "WindowedValue requires windows, but there were none");
67147

68148
if (windows.size() == 1) {
69-
return of(value, timestamp, windows.iterator().next(), pane);
149+
return of(value, timestamp, windows.iterator().next(), paneInfo);
70150
} else {
71-
return new TimestampedValueInMultipleWindows<>(value, timestamp, windows, pane);
151+
return new TimestampedValueInMultipleWindows<>(value, timestamp, windows, paneInfo);
72152
}
73153
}
74154

0 commit comments

Comments
 (0)