Skip to content

Commit 0884c46

Browse files
authored
Optimze away WindmillWatermarkHold::clear when the cached hold is empty (#38297)
1 parent c5d0ab1 commit 0884c46

3 files changed

Lines changed: 42 additions & 10 deletions

File tree

runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillWatermarkHold.java

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -72,9 +72,13 @@ public class WindmillWatermarkHold extends WindmillState implements WatermarkHol
7272

7373
@Override
7474
public void clear() {
75+
localAdditions = null;
76+
if (cachedValue != null && !cachedValue.isPresent()) {
77+
// No need to clear the backend as it is known empty.
78+
return;
79+
}
7580
cleared = true;
7681
cachedValue = Optional.absent();
77-
localAdditions = null;
7882
}
7983

8084
@Override

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

Lines changed: 3 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -2374,9 +2374,6 @@ public void testMergeSessionWindows_singleLateWindow() throws Exception {
23742374
new Action(
23752375
buildSessionInput(
23762376
1, 40, 0, Collections.singletonList(1L), Collections.EMPTY_LIST))
2377-
.withHolds(
2378-
buildHold("/gAAAAAAAAAsK/+uhold", -1, true),
2379-
buildHold("/gAAAAAAAAAsK/+uextra", -1, true))
23802377
.withTimers(buildWatermarkTimer("/s/gAAAAAAAAAsK/+0", 3600010))));
23812378
}
23822379

@@ -2404,10 +2401,7 @@ public void testMergeSessionWindows() throws Exception {
24042401
0,
24052402
Collections.EMPTY_LIST,
24062403
Collections.singletonList(buildWatermarkTimer("/s/gAAAAAAAAAsK/+0", 10))))
2407-
.withTimers(buildWatermarkTimer("/s/gAAAAAAAAAsK/+0", 3600010))
2408-
.withHolds(
2409-
buildHold("/gAAAAAAAAAsK/+uhold", -1, true),
2410-
buildHold("/gAAAAAAAAAsK/+uextra", -1, true)),
2404+
.withTimers(buildWatermarkTimer("/s/gAAAAAAAAAsK/+0", 3600010)),
24112405
new Action(
24122406
buildSessionInput(
24132407
3, 30, 0, Collections.singletonList(8L), Collections.EMPTY_LIST))
@@ -2436,8 +2430,8 @@ public void testMergeSessionWindows() throws Exception {
24362430
.withHolds(
24372431
buildHold("/gAAAAAAAACkK/+uhold", -1, true),
24382432
buildHold("/gAAAAAAAACkK/+uextra", -1, true),
2439-
buildHold("/gAAAAAAAAAsK/+uhold", 40, true),
2440-
buildHold("/gAAAAAAAAAsK/+uextra", 3600040, true)),
2433+
buildHold("/gAAAAAAAAAsK/+uhold", 40, false),
2434+
buildHold("/gAAAAAAAAAsK/+uextra", 3600040, false)),
24412435
new Action(
24422436
buildSessionInput(
24432437
6,

runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateInternalsTest.java

Lines changed: 34 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3228,6 +3228,40 @@ public void testNewWatermarkNoFetch() throws Exception {
32283228
Mockito.verifyNoInteractions(mockReader);
32293229
}
32303230

3231+
@Test
3232+
public void testWatermarkClearNoOp() throws Exception {
3233+
StateTag<WatermarkHoldState> addr =
3234+
StateTags.watermarkStateInternal("watermark", TimestampCombiner.EARLIEST);
3235+
WatermarkHoldState hold = underTestNewKey.state(NAMESPACE, addr);
3236+
3237+
hold.clear();
3238+
3239+
Windmill.WorkItemCommitRequest.Builder commitBuilder =
3240+
Windmill.WorkItemCommitRequest.newBuilder();
3241+
underTestNewKey.persist(commitBuilder);
3242+
3243+
assertEquals(0, commitBuilder.getWatermarkHoldsCount());
3244+
assertBuildable(commitBuilder);
3245+
}
3246+
3247+
@Test
3248+
public void testWatermarkClearWithLocalAdditionsNoop() throws Exception {
3249+
StateTag<WatermarkHoldState> addr =
3250+
StateTags.watermarkStateInternal("watermark", TimestampCombiner.EARLIEST);
3251+
WatermarkHoldState hold = underTestNewKey.state(NAMESPACE, addr);
3252+
3253+
hold.add(new Instant(500));
3254+
3255+
hold.clear();
3256+
3257+
Windmill.WorkItemCommitRequest.Builder commitBuilder =
3258+
Windmill.WorkItemCommitRequest.newBuilder();
3259+
underTestNewKey.persist(commitBuilder);
3260+
3261+
assertEquals(0, commitBuilder.getWatermarkHoldsCount());
3262+
assertBuildable(commitBuilder);
3263+
}
3264+
32313265
@Test
32323266
public void testValueSetBeforeRead() throws Exception {
32333267
StateTag<ValueState<String>> addr = StateTags.value("value", StringUtf8Coder.of());

0 commit comments

Comments
 (0)