Skip to content

Commit 399d9d7

Browse files
authored
Introduce ValueKind to Java and add to WindowedValue (#38315)
* introduce ValueKind to Java and add to WindowedValue * map UNSPECIFIED to INSERT; cleanup * rebase remnants * adjust equals and hashcode methods to include valuekind * use existing coder test * compile fixes * compile fixes in dataflow * compile fixes in dataflow * add test to WindmillKeyedWorkItem * compile fixes
1 parent a86f2ec commit 399d9d7

17 files changed

Lines changed: 428 additions & 59 deletions

File tree

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

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -152,7 +152,8 @@ public <K, InputT> Iterable<WindowedValue<InputT>> filter(
152152
element.getRecordId(),
153153
element.getRecordOffset(),
154154
element.causedByDrain(),
155-
element.getOpenTelemetryContext()));
155+
element.getOpenTelemetryContext(),
156+
element.getValueKind()));
156157
}
157158
}
158159
}

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

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -452,7 +452,8 @@ public <T> void outputWithTimestamp(TupleTag<T> tag, T value, Instant timestamp)
452452
element.getRecordId(),
453453
element.getRecordOffset(),
454454
element.causedByDrain(),
455-
element.getOpenTelemetryContext()));
455+
element.getOpenTelemetryContext(),
456+
element.getValueKind()));
456457
}
457458

458459
@Override
@@ -476,7 +477,8 @@ public <T> void outputWindowedValue(
476477
element.getRecordId(),
477478
element.getRecordOffset(),
478479
element.causedByDrain(),
479-
element.getOpenTelemetryContext()));
480+
element.getOpenTelemetryContext(),
481+
element.getValueKind()));
480482
}
481483

482484
@Override

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

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -473,7 +473,8 @@ public String getErrorContext() {
473473
read.getRecordId(),
474474
read.getRecordOffset(),
475475
CausedByDrain.CAUSED_BY_DRAIN,
476-
read.getOpenTelemetryContext());
476+
read.getOpenTelemetryContext(),
477+
read.getValueKind());
477478
}
478479
elementAndRestriction = KV.of(read, restrictionState.read());
479480
watermarkEstimatorStateT = watermarkEstimatorState.read();

runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/BatchViewOverrides.java

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -74,6 +74,7 @@
7474
import org.apache.beam.sdk.values.PCollectionView;
7575
import org.apache.beam.sdk.values.TupleTag;
7676
import org.apache.beam.sdk.values.TupleTagList;
77+
import org.apache.beam.sdk.values.ValueKind;
7778
import org.apache.beam.sdk.values.WindowedValue;
7879
import org.apache.beam.sdk.values.WindowedValues;
7980
import org.apache.beam.sdk.values.WindowedValues.FullWindowedValueCoder;
@@ -1415,6 +1416,11 @@ public PaneInfo getPaneInfo() {
14151416
return null;
14161417
}
14171418

1419+
@Override
1420+
public ValueKind getValueKind() {
1421+
return ValueKind.INSERT;
1422+
}
1423+
14181424
@Override
14191425
public Iterable<WindowedValue<T>> explodeWindows() {
14201426
return Collections.emptyList();

runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/UngroupedWindmillReader.java

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -38,6 +38,8 @@
3838
import org.apache.beam.sdk.transforms.windowing.PaneInfo;
3939
import org.apache.beam.sdk.values.CausedByDrain;
4040
import org.apache.beam.sdk.values.KV;
41+
import org.apache.beam.sdk.values.ValueKind;
42+
import org.apache.beam.sdk.values.ValueKindUtil;
4143
import org.apache.beam.sdk.values.WindowedValue;
4244
import org.apache.beam.sdk.values.WindowedValues;
4345
import org.apache.beam.sdk.values.WindowedValues.FullWindowedValueCoder;
@@ -139,13 +141,15 @@ protected WindowedValue<T> decodeMessage(Windmill.Message message) throws IOExce
139141
* drain happened upstream
140142
*/
141143
CausedByDrain drainingValueFromUpstream = CausedByDrain.NORMAL;
144+
ValueKind valueKind = ValueKind.INSERT;
142145
if (WindowedValues.WindowedValueCoder.isMetadataSupported()) {
143146
BeamFnApi.Elements.ElementMetadata elementMetadata =
144147
WindmillSink.decodeAdditionalMetadata(windowsCoder, message.getMetadata());
145148
drainingValueFromUpstream =
146149
elementMetadata.getDrain() == BeamFnApi.Elements.DrainMode.Enum.DRAINING
147150
? CausedByDrain.CAUSED_BY_DRAIN
148151
: CausedByDrain.NORMAL;
152+
valueKind = ValueKindUtil.fromProto(elementMetadata.getValueKind());
149153
}
150154
if (valueCoder instanceof KvCoder) {
151155
KvCoder<?, ?> kvCoder = (KvCoder<?, ?>) valueCoder;
@@ -164,7 +168,8 @@ protected WindowedValue<T> decodeMessage(Windmill.Message message) throws IOExce
164168
null,
165169
null,
166170
drainingValueFromUpstream,
167-
null);
171+
null,
172+
valueKind);
168173
} else {
169174
notifyElementRead(data.available() + metadata.available());
170175
// todo #37030 parse context from previous stage
@@ -176,7 +181,8 @@ protected WindowedValue<T> decodeMessage(Windmill.Message message) throws IOExce
176181
null,
177182
null,
178183
drainingValueFromUpstream,
179-
null);
184+
null,
185+
valueKind);
180186
}
181187
}
182188

runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItem.java

Lines changed: 13 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -42,6 +42,8 @@
4242
import org.apache.beam.sdk.transforms.windowing.PaneInfo;
4343
import org.apache.beam.sdk.util.common.ElementByteSizeObserver;
4444
import org.apache.beam.sdk.values.CausedByDrain;
45+
import org.apache.beam.sdk.values.ValueKind;
46+
import org.apache.beam.sdk.values.ValueKindUtil;
4547
import org.apache.beam.sdk.values.WindowedValue;
4648
import org.apache.beam.sdk.values.WindowedValues;
4749
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Predicate;
@@ -148,18 +150,28 @@ public Iterable<TimerData> timersIterable() {
148150
* drain happened upstream
149151
*/
150152
CausedByDrain drainingValueFromUpstream = CausedByDrain.NORMAL;
153+
ValueKind valueKind = ValueKind.INSERT;
151154
if (WindowedValues.WindowedValueCoder.isMetadataSupported()) {
152155
BeamFnApi.Elements.ElementMetadata elementMetadata =
153156
WindmillSink.decodeAdditionalMetadata(windowsCoder, message.getMetadata());
154157
drainingValueFromUpstream =
155158
elementMetadata.getDrain() == BeamFnApi.Elements.DrainMode.Enum.DRAINING
156159
? CausedByDrain.CAUSED_BY_DRAIN
157160
: CausedByDrain.NORMAL;
161+
valueKind = ValueKindUtil.fromProto(elementMetadata.getValueKind());
158162
}
159163
InputStream inputStream = message.getData().newInput();
160164
ElemT value = valueCoder.decode(inputStream, Coder.Context.OUTER);
161165
return WindowedValues.of(
162-
value, timestamp, windows, paneInfo, null, null, drainingValueFromUpstream, null);
166+
value,
167+
timestamp,
168+
windows,
169+
paneInfo,
170+
null,
171+
null,
172+
drainingValueFromUpstream,
173+
null,
174+
valueKind);
163175
} catch (RuntimeException | IOException e) {
164176
if (!skipUndecodableElements) {
165177
throw new RuntimeException(e);

runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/ValueInEmptyWindows.java

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,7 @@
2424
import org.apache.beam.sdk.transforms.windowing.BoundedWindow;
2525
import org.apache.beam.sdk.transforms.windowing.PaneInfo;
2626
import org.apache.beam.sdk.values.CausedByDrain;
27+
import org.apache.beam.sdk.values.ValueKind;
2728
import org.apache.beam.sdk.values.WindowedValue;
2829
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.MoreObjects;
2930
import org.checkerframework.checker.nullness.qual.Nullable;
@@ -71,6 +72,11 @@ public CausedByDrain causedByDrain() {
7172
return null;
7273
}
7374

75+
@Override
76+
public ValueKind getValueKind() {
77+
return ValueKind.INSERT;
78+
}
79+
7480
@Override
7581
public Iterable<WindowedValue<T>> explodeWindows() {
7682
return Collections.emptyList();

runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItemTest.java

Lines changed: 43 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -46,6 +46,7 @@
4646
import org.apache.beam.sdk.transforms.windowing.PaneInfo;
4747
import org.apache.beam.sdk.transforms.windowing.PaneInfo.Timing;
4848
import org.apache.beam.sdk.values.CausedByDrain;
49+
import org.apache.beam.sdk.values.ValueKind;
4950
import org.apache.beam.sdk.values.WindowedValue;
5051
import org.apache.beam.sdk.values.WindowedValues;
5152
import org.apache.beam.vendor.grpc.v1p69p0.com.google.protobuf.ByteString;
@@ -374,6 +375,48 @@ public void testDrainPropagated() throws Exception {
374375
assertThat(
375376
keyedWorkItem.timersIterable(),
376377
Matchers.contains(makeDrainingTimer(STATE_NAMESPACE_2, 3, TimeDomain.EVENT_TIME)));
378+
WindowedValues.WindowedValueCoder.setMetadataNotSupported();
379+
}
380+
381+
@Test
382+
public void testValueKindPropagated() throws Exception {
383+
WindowedValues.WindowedValueCoder.setMetadataSupported();
384+
Windmill.WorkItem.Builder workItem =
385+
Windmill.WorkItem.newBuilder().setKey(SERIALIZED_KEY).setWorkToken(17);
386+
Windmill.InputMessageBundle.Builder chunk1 = workItem.addMessageBundlesBuilder();
387+
chunk1.setSourceComputationId("computation");
388+
addElementWithMetadata(
389+
chunk1,
390+
5,
391+
"hello",
392+
WINDOW_1,
393+
paneInfo(0),
394+
BeamFnApi.Elements.ElementMetadata.newBuilder()
395+
.setValueKind(BeamFnApi.Elements.ValueKind.Enum.UPDATE_AFTER)
396+
.build());
397+
addElementWithMetadata(
398+
chunk1,
399+
7,
400+
"world",
401+
WINDOW_2,
402+
paneInfo(2),
403+
BeamFnApi.Elements.ElementMetadata.newBuilder()
404+
.setValueKind(BeamFnApi.Elements.ValueKind.Enum.DELETE)
405+
.build());
406+
KeyedWorkItem<String, String> keyedWorkItem =
407+
new WindmillKeyedWorkItem<>(
408+
KEY,
409+
workItem.build(),
410+
WINDOW_CODER,
411+
WINDOWS_CODER,
412+
VALUE_CODER,
413+
windmillTagEncoding,
414+
true);
415+
416+
Iterator<WindowedValue<String>> iterator = keyedWorkItem.elementsIterable().iterator();
417+
Assert.assertEquals(ValueKind.UPDATE_AFTER, iterator.next().getValueKind());
418+
Assert.assertEquals(ValueKind.DELETE, iterator.next().getValueKind());
419+
WindowedValues.WindowedValueCoder.setMetadataNotSupported();
377420
}
378421

379422
private static TimeDomain timerTypeToTimeDomain(Windmill.Timer.Type type) {

runners/spark/src/main/java/org/apache/beam/runners/spark/util/TimerUtils.java

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -35,6 +35,7 @@
3535
import org.apache.beam.sdk.transforms.windowing.PaneInfo;
3636
import org.apache.beam.sdk.values.CausedByDrain;
3737
import org.apache.beam.sdk.values.KV;
38+
import org.apache.beam.sdk.values.ValueKind;
3839
import org.apache.beam.sdk.values.WindowedValue;
3940
import org.apache.beam.sdk.values.WindowingStrategy;
4041
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.MoreObjects;
@@ -127,6 +128,11 @@ public CausedByDrain causedByDrain() {
127128
return CausedByDrain.NORMAL;
128129
}
129130

131+
@Override
132+
public ValueKind getValueKind() {
133+
return ValueKind.INSERT;
134+
}
135+
130136
@Override
131137
public @Nullable Long getRecordOffset() {
132138
return null;

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
@@ -53,5 +53,7 @@ public interface OutputBuilder<T> extends WindowedValue<T> {
5353

5454
OutputBuilder<T> setOpenTelemetryContext(@Nullable Context openTelemetryContext);
5555

56+
OutputBuilder<T> setValueKind(ValueKind valueKind);
57+
5658
void output();
5759
}

0 commit comments

Comments
 (0)