Skip to content

Commit 7904ba3

Browse files
committed
Make WindowedValue a public interface
The following mostly-automated changes: - Moved WindowedValue from util to values package - Make WindowedValue an interface with companion class WindowedValues
1 parent 60307b4 commit 7904ba3

465 files changed

Lines changed: 2217 additions & 1930 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/DoFnRunner.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,7 @@
2020
import org.apache.beam.sdk.state.TimeDomain;
2121
import org.apache.beam.sdk.transforms.DoFn;
2222
import org.apache.beam.sdk.transforms.windowing.BoundedWindow;
23-
import org.apache.beam.sdk.util.WindowedValue;
23+
import org.apache.beam.sdk.values.WindowedValue;
2424
import org.checkerframework.checker.nullness.qual.Nullable;
2525
import org.joda.time.Instant;
2626

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -29,10 +29,10 @@
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.WindowedValue;
3332
import org.apache.beam.sdk.values.KV;
3433
import org.apache.beam.sdk.values.PCollectionView;
3534
import org.apache.beam.sdk.values.TupleTag;
35+
import org.apache.beam.sdk.values.WindowedValue;
3636
import org.apache.beam.sdk.values.WindowingStrategy;
3737
import org.checkerframework.checker.nullness.qual.Nullable;
3838

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

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -25,10 +25,10 @@
2525
import org.apache.beam.sdk.transforms.windowing.BoundedWindow;
2626
import org.apache.beam.sdk.transforms.windowing.PaneInfo;
2727
import org.apache.beam.sdk.util.SystemDoFnInternal;
28-
import org.apache.beam.sdk.util.WindowedValue;
2928
import org.apache.beam.sdk.util.construction.TriggerTranslation;
3029
import org.apache.beam.sdk.values.KV;
3130
import org.apache.beam.sdk.values.TupleTag;
31+
import org.apache.beam.sdk.values.WindowedValues;
3232
import org.apache.beam.sdk.values.WindowingStrategy;
3333
import org.joda.time.Instant;
3434

@@ -99,7 +99,7 @@ public void outputWindowedValue(
9999
Instant timestamp,
100100
Collection<? extends BoundedWindow> windows,
101101
PaneInfo pane) {
102-
outputManager.output(mainTag, WindowedValue.of(output, timestamp, windows, pane));
102+
outputManager.output(mainTag, WindowedValues.of(output, timestamp, windows, pane));
103103
}
104104

105105
@Override
@@ -109,7 +109,7 @@ public <AdditionalOutputT> void outputWindowedValue(
109109
Instant timestamp,
110110
Collection<? extends BoundedWindow> windows,
111111
PaneInfo pane) {
112-
outputManager.output(tag, WindowedValue.of(output, timestamp, windows, pane));
112+
outputManager.output(tag, WindowedValues.of(output, timestamp, windows, pane));
113113
}
114114
};
115115
}

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

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -29,10 +29,10 @@
2929
import org.apache.beam.sdk.transforms.GroupByKey;
3030
import org.apache.beam.sdk.transforms.PTransform;
3131
import org.apache.beam.sdk.transforms.ParDo;
32-
import org.apache.beam.sdk.util.WindowedValue;
33-
import org.apache.beam.sdk.util.WindowedValue.WindowedValueCoder;
3432
import org.apache.beam.sdk.values.KV;
3533
import org.apache.beam.sdk.values.PCollection;
34+
import org.apache.beam.sdk.values.WindowedValue;
35+
import org.apache.beam.sdk.values.WindowedValues.WindowedValueCoder;
3636
import org.apache.beam.sdk.values.WindowingStrategy;
3737

3838
/**

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,7 @@
1818
package org.apache.beam.runners.core;
1919

2020
import org.apache.beam.runners.core.TimerInternals.TimerData;
21-
import org.apache.beam.sdk.util.WindowedValue;
21+
import org.apache.beam.sdk.values.WindowedValue;
2222

2323
/**
2424
* Interface that contains all the timers and elements associated with a specific work item.

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

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -28,8 +28,8 @@
2828
import org.apache.beam.sdk.coders.IterableCoder;
2929
import org.apache.beam.sdk.coders.StructuredCoder;
3030
import org.apache.beam.sdk.transforms.windowing.BoundedWindow;
31-
import org.apache.beam.sdk.util.WindowedValue;
32-
import org.apache.beam.sdk.util.WindowedValue.FullWindowedValueCoder;
31+
import org.apache.beam.sdk.values.WindowedValue;
32+
import org.apache.beam.sdk.values.WindowedValues.FullWindowedValueCoder;
3333
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList;
3434

3535
/** A {@link Coder} for {@link KeyedWorkItem KeyedWorkItems}. */

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,7 @@
2020
import java.util.Collections;
2121
import java.util.Objects;
2222
import org.apache.beam.runners.core.TimerInternals.TimerData;
23-
import org.apache.beam.sdk.util.WindowedValue;
23+
import org.apache.beam.sdk.values.WindowedValue;
2424
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.MoreObjects;
2525
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Iterables;
2626
import org.checkerframework.checker.nullness.qual.Nullable;

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

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -25,8 +25,9 @@
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;
2928
import org.apache.beam.sdk.values.KV;
29+
import org.apache.beam.sdk.values.WindowedValue;
30+
import org.apache.beam.sdk.values.WindowedValues;
3031
import org.apache.beam.sdk.values.WindowingStrategy;
3132
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting;
3233
import org.joda.time.Instant;
@@ -140,7 +141,7 @@ public <K, InputT> Iterable<WindowedValue<InputT>> filter(
140141
timerInternals.currentOutputWatermarkTime());
141142
} else {
142143
nonLateElements.add(
143-
WindowedValue.of(
144+
WindowedValues.of(
144145
element.getValue(), element.getTimestamp(), window, element.getPane()));
145146
}
146147
}

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,7 @@
2121
import org.apache.beam.sdk.transforms.windowing.BoundedWindow;
2222
import org.apache.beam.sdk.transforms.windowing.GlobalWindow;
2323
import org.apache.beam.sdk.util.WindowTracing;
24-
import org.apache.beam.sdk.util.WindowedValue;
24+
import org.apache.beam.sdk.values.WindowedValue;
2525
import org.apache.beam.sdk.values.WindowingStrategy;
2626
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.FluentIterable;
2727
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Iterables;

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -45,11 +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.WindowedValue;
4948
import org.apache.beam.sdk.values.KV;
5049
import org.apache.beam.sdk.values.PCollectionView;
5150
import org.apache.beam.sdk.values.Row;
5251
import org.apache.beam.sdk.values.TupleTag;
52+
import org.apache.beam.sdk.values.WindowedValue;
5353
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Iterables;
5454
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.util.concurrent.Futures;
5555
import org.checkerframework.checker.nullness.qual.Nullable;

0 commit comments

Comments
 (0)