Skip to content

Commit 28957f5

Browse files
committed
Use trigger state to know if a window is new
1 parent aaecfa8 commit 28957f5

8 files changed

Lines changed: 176 additions & 10 deletions

File tree

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

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -274,6 +274,16 @@ boolean hasNoActiveWindows() {
274274
return activeWindows.getActiveAndNewWindows().isEmpty();
275275
}
276276

277+
@VisibleForTesting
278+
TriggerStateMachineRunner<W> getTriggerRunner() {
279+
return triggerRunner;
280+
}
281+
282+
@VisibleForTesting
283+
ReduceFnContextFactory<K, InputT, OutputT, W> getContextFactory() {
284+
return contextFactory;
285+
}
286+
277287
private Set<W> windowsThatAreOpen(Collection<W> windows) {
278288
Set<W> result = new HashSet<>();
279289
for (W window : windows) {
@@ -603,6 +613,14 @@ private void processElement(Map<W, W> windowToMergeResult, WindowedValue<InputT>
603613
contextFactory.forValue(
604614
window, value.getValue(), value.getTimestamp(), StateStyle.RENAMED);
605615

616+
if (triggerRunner.isNew(directContext.state())) {
617+
// Blindly clear state to ensure Windmill doesn't do unnecessary reads.
618+
reduceFn.clearState(renamedContext);
619+
paneInfoTracker.clear(directContext.state());
620+
watermarkHold.setKnownEmpty(renamedContext);
621+
nonEmptyPanes.clearPane(renamedContext.state());
622+
}
623+
606624
nonEmptyPanes.recordContent(renamedContext.state());
607625
scheduleGarbageCollectionTimer(directContext);
608626

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

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -466,6 +466,23 @@ public void clearHolds(ReduceFn<?, ?, ?, W>.Context context) {
466466
context.state().access(EXTRA_HOLD_TAG).clear();
467467
}
468468

469+
/**
470+
* <b><i>For internal use only; no backwards-compatibility guarantees.</i></b>
471+
*
472+
* <p>Permit marking the watermark holds as empty locally, without necessarily clearing them in
473+
* the backend.
474+
*/
475+
public void setKnownEmpty(ReduceFn<?, ?, ?, W>.Context context) {
476+
WindowTracing.debug(
477+
"WatermarkHold.setKnownEmpty: For key:{}; window:{}; inputWatermark:{}; outputWatermark:{}",
478+
context.key(),
479+
context.window(),
480+
timerInternals.currentInputWatermarkTime(),
481+
timerInternals.currentOutputWatermarkTime());
482+
context.state().access(elementHoldTag).setKnownEmpty();
483+
context.state().access(EXTRA_HOLD_TAG).setKnownEmpty();
484+
}
485+
469486
/** Return the current data hold, or null if none. Does not clear. For debugging only. */
470487
public @Nullable Instant getDataCurrent(ReduceFn<?, ?, ?, W>.Context context) {
471488
return context.state().access(elementHoldTag).read();

runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/FinishedTriggersBitSet.java

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@
1818
package org.apache.beam.runners.core.triggers;
1919

2020
import java.util.BitSet;
21+
import org.checkerframework.checker.nullness.qual.Nullable;
2122

2223
/** A {@link FinishedTriggers} implementation based on an underlying {@link BitSet}. */
2324
public class FinishedTriggersBitSet implements FinishedTriggers {
@@ -60,4 +61,17 @@ public void clearRecursively(ExecutableTriggerStateMachine trigger) {
6061
public FinishedTriggersBitSet copy() {
6162
return new FinishedTriggersBitSet((BitSet) bitSet.clone());
6263
}
64+
65+
@Override
66+
public boolean equals(@Nullable Object obj) {
67+
if (!(obj instanceof FinishedTriggersBitSet)) {
68+
return false;
69+
}
70+
return bitSet.equals(((FinishedTriggersBitSet) obj).bitSet);
71+
}
72+
73+
@Override
74+
public int hashCode() {
75+
return bitSet.hashCode();
76+
}
6377
}

runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/TriggerStateMachineRunner.java

Lines changed: 18 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -81,9 +81,11 @@ private FinishedTriggersBitSet readFinishedBits(ValueState<BitSet> state) {
8181
}
8282

8383
@Nullable BitSet bitSet = state.read();
84-
return bitSet == null
85-
? FinishedTriggersBitSet.emptyWithCapacity(rootTrigger.getFirstIndexAfterSubtree())
86-
: FinishedTriggersBitSet.fromBitSet(bitSet);
84+
if (bitSet == null) {
85+
return FinishedTriggersBitSet.emptyWithCapacity(rootTrigger.getFirstIndexAfterSubtree());
86+
}
87+
88+
return FinishedTriggersBitSet.fromBitSet(bitSet);
8789
}
8890

8991
private void clearFinishedBits(ValueState<BitSet> state) {
@@ -99,6 +101,16 @@ public boolean isClosed(StateAccessor<?> state) {
99101
return readFinishedBits(state.access(FINISHED_BITS_TAG)).isFinished(rootTrigger);
100102
}
101103

104+
/** Return true if the window is new (no trigger state has ever been persisted). */
105+
public boolean isNew(StateAccessor<?> state) {
106+
return isFinishedSetNeeded() && state.access(FINISHED_BITS_TAG).read() == null;
107+
}
108+
109+
@VisibleForTesting
110+
public BitSet getFinishedBits(StateAccessor<?> state) {
111+
return readFinishedBits(state.access(FINISHED_BITS_TAG)).getBitSet();
112+
}
113+
102114
public void prefetchIsClosed(StateAccessor<?> state) {
103115
if (isFinishedSetNeeded()) {
104116
state.access(FINISHED_BITS_TAG).readLater();
@@ -187,12 +199,9 @@ private void persistFinishedSet(
187199
}
188200

189201
ValueState<BitSet> finishedSetState = state.access(FINISHED_BITS_TAG);
190-
if (!readFinishedBits(finishedSetState).equals(modifiedFinishedSet)) {
191-
if (modifiedFinishedSet.getBitSet().isEmpty()) {
192-
finishedSetState.clear();
193-
} else {
194-
finishedSetState.write(modifiedFinishedSet.getBitSet());
195-
}
202+
@Nullable BitSet currentBits = finishedSetState.read();
203+
if (currentBits == null || !currentBits.equals(modifiedFinishedSet.getBitSet())) {
204+
finishedSetState.write(modifiedFinishedSet.getBitSet());
196205
}
197206
}
198207

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

Lines changed: 48 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -40,9 +40,11 @@
4040

4141
import java.util.ArrayList;
4242
import java.util.Arrays;
43+
import java.util.BitSet;
4344
import java.util.List;
4445
import java.util.Random;
4546
import java.util.concurrent.ThreadLocalRandom;
47+
import org.apache.beam.runners.core.ReduceFnContextFactory.StateStyle;
4648
import org.apache.beam.runners.core.metrics.MetricsContainerImpl;
4749
import org.apache.beam.runners.core.triggers.DefaultTriggerStateMachine;
4850
import org.apache.beam.runners.core.triggers.TriggerStateMachine;
@@ -2343,4 +2345,50 @@ public interface TestOptions extends PipelineOptions {
23432345

23442346
void setValue(int value);
23452347
}
2348+
2349+
@Test
2350+
public void testNewWindowOptimization() throws Exception {
2351+
WindowingStrategy<?, IntervalWindow> strategy =
2352+
WindowingStrategy.of(FixedWindows.of(Duration.millis(10)))
2353+
.withTrigger(AfterPane.elementCountAtLeast(2))
2354+
.withMode(AccumulationMode.ACCUMULATING_FIRED_PANES);
2355+
2356+
ReduceFnTester<Integer, Iterable<Integer>, IntervalWindow> tester =
2357+
ReduceFnTester.nonCombining(strategy);
2358+
2359+
IntervalWindow window = new IntervalWindow(new Instant(0), new Instant(10));
2360+
2361+
// 1. First element for a new window.
2362+
tester.injectElements(TimestampedValue.of(1, new Instant(1)));
2363+
2364+
// Verify sentinel bit is written.
2365+
BitSet bitSet =
2366+
tester
2367+
.createRunner()
2368+
.getTriggerRunner()
2369+
.getFinishedBits(
2370+
tester.createRunner().getContextFactory().base(window, StateStyle.DIRECT).state());
2371+
2372+
// We expect the bitset to be empty (the sentinel bit is no longer used).
2373+
assertTrue("Bitset should be empty", bitSet.isEmpty());
2374+
// And trigger not finished.
2375+
assertFalse("Trigger should not be finished", bitSet.get(0));
2376+
2377+
// And verify that it is no longer "new".
2378+
assertFalse(
2379+
"Window should no longer be new",
2380+
tester
2381+
.createRunner()
2382+
.getTriggerRunner()
2383+
.isNew(
2384+
tester.createRunner().getContextFactory().base(window, StateStyle.DIRECT).state()));
2385+
2386+
// 2. Second element for the same window.
2387+
// We want to verify it doesn't clear the first element.
2388+
tester.injectElements(TimestampedValue.of(2, new Instant(2)));
2389+
2390+
// Extract output.
2391+
List<WindowedValue<Iterable<Integer>>> output = tester.extractOutput();
2392+
assertThat(output, contains(isSingleWindowedValue(containsInAnyOrder(1, 2), 9, 0, 10)));
2393+
}
23462394
}

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

Lines changed: 14 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -37,6 +37,7 @@
3737
"nullness" // TODO(https://github.com/apache/beam/issues/20497)
3838
})
3939
public class WindmillWatermarkHold extends WindmillState implements WatermarkHoldState {
40+
4041
// The encoded size of an Instant.
4142
private static final int ENCODED_SIZE = 8;
4243

@@ -46,6 +47,7 @@ public class WindmillWatermarkHold extends WindmillState implements WatermarkHol
4647
private final String stateFamily;
4748

4849
private boolean cleared = false;
50+
private boolean knownEmpty = false;
4951
/**
5052
* If non-{@literal null}, the known current hold value, or absent if we know there are no output
5153
* watermark holds. If {@literal null}, the current hold value could depend on holds in Windmill
@@ -77,6 +79,13 @@ public void clear() {
7779
localAdditions = null;
7880
}
7981

82+
@Override
83+
public void setKnownEmpty() {
84+
cachedValue = Optional.absent();
85+
localAdditions = null;
86+
knownEmpty = true;
87+
}
88+
8089
@Override
8190
@SuppressWarnings("FutureReturnValueIgnored")
8291
public WindmillWatermarkHold readLater() {
@@ -133,7 +142,7 @@ public Future<Windmill.WorkItemCommitRequest> persist(
133142

134143
Future<Windmill.WorkItemCommitRequest> result;
135144

136-
if (!cleared && localAdditions == null) {
145+
if (!knownEmpty && !cleared && localAdditions == null) {
137146
// No changes, so no need to update Windmill and no need to cache any value.
138147
return Futures.immediateFuture(Windmill.WorkItemCommitRequest.newBuilder().buildPartial());
139148
}
@@ -166,15 +175,19 @@ public Future<Windmill.WorkItemCommitRequest> persist(
166175
} else if (!cleared && localAdditions != null) {
167176
// Otherwise, we need to combine the local additions with the already persisted data
168177
result = combineWithPersisted();
178+
} else if (knownEmpty) {
179+
result = Futures.immediateFuture(Windmill.WorkItemCommitRequest.newBuilder().buildPartial());
169180
} else {
170181
throw new IllegalStateException("Unreachable condition");
171182
}
172183

173184
final int estimatedByteSize = ENCODED_SIZE + stateKey.byteString().size();
185+
174186
return Futures.lazyTransform(
175187
result,
176188
result1 -> {
177189
cleared = false;
190+
knownEmpty = false;
178191
localAdditions = null;
179192
if (cachedValue != null) {
180193
cache.put(namespace, stateKey, WindmillWatermarkHold.this, estimatedByteSize);

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

Lines changed: 37 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3037,6 +3037,43 @@ public void testWatermarkClearBeforeRead() throws Exception {
30373037
Mockito.verifyNoMoreInteractions(mockReader);
30383038
}
30393039

3040+
@Test
3041+
public void testWatermarkSetKnownEmptyBeforeRead() throws Exception {
3042+
StateTag<WatermarkHoldState> addr =
3043+
StateTags.watermarkStateInternal("watermark", TimestampCombiner.EARLIEST);
3044+
3045+
WatermarkHoldState bag = underTest.state(NAMESPACE, addr);
3046+
3047+
bag.setKnownEmpty();
3048+
assertThat(bag.read(), Matchers.nullValue());
3049+
3050+
bag.add(new Instant(300));
3051+
assertThat(bag.read(), Matchers.equalTo(new Instant(300)));
3052+
3053+
// Shouldn't need to read from windmill because the value is already available.
3054+
Mockito.verifyNoMoreInteractions(mockReader);
3055+
}
3056+
3057+
@Test
3058+
public void testWatermarkSetKnownEmptyPersist() throws Exception {
3059+
StateTag<WatermarkHoldState> addr =
3060+
StateTags.watermarkStateInternal("watermark", TimestampCombiner.EARLIEST);
3061+
3062+
WatermarkHoldState bag = underTest.state(NAMESPACE, addr);
3063+
3064+
bag.add(new Instant(1000));
3065+
bag.setKnownEmpty();
3066+
3067+
Windmill.WorkItemCommitRequest.Builder commitBuilder =
3068+
Windmill.WorkItemCommitRequest.newBuilder();
3069+
underTest.persist(commitBuilder);
3070+
3071+
// Should be a no-op, no reset, no adds.
3072+
assertEquals(0, commitBuilder.getWatermarkHoldsCount());
3073+
3074+
Mockito.verifyNoMoreInteractions(mockReader);
3075+
}
3076+
30403077
@Test
30413078
public void testWatermarkPersistEarliest() throws Exception {
30423079
StateTag<WatermarkHoldState> addr =

sdks/java/core/src/main/java/org/apache/beam/sdk/state/WatermarkHoldState.java

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -38,4 +38,14 @@ public interface WatermarkHoldState extends GroupingState<Instant, Instant> {
3838

3939
@Override
4040
WatermarkHoldState readLater();
41+
42+
/**
43+
* <b><i>For internal use only; no backwards-compatibility guarantees.</i></b>
44+
*
45+
* <p>Permit marking the state as empty locally, without necessarily clearing it in the backend.
46+
*
47+
* <p>This may be used by runners to optimize out unnecessary state reads.
48+
*/
49+
@Internal
50+
default void setKnownEmpty() {}
4151
}

0 commit comments

Comments
 (0)