Skip to content

Commit 57a6f14

Browse files
committed
Introduce WindowedValue receivers and consolidate runner code to them
We need a receiver for WindowedValue to create an OutputBuilder. This change introduces the interface and uses it in many places where it is appropriate, replacing and simplifying internal runner code.
1 parent deff583 commit 57a6f14

63 files changed

Lines changed: 407 additions & 716 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

runners/core-java/src/main/java/org/apache/beam/runners/core/DoFnRunners.java

Lines changed: 3 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,7 @@
2929
import org.apache.beam.sdk.transforms.DoFnSchemaInformation;
3030
import org.apache.beam.sdk.transforms.reflect.DoFnSignatures;
3131
import org.apache.beam.sdk.transforms.windowing.BoundedWindow;
32+
import org.apache.beam.sdk.util.WindowedValueMultiReceiver;
3233
import org.apache.beam.sdk.values.KV;
3334
import org.apache.beam.sdk.values.PCollectionView;
3435
import org.apache.beam.sdk.values.TupleTag;
@@ -41,12 +42,6 @@
4142
"nullness" // TODO(https://github.com/apache/beam/issues/20497)
4243
})
4344
public class DoFnRunners {
44-
/** Information about how to create output receivers and output to them. */
45-
public interface OutputManager {
46-
/** Outputs a single element to the receiver indicated by the given {@link TupleTag}. */
47-
<T> void output(TupleTag<T> tag, WindowedValue<T> output);
48-
}
49-
5045
/**
5146
* Returns an implementation of {@link DoFnRunner} that for a {@link DoFn}.
5247
*
@@ -58,7 +53,7 @@ public static <InputT, OutputT> DoFnRunner<InputT, OutputT> simpleRunner(
5853
PipelineOptions options,
5954
DoFn<InputT, OutputT> fn,
6055
SideInputReader sideInputReader,
61-
OutputManager outputManager,
56+
WindowedValueMultiReceiver outputManager,
6257
TupleTag<OutputT> mainOutputTag,
6358
List<TupleTag<?>> additionalOutputTags,
6459
StepContext stepContext,
@@ -168,7 +163,7 @@ ProcessFnRunner<InputT, OutputT, RestrictionT> newProcessFnRunner(
168163
PipelineOptions options,
169164
Collection<PCollectionView<?>> views,
170165
ReadyCheckingSideInputReader sideInputReader,
171-
OutputManager outputManager,
166+
WindowedValueMultiReceiver outputManager,
172167
TupleTag<OutputT> mainOutputTag,
173168
List<TupleTag<?>> additionalOutputTags,
174169
StepContext stepContext,

runners/core-java/src/main/java/org/apache/beam/runners/core/GroupAlsoByWindowViaWindowSetNewDoFn.java

Lines changed: 5 additions & 31 deletions
Original file line numberDiff line numberDiff line change
@@ -17,20 +17,17 @@
1717
*/
1818
package org.apache.beam.runners.core;
1919

20-
import java.util.Collection;
2120
import org.apache.beam.model.pipeline.v1.RunnerApi;
2221
import org.apache.beam.runners.core.triggers.ExecutableTriggerStateMachine;
2322
import org.apache.beam.runners.core.triggers.TriggerStateMachines;
2423
import org.apache.beam.sdk.transforms.DoFn;
2524
import org.apache.beam.sdk.transforms.windowing.BoundedWindow;
26-
import org.apache.beam.sdk.transforms.windowing.PaneInfo;
2725
import org.apache.beam.sdk.util.SystemDoFnInternal;
26+
import org.apache.beam.sdk.util.WindowedValueMultiReceiver;
2827
import org.apache.beam.sdk.util.construction.TriggerTranslation;
2928
import org.apache.beam.sdk.values.KV;
3029
import org.apache.beam.sdk.values.TupleTag;
31-
import org.apache.beam.sdk.values.WindowedValues;
3230
import org.apache.beam.sdk.values.WindowingStrategy;
33-
import org.joda.time.Instant;
3431

3532
/**
3633
* A general {@link GroupAlsoByWindowsAggregators}. This delegates all of the logic to the {@link
@@ -51,7 +48,7 @@ DoFn<KeyedWorkItem<K, InputT>, KV<K, OutputT>> create(
5148
TimerInternalsFactory<K> timerInternalsFactory,
5249
SideInputReader sideInputReader,
5350
SystemReduceFn<K, InputT, ?, OutputT, W> reduceFn,
54-
DoFnRunners.OutputManager outputManager,
51+
WindowedValueMultiReceiver outputManager,
5552
TupleTag<KV<K, OutputT>> mainTag) {
5653
return new GroupAlsoByWindowViaWindowSetNewDoFn<>(
5754
strategy,
@@ -68,7 +65,7 @@ DoFn<KeyedWorkItem<K, InputT>, KV<K, OutputT>> create(
6865
private transient StateInternalsFactory<K> stateInternalsFactory;
6966
private transient TimerInternalsFactory<K> timerInternalsFactory;
7067
private transient SideInputReader sideInputReader;
71-
private transient DoFnRunners.OutputManager outputManager;
68+
private transient WindowedValueMultiReceiver outputManager;
7269
private TupleTag<KV<K, OutputT>> mainTag;
7370

7471
public GroupAlsoByWindowViaWindowSetNewDoFn(
@@ -77,7 +74,7 @@ public GroupAlsoByWindowViaWindowSetNewDoFn(
7774
TimerInternalsFactory<K> timerInternalsFactory,
7875
SideInputReader sideInputReader,
7976
SystemReduceFn<K, InputT, ?, OutputT, W> reduceFn,
80-
DoFnRunners.OutputManager outputManager,
77+
WindowedValueMultiReceiver outputManager,
8178
TupleTag<KV<K, OutputT>> mainTag) {
8279
this.timerInternalsFactory = timerInternalsFactory;
8380
this.sideInputReader = sideInputReader;
@@ -91,29 +88,6 @@ public GroupAlsoByWindowViaWindowSetNewDoFn(
9188
this.triggerProto = TriggerTranslation.toProto(windowingStrategy.getTrigger());
9289
}
9390

94-
private OutputWindowedValue<KV<K, OutputT>> outputWindowedValue() {
95-
return new OutputWindowedValue<KV<K, OutputT>>() {
96-
@Override
97-
public void outputWindowedValue(
98-
KV<K, OutputT> output,
99-
Instant timestamp,
100-
Collection<? extends BoundedWindow> windows,
101-
PaneInfo pane) {
102-
outputManager.output(mainTag, WindowedValues.of(output, timestamp, windows, pane));
103-
}
104-
105-
@Override
106-
public <AdditionalOutputT> void outputWindowedValue(
107-
TupleTag<AdditionalOutputT> tag,
108-
AdditionalOutputT output,
109-
Instant timestamp,
110-
Collection<? extends BoundedWindow> windows,
111-
PaneInfo pane) {
112-
outputManager.output(tag, WindowedValues.of(output, timestamp, windows, pane));
113-
}
114-
};
115-
}
116-
11791
@ProcessElement
11892
public void processElement(ProcessContext c) throws Exception {
11993
KeyedWorkItem<K, InputT> keyedWorkItem = c.element();
@@ -130,7 +104,7 @@ public void processElement(ProcessContext c) throws Exception {
130104
TriggerStateMachines.stateMachineForTrigger(triggerProto)),
131105
stateInternals,
132106
timerInternals,
133-
outputWindowedValue(),
107+
windowedValue -> outputManager.output(mainTag, windowedValue),
134108
sideInputReader,
135109
reduceFn,
136110
c.getPipelineOptions());

runners/core-java/src/main/java/org/apache/beam/runners/core/OutputAndTimeBoundedSplittableProcessElementInvoker.java

Lines changed: 13 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -45,11 +45,13 @@
4545
import org.apache.beam.sdk.transforms.splittabledofn.WatermarkEstimator;
4646
import org.apache.beam.sdk.transforms.windowing.BoundedWindow;
4747
import org.apache.beam.sdk.transforms.windowing.PaneInfo;
48+
import org.apache.beam.sdk.util.WindowedValueMultiReceiver;
4849
import org.apache.beam.sdk.values.KV;
4950
import org.apache.beam.sdk.values.PCollectionView;
5051
import org.apache.beam.sdk.values.Row;
5152
import org.apache.beam.sdk.values.TupleTag;
5253
import org.apache.beam.sdk.values.WindowedValue;
54+
import org.apache.beam.sdk.values.WindowedValues;
5355
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Iterables;
5456
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.util.concurrent.Futures;
5557
import org.checkerframework.checker.nullness.qual.Nullable;
@@ -72,19 +74,20 @@ public class OutputAndTimeBoundedSplittableProcessElementInvoker<
7274
InputT, OutputT, RestrictionT, PositionT, WatermarkEstimatorStateT> {
7375
private final DoFn<InputT, OutputT> fn;
7476
private final PipelineOptions pipelineOptions;
75-
private final OutputWindowedValue<OutputT> output;
77+
private final WindowedValueMultiReceiver outputReceiver;
7678
private final SideInputReader sideInputReader;
7779
private final ScheduledExecutorService executor;
7880
private final int maxNumOutputs;
7981
private final Duration maxDuration;
8082
private final Supplier<BundleFinalizer> bundleFinalizer;
83+
private final TupleTag<OutputT> mainOutputTag;
8184

8285
/**
8386
* Creates a new invoker from components.
8487
*
8588
* @param fn The original {@link DoFn}.
8689
* @param pipelineOptions {@link PipelineOptions} to include in the {@link DoFn.ProcessContext}.
87-
* @param output Hook for outputting from the {@link DoFn.ProcessElement} method.
90+
* @param outputReceiver Hook for outputting from the {@link DoFn.ProcessElement} method.
8891
* @param sideInputReader Hook for accessing side inputs.
8992
* @param executor Executor on which a checkpoint will be scheduled after the given duration.
9093
* @param maxNumOutputs Maximum number of outputs, in total over all output tags, after which a
@@ -98,15 +101,17 @@ public class OutputAndTimeBoundedSplittableProcessElementInvoker<
98101
public OutputAndTimeBoundedSplittableProcessElementInvoker(
99102
DoFn<InputT, OutputT> fn,
100103
PipelineOptions pipelineOptions,
101-
OutputWindowedValue<OutputT> output,
104+
WindowedValueMultiReceiver outputReceiver,
105+
TupleTag<OutputT> mainOutputTag,
102106
SideInputReader sideInputReader,
103107
ScheduledExecutorService executor,
104108
int maxNumOutputs,
105109
Duration maxDuration,
106110
Supplier<BundleFinalizer> bundleFinalizer) {
107111
this.fn = fn;
108112
this.pipelineOptions = pipelineOptions;
109-
this.output = output;
113+
this.outputReceiver = outputReceiver;
114+
this.mainOutputTag = mainOutputTag;
110115
this.sideInputReader = sideInputReader;
111116
this.executor = executor;
112117
this.maxNumOutputs = maxNumOutputs;
@@ -403,7 +408,7 @@ public void outputWindowedValue(
403408
if (watermarkEstimator instanceof TimestampObservingWatermarkEstimator) {
404409
((TimestampObservingWatermarkEstimator) watermarkEstimator).observeTimestamp(timestamp);
405410
}
406-
output.outputWindowedValue(value, timestamp, windows, paneInfo);
411+
outputReceiver.output(mainOutputTag, WindowedValues.of(value, timestamp, windows, paneInfo));
407412
}
408413

409414
@Override
@@ -413,7 +418,8 @@ public <T> void output(TupleTag<T> tag, T value) {
413418

414419
@Override
415420
public <T> void outputWithTimestamp(TupleTag<T> tag, T value, Instant timestamp) {
416-
outputWindowedValue(tag, value, timestamp, element.getWindows(), element.getPane());
421+
outputReceiver.output(
422+
tag, WindowedValues.of(value, timestamp, element.getWindows(), element.getPane()));
417423
}
418424

419425
@Override
@@ -427,7 +433,7 @@ public <T> void outputWindowedValue(
427433
if (watermarkEstimator instanceof TimestampObservingWatermarkEstimator) {
428434
((TimestampObservingWatermarkEstimator) watermarkEstimator).observeTimestamp(timestamp);
429435
}
430-
output.outputWindowedValue(tag, value, timestamp, windows, paneInfo);
436+
outputReceiver.output(tag, WindowedValues.of(value, timestamp, windows, paneInfo));
431437
}
432438

433439
private void noteOutput() {

runners/core-java/src/main/java/org/apache/beam/runners/core/OutputWindowedValue.java

Lines changed: 0 additions & 45 deletions
This file was deleted.

runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java

Lines changed: 6 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -47,9 +47,11 @@
4747
import org.apache.beam.sdk.transforms.windowing.Window.ClosingBehavior;
4848
import org.apache.beam.sdk.transforms.windowing.WindowFn;
4949
import org.apache.beam.sdk.util.WindowTracing;
50+
import org.apache.beam.sdk.util.WindowedValueReceiver;
5051
import org.apache.beam.sdk.values.KV;
5152
import org.apache.beam.sdk.values.PCollection;
5253
import org.apache.beam.sdk.values.WindowedValue;
54+
import org.apache.beam.sdk.values.WindowedValues;
5355
import org.apache.beam.sdk.values.WindowingStrategy;
5456
import org.apache.beam.sdk.values.WindowingStrategy.AccumulationMode;
5557
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting;
@@ -106,7 +108,7 @@ public class ReduceFnRunner<K, InputT, OutputT, W extends BoundedWindow> {
106108
*/
107109
private final WindowingStrategy<Object, W> windowingStrategy;
108110

109-
private final OutputWindowedValue<KV<K, OutputT>> outputter;
111+
private final WindowedValueReceiver<KV<K, OutputT>> outputter;
110112

111113
private final StateInternals stateInternals;
112114

@@ -214,7 +216,7 @@ public ReduceFnRunner(
214216
ExecutableTriggerStateMachine triggerStateMachine,
215217
StateInternals stateInternals,
216218
TimerInternals timerInternals,
217-
OutputWindowedValue<KV<K, OutputT>> outputter,
219+
WindowedValueReceiver<KV<K, OutputT>> outputter,
218220
@Nullable SideInputReader sideInputReader,
219221
ReduceFn<K, InputT, OutputT, W> reduceFn,
220222
@Nullable PipelineOptions options) {
@@ -1055,7 +1057,8 @@ private void prefetchOnTrigger(
10551057
}
10561058

10571059
// Output the actual value.
1058-
outputter.outputWindowedValue(KV.of(key, toOutput), outputTimestamp, windows, pane);
1060+
outputter.output(
1061+
WindowedValues.of(KV.of(key, toOutput), outputTimestamp, windows, pane));
10591062
});
10601063

10611064
reduceFn.onTrigger(renamedTriggerContext);

runners/core-java/src/main/java/org/apache/beam/runners/core/SimpleDoFnRunner.java

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -26,7 +26,6 @@
2626
import java.util.List;
2727
import java.util.Map;
2828
import java.util.Set;
29-
import org.apache.beam.runners.core.DoFnRunners.OutputManager;
3029
import org.apache.beam.sdk.coders.Coder;
3130
import org.apache.beam.sdk.options.PipelineOptions;
3231
import org.apache.beam.sdk.schemas.SchemaCoder;
@@ -54,6 +53,7 @@
5453
import org.apache.beam.sdk.transforms.windowing.PaneInfo;
5554
import org.apache.beam.sdk.util.SystemDoFnInternal;
5655
import org.apache.beam.sdk.util.UserCodeException;
56+
import org.apache.beam.sdk.util.WindowedValueMultiReceiver;
5757
import org.apache.beam.sdk.values.PCollectionView;
5858
import org.apache.beam.sdk.values.Row;
5959
import org.apache.beam.sdk.values.TupleTag;
@@ -94,7 +94,7 @@ public class SimpleDoFnRunner<InputT, OutputT> implements DoFnRunner<InputT, Out
9494
private final DoFnInvoker<InputT, OutputT> invoker;
9595

9696
private final SideInputReader sideInputReader;
97-
private final OutputManager outputManager;
97+
private final WindowedValueMultiReceiver outputManager;
9898

9999
private final TupleTag<OutputT> mainOutputTag;
100100
/** The set of known output tags. */
@@ -124,7 +124,7 @@ public SimpleDoFnRunner(
124124
PipelineOptions options,
125125
DoFn<InputT, OutputT> fn,
126126
SideInputReader sideInputReader,
127-
OutputManager outputManager,
127+
WindowedValueMultiReceiver outputManager,
128128
TupleTag<OutputT> mainOutputTag,
129129
List<TupleTag<?>> additionalOutputTags,
130130
StepContext stepContext,

runners/core-java/src/test/java/org/apache/beam/runners/core/OutputAndTimeBoundedSplittableProcessElementInvokerTest.java

Lines changed: 7 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -27,7 +27,6 @@
2727
import static org.junit.Assert.assertNull;
2828
import static org.junit.Assert.assertTrue;
2929

30-
import java.util.Collection;
3130
import java.util.Collections;
3231
import java.util.concurrent.Executors;
3332
import java.util.concurrent.TimeUnit;
@@ -38,10 +37,11 @@
3837
import org.apache.beam.sdk.transforms.splittabledofn.OffsetRangeTracker;
3938
import org.apache.beam.sdk.transforms.splittabledofn.RestrictionTracker;
4039
import org.apache.beam.sdk.transforms.splittabledofn.WatermarkEstimator;
41-
import org.apache.beam.sdk.transforms.windowing.BoundedWindow;
4240
import org.apache.beam.sdk.transforms.windowing.GlobalWindow;
4341
import org.apache.beam.sdk.transforms.windowing.PaneInfo;
42+
import org.apache.beam.sdk.util.WindowedValueMultiReceiver;
4443
import org.apache.beam.sdk.values.TupleTag;
44+
import org.apache.beam.sdk.values.WindowedValue;
4545
import org.apache.beam.sdk.values.WindowedValues;
4646
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.util.concurrent.Uninterruptibles;
4747
import org.joda.time.Duration;
@@ -112,22 +112,13 @@ private SplittableProcessElementInvoker<Void, String, OffsetRange, Long, Void>.R
112112
new OutputAndTimeBoundedSplittableProcessElementInvoker<>(
113113
fn,
114114
PipelineOptionsFactory.create(),
115-
new OutputWindowedValue<String>() {
115+
new WindowedValueMultiReceiver() {
116116
@Override
117-
public void outputWindowedValue(
118-
String output,
119-
Instant timestamp,
120-
Collection<? extends BoundedWindow> windows,
121-
PaneInfo pane) {}
122-
123-
@Override
124-
public <AdditionalOutputT> void outputWindowedValue(
125-
TupleTag<AdditionalOutputT> tag,
126-
AdditionalOutputT output,
127-
Instant timestamp,
128-
Collection<? extends BoundedWindow> windows,
129-
PaneInfo pane) {}
117+
public <OutputT> void output(TupleTag<OutputT> tag, WindowedValue<OutputT> output) {
118+
// discard
119+
}
130120
},
121+
null,
131122
NullSideInputReader.empty(),
132123
Executors.newSingleThreadScheduledExecutor(),
133124
1000,

0 commit comments

Comments
 (0)