Skip to content

Commit 32fc081

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

9 files changed

Lines changed: 174 additions & 2 deletions

File tree

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

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -272,6 +272,11 @@ public void clear() {
272272
combinedHold = null;
273273
}
274274

275+
@Override
276+
public void setKnownEmpty() {
277+
combinedHold = null;
278+
}
279+
275280
@Override
276281
public Instant read() {
277282
return combinedHold;

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: 28 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -99,6 +99,16 @@ public boolean isClosed(StateAccessor<?> state) {
9999
return readFinishedBits(state.access(FINISHED_BITS_TAG)).isFinished(rootTrigger);
100100
}
101101

102+
/** Return true if the window is new (no trigger state has ever been persisted). */
103+
public boolean isNew(StateAccessor<?> state) {
104+
return isFinishedSetNeeded() && state.access(FINISHED_BITS_TAG).read() == null;
105+
}
106+
107+
@VisibleForTesting
108+
public BitSet getFinishedBits(StateAccessor<?> state) {
109+
return readFinishedBits(state.access(FINISHED_BITS_TAG)).getBitSet();
110+
}
111+
102112
public void prefetchIsClosed(StateAccessor<?> state) {
103113
if (isFinishedSetNeeded()) {
104114
state.access(FINISHED_BITS_TAG).readLater();
@@ -187,15 +197,31 @@ private void persistFinishedSet(
187197
}
188198

189199
ValueState<BitSet> finishedSetState = state.access(FINISHED_BITS_TAG);
190-
if (!readFinishedBits(finishedSetState).equals(modifiedFinishedSet)) {
200+
@Nullable BitSet currentBits = finishedSetState.read();
201+
if (currentBits == null || !isEquivalent(currentBits, modifiedFinishedSet.getBitSet())) {
191202
if (modifiedFinishedSet.getBitSet().isEmpty()) {
192-
finishedSetState.clear();
203+
// To distinguish between a "new" window and a "seen but empty" window, we
204+
// write a BitSet with a sentinel bit at an index that will never be used by
205+
// any trigger in the tree.
206+
BitSet sentinel = new BitSet();
207+
sentinel.set(rootTrigger.getFirstIndexAfterSubtree());
208+
finishedSetState.write(sentinel);
193209
} else {
194210
finishedSetState.write(modifiedFinishedSet.getBitSet());
195211
}
196212
}
197213
}
198214

215+
private boolean isEquivalent(BitSet currentBits, BitSet modifiedBits) {
216+
if (currentBits.equals(modifiedBits)) {
217+
return true;
218+
}
219+
// They might only differ by the sentinel bit.
220+
BitSet currentBitsCopy = (BitSet) currentBits.clone();
221+
currentBitsCopy.clear(rootTrigger.getFirstIndexAfterSubtree());
222+
return currentBitsCopy.equals(modifiedBits);
223+
}
224+
199225
/** Clear the finished bits. */
200226
public void clearFinished(StateAccessor<?> state) {
201227
clearFinishedBits(state.access(FINISHED_BITS_TAG));

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

Lines changed: 39 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,41 @@ 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 sentinel bit to be set.
2373+
assertTrue("Sentinel bit should be set", bitSet.get(1));
2374+
// And trigger not finished.
2375+
assertFalse("Trigger should not be finished", bitSet.get(0));
2376+
2377+
// 2. Second element for the same window.
2378+
// We want to verify it doesn't clear the first element.
2379+
tester.injectElements(TimestampedValue.of(2, new Instant(2)));
2380+
2381+
// Extract output.
2382+
List<WindowedValue<Iterable<Integer>>> output = tester.extractOutput();
2383+
assertThat(output, contains(isSingleWindowedValue(containsInAnyOrder(1, 2), 9, 0, 10)));
2384+
}
23462385
}

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

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -77,6 +77,12 @@ public void clear() {
7777
localAdditions = null;
7878
}
7979

80+
@Override
81+
public void setKnownEmpty() {
82+
cachedValue = Optional.absent();
83+
localAdditions = null;
84+
}
85+
8086
@Override
8187
@SuppressWarnings("FutureReturnValueIgnored")
8288
public WindmillWatermarkHold readLater() {

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)