Skip to content

Commit fe6043f

Browse files
committed
Remove WindowedValue.getPane() in preparation for making WindowedValue a user-facing interface
1 parent 7995a5e commit fe6043f

47 files changed

Lines changed: 500 additions & 479 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/LateDataDroppingDoFnRunner.java

Lines changed: 3 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -25,9 +25,8 @@
2525
import org.apache.beam.sdk.transforms.DoFn;
2626
import org.apache.beam.sdk.transforms.windowing.BoundedWindow;
2727
import org.apache.beam.sdk.util.WindowTracing;
28+
import org.apache.beam.sdk.util.WindowedValue;
2829
import org.apache.beam.sdk.values.KV;
29-
import org.apache.beam.sdk.values.WindowedValue;
30-
import org.apache.beam.sdk.values.WindowedValues;
3130
import org.apache.beam.sdk.values.WindowingStrategy;
3231
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting;
3332
import org.joda.time.Instant;
@@ -141,8 +140,8 @@ public <K, InputT> Iterable<WindowedValue<InputT>> filter(
141140
timerInternals.currentOutputWatermarkTime());
142141
} else {
143142
nonLateElements.add(
144-
WindowedValues.of(
145-
element.getValue(), element.getTimestamp(), window, element.getPane()));
143+
WindowedValue.of(
144+
element.getValue(), element.getTimestamp(), window, element.getPaneInfo()));
146145
}
147146
}
148147
}

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

Lines changed: 10 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -45,13 +45,11 @@
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;
48+
import org.apache.beam.sdk.util.WindowedValue;
4949
import org.apache.beam.sdk.values.KV;
5050
import org.apache.beam.sdk.values.PCollectionView;
5151
import org.apache.beam.sdk.values.Row;
5252
import org.apache.beam.sdk.values.TupleTag;
53-
import org.apache.beam.sdk.values.WindowedValue;
54-
import org.apache.beam.sdk.values.WindowedValues;
5553
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Iterables;
5654
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.util.concurrent.Futures;
5755
import org.checkerframework.checker.nullness.qual.Nullable;
@@ -74,20 +72,19 @@ public class OutputAndTimeBoundedSplittableProcessElementInvoker<
7472
InputT, OutputT, RestrictionT, PositionT, WatermarkEstimatorStateT> {
7573
private final DoFn<InputT, OutputT> fn;
7674
private final PipelineOptions pipelineOptions;
77-
private final WindowedValueMultiReceiver outputReceiver;
75+
private final OutputWindowedValue<OutputT> output;
7876
private final SideInputReader sideInputReader;
7977
private final ScheduledExecutorService executor;
8078
private final int maxNumOutputs;
8179
private final Duration maxDuration;
8280
private final Supplier<BundleFinalizer> bundleFinalizer;
83-
private final TupleTag<OutputT> mainOutputTag;
8481

8582
/**
8683
* Creates a new invoker from components.
8784
*
8885
* @param fn The original {@link DoFn}.
8986
* @param pipelineOptions {@link PipelineOptions} to include in the {@link DoFn.ProcessContext}.
90-
* @param outputReceiver Hook for outputting from the {@link DoFn.ProcessElement} method.
87+
* @param output Hook for outputting from the {@link DoFn.ProcessElement} method.
9188
* @param sideInputReader Hook for accessing side inputs.
9289
* @param executor Executor on which a checkpoint will be scheduled after the given duration.
9390
* @param maxNumOutputs Maximum number of outputs, in total over all output tags, after which a
@@ -101,17 +98,15 @@ public class OutputAndTimeBoundedSplittableProcessElementInvoker<
10198
public OutputAndTimeBoundedSplittableProcessElementInvoker(
10299
DoFn<InputT, OutputT> fn,
103100
PipelineOptions pipelineOptions,
104-
WindowedValueMultiReceiver outputReceiver,
105-
TupleTag<OutputT> mainOutputTag,
101+
OutputWindowedValue<OutputT> output,
106102
SideInputReader sideInputReader,
107103
ScheduledExecutorService executor,
108104
int maxNumOutputs,
109105
Duration maxDuration,
110106
Supplier<BundleFinalizer> bundleFinalizer) {
111107
this.fn = fn;
112108
this.pipelineOptions = pipelineOptions;
113-
this.outputReceiver = outputReceiver;
114-
this.mainOutputTag = mainOutputTag;
109+
this.output = output;
115110
this.sideInputReader = sideInputReader;
116111
this.executor = executor;
117112
this.maxNumOutputs = maxNumOutputs;
@@ -380,7 +375,7 @@ public Instant timestamp() {
380375

381376
@Override
382377
public PaneInfo pane() {
383-
return element.getPane();
378+
return element.getPaneInfo();
384379
}
385380

386381
@Override
@@ -395,7 +390,7 @@ public void output(OutputT output) {
395390

396391
@Override
397392
public void outputWithTimestamp(OutputT value, Instant timestamp) {
398-
outputWindowedValue(value, timestamp, element.getWindows(), element.getPane());
393+
outputWindowedValue(value, timestamp, element.getWindows(), element.getPaneInfo());
399394
}
400395

401396
@Override
@@ -408,7 +403,7 @@ public void outputWindowedValue(
408403
if (watermarkEstimator instanceof TimestampObservingWatermarkEstimator) {
409404
((TimestampObservingWatermarkEstimator) watermarkEstimator).observeTimestamp(timestamp);
410405
}
411-
outputReceiver.output(mainOutputTag, WindowedValues.of(value, timestamp, windows, paneInfo));
406+
output.outputWindowedValue(value, timestamp, windows, paneInfo);
412407
}
413408

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

419414
@Override
420415
public <T> void outputWithTimestamp(TupleTag<T> tag, T value, Instant timestamp) {
421-
outputReceiver.output(
422-
tag, WindowedValues.of(value, timestamp, element.getWindows(), element.getPane()));
416+
outputWindowedValue(tag, value, timestamp, element.getWindows(), element.getPaneInfo());
423417
}
424418

425419
@Override
@@ -433,7 +427,7 @@ public <T> void outputWindowedValue(
433427
if (watermarkEstimator instanceof TimestampObservingWatermarkEstimator) {
434428
((TimestampObservingWatermarkEstimator) watermarkEstimator).observeTimestamp(timestamp);
435429
}
436-
outputReceiver.output(tag, WindowedValues.of(value, timestamp, windows, paneInfo));
430+
output.outputWindowedValue(tag, value, timestamp, windows, paneInfo);
437431
}
438432

439433
private void noteOutput() {

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

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -405,7 +405,7 @@ public <T> T sideInput(PCollectionView<T> view) {
405405

406406
@Override
407407
public PaneInfo pane() {
408-
return elem.getPane();
408+
return elem.getPaneInfo();
409409
}
410410

411411
@Override
@@ -437,7 +437,7 @@ public <T> void output(TupleTag<T> tag, T output) {
437437
public <T> void outputWithTimestamp(TupleTag<T> tag, T output, Instant timestamp) {
438438
checkNotNull(tag, "Tag passed to outputWithTimestamp cannot be null");
439439
checkTimestamp(elem.getTimestamp(), timestamp);
440-
outputWindowedValue(tag, output, timestamp, elem.getWindows(), elem.getPane());
440+
outputWindowedValue(tag, output, timestamp, elem.getWindows(), elem.getPaneInfo());
441441
}
442442

443443
@Override

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

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -431,7 +431,7 @@ public PipelineOptions pipelineOptions() {
431431

432432
@Override
433433
public PaneInfo paneInfo(DoFn<InputT, OutputT> doFn) {
434-
return elementAndRestriction.getKey().getPane();
434+
return elementAndRestriction.getKey().getPaneInfo();
435435
}
436436

437437
@Override
@@ -491,7 +491,7 @@ public PipelineOptions pipelineOptions() {
491491

492492
@Override
493493
public PaneInfo paneInfo(DoFn<InputT, OutputT> doFn) {
494-
return elementAndRestriction.getKey().getPane();
494+
return elementAndRestriction.getKey().getPaneInfo();
495495
}
496496

497497
@Override
@@ -545,7 +545,7 @@ public PipelineOptions pipelineOptions() {
545545

546546
@Override
547547
public PaneInfo paneInfo(DoFn<InputT, OutputT> doFn) {
548-
return elementAndRestriction.getKey().getPane();
548+
return elementAndRestriction.getKey().getPaneInfo();
549549
}
550550

551551
@Override

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

Lines changed: 29 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -759,7 +759,8 @@ public void testWatermarkHoldAndLateData() throws Exception {
759759
equalTo(new Instant(1)),
760760
equalTo((BoundedWindow) expectedWindow))));
761761
assertThat(
762-
output.get(0).getPane(), equalTo(PaneInfo.createPane(true, false, Timing.EARLY, 0, -1)));
762+
output.get(0).getPaneInfo(),
763+
equalTo(PaneInfo.createPane(true, false, Timing.EARLY, 0, -1)));
763764

764765
// There is no end-of-window hold, but the timer set by the trigger holds the watermark
765766
assertThat(tester.getWatermarkHold(), nullValue());
@@ -805,7 +806,8 @@ public void testWatermarkHoldAndLateData() throws Exception {
805806
0, // window start
806807
10))); // window end
807808
assertThat(
808-
output.get(0).getPane(), equalTo(PaneInfo.createPane(false, false, Timing.EARLY, 1, -1)));
809+
output.get(0).getPaneInfo(),
810+
equalTo(PaneInfo.createPane(false, false, Timing.EARLY, 1, -1)));
809811

810812
// Since the element hold is cleared, there is no hold remaining
811813
assertThat(tester.getWatermarkHold(), nullValue());
@@ -846,7 +848,8 @@ public void testWatermarkHoldAndLateData() throws Exception {
846848
0, // window start
847849
10))); // window end
848850
assertThat(
849-
output.get(0).getPane(), equalTo(PaneInfo.createPane(false, false, Timing.ON_TIME, 2, 0)));
851+
output.get(0).getPaneInfo(),
852+
equalTo(PaneInfo.createPane(false, false, Timing.ON_TIME, 2, 0)));
850853

851854
tester.setAutoAdvanceOutputWatermark(true);
852855

@@ -879,7 +882,7 @@ public void testWatermarkHoldAndLateData() throws Exception {
879882
0, // window start
880883
10))); // window end
881884
assertThat(
882-
output.get(0).getPane(), equalTo(PaneInfo.createPane(false, true, Timing.LATE, 3, 1)));
885+
output.get(0).getPaneInfo(), equalTo(PaneInfo.createPane(false, true, Timing.LATE, 3, 1)));
883886
assertEquals(new Instant(50), tester.getOutputWatermark());
884887
assertEquals(null, tester.getWatermarkHold());
885888

@@ -1486,7 +1489,8 @@ public void testMergeBeforeFinalizing() throws Exception {
14861489
1, // window start
14871490
20)); // window end
14881491
assertThat(
1489-
output.get(0).getPane(), equalTo(PaneInfo.createPane(true, true, Timing.ON_TIME, 0, 0)));
1492+
output.get(0).getPaneInfo(),
1493+
equalTo(PaneInfo.createPane(true, true, Timing.ON_TIME, 0, 0)));
14901494
}
14911495

14921496
/**
@@ -1527,7 +1531,8 @@ public void testMergingWithCloseBeforeGC() throws Exception {
15271531
1, // window start
15281532
20)); // window end
15291533
assertThat(
1530-
output.get(0).getPane(), equalTo(PaneInfo.createPane(true, true, Timing.ON_TIME, 0, 0)));
1534+
output.get(0).getPaneInfo(),
1535+
equalTo(PaneInfo.createPane(true, true, Timing.ON_TIME, 0, 0)));
15311536
}
15321537

15331538
/** Ensure a closed trigger has its state recorded in the merge result window. */
@@ -1614,7 +1619,8 @@ public void testMergingWithReusedWindow() throws Exception {
16141619
equalTo((BoundedWindow) mergedWindow)));
16151620

16161621
assertThat(
1617-
output.get(0).getPane(), equalTo(PaneInfo.createPane(true, true, Timing.ON_TIME, 0, 0)));
1622+
output.get(0).getPaneInfo(),
1623+
equalTo(PaneInfo.createPane(true, true, Timing.ON_TIME, 0, 0)));
16181624
}
16191625

16201626
/**
@@ -1661,7 +1667,7 @@ public void testMergingWithClosedRepresentative() throws Exception {
16611667
1, // window start
16621668
18)); // window end
16631669
assertThat(
1664-
output.get(0).getPane(), equalTo(PaneInfo.createPane(true, true, Timing.EARLY, 0, 0)));
1670+
output.get(0).getPaneInfo(), equalTo(PaneInfo.createPane(true, true, Timing.EARLY, 0, 0)));
16651671
}
16661672

16671673
/**
@@ -1702,7 +1708,7 @@ public void testMergingWithClosedDoesNotPoison() throws Exception {
17021708
2, // window start
17031709
12)); // window end
17041710
assertThat(
1705-
output.get(0).getPane(), equalTo(PaneInfo.createPane(true, true, Timing.EARLY, 0, 0)));
1711+
output.get(0).getPaneInfo(), equalTo(PaneInfo.createPane(true, true, Timing.EARLY, 0, 0)));
17061712
assertThat(
17071713
output.get(1),
17081714
isSingleWindowedValue(
@@ -1711,7 +1717,8 @@ public void testMergingWithClosedDoesNotPoison() throws Exception {
17111717
1, // window start
17121718
13)); // window end
17131719
assertThat(
1714-
output.get(1).getPane(), equalTo(PaneInfo.createPane(true, true, Timing.ON_TIME, 0, 0)));
1720+
output.get(1).getPaneInfo(),
1721+
equalTo(PaneInfo.createPane(true, true, Timing.ON_TIME, 0, 0)));
17151722
}
17161723

17171724
/**
@@ -1811,7 +1818,7 @@ public void testIdempotentEmptyPanesDiscarding() throws Exception {
18111818
// The late pane has the correct indices.
18121819
assertThat(output.get(1).getValue(), contains(3));
18131820
assertThat(
1814-
output.get(1).getPane(), equalTo(PaneInfo.createPane(false, true, Timing.LATE, 1, 1)));
1821+
output.get(1).getPaneInfo(), equalTo(PaneInfo.createPane(false, true, Timing.LATE, 1, 1)));
18151822

18161823
assertTrue(tester.isMarkedFinished(firstWindow));
18171824
tester.assertHasOnlyGlobalAndFinishedSetsFor(firstWindow);
@@ -1850,7 +1857,8 @@ public void testIdempotentEmptyPanesAccumulating() throws Exception {
18501857
assertThat(output.size(), equalTo(1));
18511858
assertThat(output.get(0), isSingleWindowedValue(containsInAnyOrder(1, 2), 1, 0, 10));
18521859
assertThat(
1853-
output.get(0).getPane(), equalTo(PaneInfo.createPane(true, false, Timing.ON_TIME, 0, 0)));
1860+
output.get(0).getPaneInfo(),
1861+
equalTo(PaneInfo.createPane(true, false, Timing.ON_TIME, 0, 0)));
18541862

18551863
// Fire another timer with no data; the empty pane should not be output even though the
18561864
// trigger is ready to fire
@@ -1868,7 +1876,7 @@ public void testIdempotentEmptyPanesAccumulating() throws Exception {
18681876
// The late pane has the correct indices.
18691877
assertThat(output.get(0).getValue(), containsInAnyOrder(1, 2, 3));
18701878
assertThat(
1871-
output.get(0).getPane(), equalTo(PaneInfo.createPane(false, true, Timing.LATE, 1, 1)));
1879+
output.get(0).getPaneInfo(), equalTo(PaneInfo.createPane(false, true, Timing.LATE, 1, 1)));
18721880

18731881
assertTrue(tester.isMarkedFinished(firstWindow));
18741882
tester.assertHasOnlyGlobalAndFinishedSetsFor(firstWindow);
@@ -2193,17 +2201,17 @@ public void fireNonEmptyOnDrainInGlobalWindow() throws Exception {
21932201
List<WindowedValue<Iterable<Integer>>> output = tester.extractOutput();
21942202
assertEquals(n / 3, output.size());
21952203
for (int i = 0; i < output.size(); i++) {
2196-
assertEquals(Timing.EARLY, output.get(i).getPane().getTiming());
2197-
assertEquals(i, output.get(i).getPane().getIndex());
2204+
assertEquals(Timing.EARLY, output.get(i).getPaneInfo().getTiming());
2205+
assertEquals(i, output.get(i).getPaneInfo().getIndex());
21982206
assertEquals(3, Iterables.size(output.get(i).getValue()));
21992207
}
22002208

22012209
tester.advanceInputWatermark(BoundedWindow.TIMESTAMP_MAX_VALUE);
22022210

22032211
output = tester.extractOutput();
22042212
assertEquals(1, output.size());
2205-
assertEquals(Timing.ON_TIME, output.get(0).getPane().getTiming());
2206-
assertEquals(n / 3, output.get(0).getPane().getIndex());
2213+
assertEquals(Timing.ON_TIME, output.get(0).getPaneInfo().getTiming());
2214+
assertEquals(n / 3, output.get(0).getPaneInfo().getIndex());
22072215
assertEquals(n - ((n / 3) * 3), Iterables.size(output.get(0).getValue()));
22082216
}
22092217

@@ -2231,17 +2239,17 @@ public void fireEmptyOnDrainInGlobalWindowIfRequested() throws Exception {
22312239
List<WindowedValue<Iterable<Integer>>> output = tester.extractOutput();
22322240
assertEquals((n + 3) / 4, output.size());
22332241
for (int i = 0; i < output.size(); i++) {
2234-
assertEquals(Timing.EARLY, output.get(i).getPane().getTiming());
2235-
assertEquals(i, output.get(i).getPane().getIndex());
2242+
assertEquals(Timing.EARLY, output.get(i).getPaneInfo().getTiming());
2243+
assertEquals(i, output.get(i).getPaneInfo().getIndex());
22362244
assertEquals(4, Iterables.size(output.get(i).getValue()));
22372245
}
22382246

22392247
tester.advanceInputWatermark(BoundedWindow.TIMESTAMP_MAX_VALUE);
22402248

22412249
output = tester.extractOutput();
22422250
assertEquals(1, output.size());
2243-
assertEquals(Timing.ON_TIME, output.get(0).getPane().getTiming());
2244-
assertEquals((n + 3) / 4, output.get(0).getPane().getIndex());
2251+
assertEquals(Timing.ON_TIME, output.get(0).getPaneInfo().getTiming());
2252+
assertEquals((n + 3) / 4, output.get(0).getPaneInfo().getIndex());
22452253
assertEquals(0, Iterables.size(output.get(0).getValue()));
22462254
}
22472255

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

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -155,13 +155,13 @@ public void describeTo(Description description) {
155155

156156
@Override
157157
protected boolean matchesSafely(WindowedValue<? extends T> item) {
158-
return Objects.equals(item.getPane(), paneInfo);
158+
return Objects.equals(item.getPaneInfo(), paneInfo);
159159
}
160160

161161
@Override
162162
protected void describeMismatchSafely(
163163
WindowedValue<? extends T> item, Description mismatchDescription) {
164-
mismatchDescription.appendValue(item.getPane());
164+
mismatchDescription.appendValue(item.getPaneInfo());
165165
}
166166
};
167167
}
@@ -212,7 +212,7 @@ protected boolean matchesSafely(WindowedValue<? extends T> windowedValue) {
212212
return valueMatcher.matches(windowedValue.getValue())
213213
&& timestampMatcher.matches(windowedValue.getTimestamp())
214214
&& windowsMatcher.matches(windowedValue.getWindows())
215-
&& paneInfoMatcher.matches(windowedValue.getPane());
215+
&& paneInfoMatcher.matches(windowedValue.getPaneInfo());
216216
}
217217
}
218218
}

runners/direct-java/src/main/java/org/apache/beam/runners/direct/SideInputContainer.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -159,7 +159,7 @@ private void updatePCollectionViewWindowValues(
159159
// the value had never been set, so we set it and are done.
160160
return;
161161
}
162-
PaneInfo newPane = windowValues.iterator().next().getPane();
162+
PaneInfo newPane = windowValues.iterator().next().getPaneInfo();
163163

164164
Iterable<? extends WindowedValue<?>> existingValues;
165165
long existingPane;
@@ -168,7 +168,7 @@ private void updatePCollectionViewWindowValues(
168168
existingPane =
169169
Iterables.isEmpty(existingValues)
170170
? -1L
171-
: existingValues.iterator().next().getPane().getIndex();
171+
: existingValues.iterator().next().getPaneInfo().getIndex();
172172
} while (newPane.getIndex() > existingPane
173173
&& !contents.compareAndSet(existingValues, windowValues));
174174
}

0 commit comments

Comments
 (0)