From 28957f5ca956ab85c4f64e13d51f66096470dfd3 Mon Sep 17 00:00:00 2001 From: Arun Pandian Date: Thu, 12 Feb 2026 01:54:55 +0000 Subject: [PATCH 01/20] Use trigger state to know if a window is new --- .../beam/runners/core/ReduceFnRunner.java | 18 +++++++ .../beam/runners/core/WatermarkHold.java | 17 +++++++ .../core/triggers/FinishedTriggersBitSet.java | 14 ++++++ .../triggers/TriggerStateMachineRunner.java | 27 +++++++---- .../beam/runners/core/ReduceFnRunnerTest.java | 48 +++++++++++++++++++ .../windmill/state/WindmillWatermarkHold.java | 15 +++++- .../state/WindmillStateInternalsTest.java | 37 ++++++++++++++ .../beam/sdk/state/WatermarkHoldState.java | 10 ++++ 8 files changed, 176 insertions(+), 10 deletions(-) diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java index b08bd42b0b22..7493bc78dce3 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java @@ -274,6 +274,16 @@ boolean hasNoActiveWindows() { return activeWindows.getActiveAndNewWindows().isEmpty(); } + @VisibleForTesting + TriggerStateMachineRunner getTriggerRunner() { + return triggerRunner; + } + + @VisibleForTesting + ReduceFnContextFactory getContextFactory() { + return contextFactory; + } + private Set windowsThatAreOpen(Collection windows) { Set result = new HashSet<>(); for (W window : windows) { @@ -603,6 +613,14 @@ private void processElement(Map windowToMergeResult, WindowedValue contextFactory.forValue( window, value.getValue(), value.getTimestamp(), StateStyle.RENAMED); + if (triggerRunner.isNew(directContext.state())) { + // Blindly clear state to ensure Windmill doesn't do unnecessary reads. + reduceFn.clearState(renamedContext); + paneInfoTracker.clear(directContext.state()); + watermarkHold.setKnownEmpty(renamedContext); + nonEmptyPanes.clearPane(renamedContext.state()); + } + nonEmptyPanes.recordContent(renamedContext.state()); scheduleGarbageCollectionTimer(directContext); diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/WatermarkHold.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/WatermarkHold.java index 15ae8dfe5f1a..b9185ccfba3f 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/WatermarkHold.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/WatermarkHold.java @@ -466,6 +466,23 @@ public void clearHolds(ReduceFn.Context context) { context.state().access(EXTRA_HOLD_TAG).clear(); } + /** + * For internal use only; no backwards-compatibility guarantees. + * + *

Permit marking the watermark holds as empty locally, without necessarily clearing them in + * the backend. + */ + public void setKnownEmpty(ReduceFn.Context context) { + WindowTracing.debug( + "WatermarkHold.setKnownEmpty: For key:{}; window:{}; inputWatermark:{}; outputWatermark:{}", + context.key(), + context.window(), + timerInternals.currentInputWatermarkTime(), + timerInternals.currentOutputWatermarkTime()); + context.state().access(elementHoldTag).setKnownEmpty(); + context.state().access(EXTRA_HOLD_TAG).setKnownEmpty(); + } + /** Return the current data hold, or null if none. Does not clear. For debugging only. */ public @Nullable Instant getDataCurrent(ReduceFn.Context context) { return context.state().access(elementHoldTag).read(); diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/FinishedTriggersBitSet.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/FinishedTriggersBitSet.java index 7eebb4474c6c..967e1ef43f08 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/FinishedTriggersBitSet.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/FinishedTriggersBitSet.java @@ -18,6 +18,7 @@ package org.apache.beam.runners.core.triggers; import java.util.BitSet; +import org.checkerframework.checker.nullness.qual.Nullable; /** A {@link FinishedTriggers} implementation based on an underlying {@link BitSet}. */ public class FinishedTriggersBitSet implements FinishedTriggers { @@ -60,4 +61,17 @@ public void clearRecursively(ExecutableTriggerStateMachine trigger) { public FinishedTriggersBitSet copy() { return new FinishedTriggersBitSet((BitSet) bitSet.clone()); } + + @Override + public boolean equals(@Nullable Object obj) { + if (!(obj instanceof FinishedTriggersBitSet)) { + return false; + } + return bitSet.equals(((FinishedTriggersBitSet) obj).bitSet); + } + + @Override + public int hashCode() { + return bitSet.hashCode(); + } } diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/TriggerStateMachineRunner.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/TriggerStateMachineRunner.java index cf29646ebaa3..db93b3b12d52 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/TriggerStateMachineRunner.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/TriggerStateMachineRunner.java @@ -81,9 +81,11 @@ private FinishedTriggersBitSet readFinishedBits(ValueState state) { } @Nullable BitSet bitSet = state.read(); - return bitSet == null - ? FinishedTriggersBitSet.emptyWithCapacity(rootTrigger.getFirstIndexAfterSubtree()) - : FinishedTriggersBitSet.fromBitSet(bitSet); + if (bitSet == null) { + return FinishedTriggersBitSet.emptyWithCapacity(rootTrigger.getFirstIndexAfterSubtree()); + } + + return FinishedTriggersBitSet.fromBitSet(bitSet); } private void clearFinishedBits(ValueState state) { @@ -99,6 +101,16 @@ public boolean isClosed(StateAccessor state) { return readFinishedBits(state.access(FINISHED_BITS_TAG)).isFinished(rootTrigger); } + /** Return true if the window is new (no trigger state has ever been persisted). */ + public boolean isNew(StateAccessor state) { + return isFinishedSetNeeded() && state.access(FINISHED_BITS_TAG).read() == null; + } + + @VisibleForTesting + public BitSet getFinishedBits(StateAccessor state) { + return readFinishedBits(state.access(FINISHED_BITS_TAG)).getBitSet(); + } + public void prefetchIsClosed(StateAccessor state) { if (isFinishedSetNeeded()) { state.access(FINISHED_BITS_TAG).readLater(); @@ -187,12 +199,9 @@ private void persistFinishedSet( } ValueState finishedSetState = state.access(FINISHED_BITS_TAG); - if (!readFinishedBits(finishedSetState).equals(modifiedFinishedSet)) { - if (modifiedFinishedSet.getBitSet().isEmpty()) { - finishedSetState.clear(); - } else { - finishedSetState.write(modifiedFinishedSet.getBitSet()); - } + @Nullable BitSet currentBits = finishedSetState.read(); + if (currentBits == null || !currentBits.equals(modifiedFinishedSet.getBitSet())) { + finishedSetState.write(modifiedFinishedSet.getBitSet()); } } diff --git a/runners/core-java/src/test/java/org/apache/beam/runners/core/ReduceFnRunnerTest.java b/runners/core-java/src/test/java/org/apache/beam/runners/core/ReduceFnRunnerTest.java index 85f6573be23e..2231a3fa1450 100644 --- a/runners/core-java/src/test/java/org/apache/beam/runners/core/ReduceFnRunnerTest.java +++ b/runners/core-java/src/test/java/org/apache/beam/runners/core/ReduceFnRunnerTest.java @@ -40,9 +40,11 @@ import java.util.ArrayList; import java.util.Arrays; +import java.util.BitSet; import java.util.List; import java.util.Random; import java.util.concurrent.ThreadLocalRandom; +import org.apache.beam.runners.core.ReduceFnContextFactory.StateStyle; import org.apache.beam.runners.core.metrics.MetricsContainerImpl; import org.apache.beam.runners.core.triggers.DefaultTriggerStateMachine; import org.apache.beam.runners.core.triggers.TriggerStateMachine; @@ -2343,4 +2345,50 @@ public interface TestOptions extends PipelineOptions { void setValue(int value); } + + @Test + public void testNewWindowOptimization() throws Exception { + WindowingStrategy strategy = + WindowingStrategy.of(FixedWindows.of(Duration.millis(10))) + .withTrigger(AfterPane.elementCountAtLeast(2)) + .withMode(AccumulationMode.ACCUMULATING_FIRED_PANES); + + ReduceFnTester, IntervalWindow> tester = + ReduceFnTester.nonCombining(strategy); + + IntervalWindow window = new IntervalWindow(new Instant(0), new Instant(10)); + + // 1. First element for a new window. + tester.injectElements(TimestampedValue.of(1, new Instant(1))); + + // Verify sentinel bit is written. + BitSet bitSet = + tester + .createRunner() + .getTriggerRunner() + .getFinishedBits( + tester.createRunner().getContextFactory().base(window, StateStyle.DIRECT).state()); + + // We expect the bitset to be empty (the sentinel bit is no longer used). + assertTrue("Bitset should be empty", bitSet.isEmpty()); + // And trigger not finished. + assertFalse("Trigger should not be finished", bitSet.get(0)); + + // And verify that it is no longer "new". + assertFalse( + "Window should no longer be new", + tester + .createRunner() + .getTriggerRunner() + .isNew( + tester.createRunner().getContextFactory().base(window, StateStyle.DIRECT).state())); + + // 2. Second element for the same window. + // We want to verify it doesn't clear the first element. + tester.injectElements(TimestampedValue.of(2, new Instant(2))); + + // Extract output. + List>> output = tester.extractOutput(); + assertThat(output, contains(isSingleWindowedValue(containsInAnyOrder(1, 2), 9, 0, 10))); + } } diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillWatermarkHold.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillWatermarkHold.java index 613d87c127b7..1e9778f0e4e4 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillWatermarkHold.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillWatermarkHold.java @@ -37,6 +37,7 @@ "nullness" // TODO(https://github.com/apache/beam/issues/20497) }) public class WindmillWatermarkHold extends WindmillState implements WatermarkHoldState { + // The encoded size of an Instant. private static final int ENCODED_SIZE = 8; @@ -46,6 +47,7 @@ public class WindmillWatermarkHold extends WindmillState implements WatermarkHol private final String stateFamily; private boolean cleared = false; + private boolean knownEmpty = false; /** * If non-{@literal null}, the known current hold value, or absent if we know there are no output * watermark holds. If {@literal null}, the current hold value could depend on holds in Windmill @@ -77,6 +79,13 @@ public void clear() { localAdditions = null; } + @Override + public void setKnownEmpty() { + cachedValue = Optional.absent(); + localAdditions = null; + knownEmpty = true; + } + @Override @SuppressWarnings("FutureReturnValueIgnored") public WindmillWatermarkHold readLater() { @@ -133,7 +142,7 @@ public Future persist( Future result; - if (!cleared && localAdditions == null) { + if (!knownEmpty && !cleared && localAdditions == null) { // No changes, so no need to update Windmill and no need to cache any value. return Futures.immediateFuture(Windmill.WorkItemCommitRequest.newBuilder().buildPartial()); } @@ -166,15 +175,19 @@ public Future persist( } else if (!cleared && localAdditions != null) { // Otherwise, we need to combine the local additions with the already persisted data result = combineWithPersisted(); + } else if (knownEmpty) { + result = Futures.immediateFuture(Windmill.WorkItemCommitRequest.newBuilder().buildPartial()); } else { throw new IllegalStateException("Unreachable condition"); } final int estimatedByteSize = ENCODED_SIZE + stateKey.byteString().size(); + return Futures.lazyTransform( result, result1 -> { cleared = false; + knownEmpty = false; localAdditions = null; if (cachedValue != null) { cache.put(namespace, stateKey, WindmillWatermarkHold.this, estimatedByteSize); diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateInternalsTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateInternalsTest.java index 7a06d3a29493..d368c52e8743 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateInternalsTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateInternalsTest.java @@ -3037,6 +3037,43 @@ public void testWatermarkClearBeforeRead() throws Exception { Mockito.verifyNoMoreInteractions(mockReader); } + @Test + public void testWatermarkSetKnownEmptyBeforeRead() throws Exception { + StateTag addr = + StateTags.watermarkStateInternal("watermark", TimestampCombiner.EARLIEST); + + WatermarkHoldState bag = underTest.state(NAMESPACE, addr); + + bag.setKnownEmpty(); + assertThat(bag.read(), Matchers.nullValue()); + + bag.add(new Instant(300)); + assertThat(bag.read(), Matchers.equalTo(new Instant(300))); + + // Shouldn't need to read from windmill because the value is already available. + Mockito.verifyNoMoreInteractions(mockReader); + } + + @Test + public void testWatermarkSetKnownEmptyPersist() throws Exception { + StateTag addr = + StateTags.watermarkStateInternal("watermark", TimestampCombiner.EARLIEST); + + WatermarkHoldState bag = underTest.state(NAMESPACE, addr); + + bag.add(new Instant(1000)); + bag.setKnownEmpty(); + + Windmill.WorkItemCommitRequest.Builder commitBuilder = + Windmill.WorkItemCommitRequest.newBuilder(); + underTest.persist(commitBuilder); + + // Should be a no-op, no reset, no adds. + assertEquals(0, commitBuilder.getWatermarkHoldsCount()); + + Mockito.verifyNoMoreInteractions(mockReader); + } + @Test public void testWatermarkPersistEarliest() throws Exception { StateTag addr = diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/state/WatermarkHoldState.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/state/WatermarkHoldState.java index 6d4183da101f..f8b09bcf03a6 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/state/WatermarkHoldState.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/state/WatermarkHoldState.java @@ -38,4 +38,14 @@ public interface WatermarkHoldState extends GroupingState { @Override WatermarkHoldState readLater(); + + /** + * For internal use only; no backwards-compatibility guarantees. + * + *

Permit marking the state as empty locally, without necessarily clearing it in the backend. + * + *

This may be used by runners to optimize out unnecessary state reads. + */ + @Internal + default void setKnownEmpty() {} } From b24494ff8471782a98c67e429bef5ca44ea18dae Mon Sep 17 00:00:00 2001 From: Arun Pandian Date: Fri, 10 Apr 2026 04:42:47 +0000 Subject: [PATCH 02/20] Change TriggerState finished bitset coder to a SentinelBitSetCoder SentinelBitSetCoder is same as BitSetCoder except that it encodes empty bitset as a single element 0 byte array. This allows checking if the finished bitset is empty or missing. SentinelBitSetCoder and BitSetCoder are state compatible. Both coders can decode encoded bytes from the other coder successfully. --- CHANGES.md | 4 + .../serialization/SentinelBitSetCoder.java | 78 +++++++++ .../triggers/TriggerStateMachineRunner.java | 4 +- .../SentinelBitSetCoderTest.java | 150 ++++++++++++++++++ 4 files changed, 234 insertions(+), 2 deletions(-) create mode 100644 runners/core-java/src/main/java/org/apache/beam/runners/core/serialization/SentinelBitSetCoder.java create mode 100644 runners/core-java/src/test/java/org/apache/beam/runners/core/serialization/SentinelBitSetCoderTest.java diff --git a/CHANGES.md b/CHANGES.md index 319520f94309..58cc8d16d088 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -69,6 +69,10 @@ ## New Features / Improvements * X feature added (Java/Python) ([#X](https://github.com/apache/beam/issues/X)). +* TriggerStateMachineRunner changes from BitSetCoder to SentinelBitSetCoder to + encode finished bitset [#38139](https://github.com/apache/beam/pull/38139). + SentinelBitSetCoder and BitSetCoder are state compatible. Both coders can + decode encoded bytes from the other coder. ## Breaking Changes diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/serialization/SentinelBitSetCoder.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/serialization/SentinelBitSetCoder.java new file mode 100644 index 000000000000..e9f0582c5ca9 --- /dev/null +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/serialization/SentinelBitSetCoder.java @@ -0,0 +1,78 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.beam.runners.core.serialization; + +import java.io.IOException; +import java.io.InputStream; +import java.io.OutputStream; +import java.util.BitSet; +import org.apache.beam.sdk.coders.AtomicCoder; +import org.apache.beam.sdk.coders.ByteArrayCoder; +import org.apache.beam.sdk.coders.CoderException; + +/** + * Coder for {@link BitSet} that stores an empty bit set as a byte array with a single 0 element. + */ +public class SentinelBitSetCoder extends AtomicCoder { + private static final SentinelBitSetCoder INSTANCE = new SentinelBitSetCoder(); + private static final ByteArrayCoder BYTE_ARRAY_CODER = ByteArrayCoder.of(); + + private SentinelBitSetCoder() {} + + public static SentinelBitSetCoder of() { + return INSTANCE; + } + + @Override + public void encode(BitSet value, OutputStream outStream) throws CoderException, IOException { + encode(value, outStream, Context.NESTED); + } + + @Override + public void encode(BitSet value, OutputStream outStream, Context context) + throws CoderException, IOException { + if (value == null) { + throw new CoderException("cannot encode a null BitSet"); + } + byte[] bytes = value.isEmpty() ? new byte[] {0} : value.toByteArray(); + BYTE_ARRAY_CODER.encodeAndOwn(bytes, outStream, context); + } + + @Override + public BitSet decode(InputStream inStream) throws CoderException, IOException { + return decode(inStream, Context.NESTED); + } + + @Override + public BitSet decode(InputStream inStream, Context context) throws CoderException, IOException { + return BitSet.valueOf(BYTE_ARRAY_CODER.decode(inStream, context)); + } + + @Override + public void verifyDeterministic() throws NonDeterministicException { + verifyDeterministic( + this, + "SentinelBitSetCoder requires its ByteArrayCoder to be deterministic.", + BYTE_ARRAY_CODER); + } + + @Override + public boolean consistentWithEquals() { + return true; + } +} diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/TriggerStateMachineRunner.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/TriggerStateMachineRunner.java index cf29646ebaa3..e3791821b728 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/TriggerStateMachineRunner.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/TriggerStateMachineRunner.java @@ -26,7 +26,7 @@ import org.apache.beam.runners.core.StateAccessor; import org.apache.beam.runners.core.StateTag; import org.apache.beam.runners.core.StateTags; -import org.apache.beam.sdk.coders.BitSetCoder; +import org.apache.beam.runners.core.serialization.SentinelBitSetCoder; import org.apache.beam.sdk.state.Timers; import org.apache.beam.sdk.state.ValueState; import org.apache.beam.sdk.transforms.windowing.BoundedWindow; @@ -59,7 +59,7 @@ public class TriggerStateMachineRunner { @VisibleForTesting public static final StateTag> FINISHED_BITS_TAG = - StateTags.makeSystemTagInternal(StateTags.value("closed", BitSetCoder.of())); + StateTags.makeSystemTagInternal(StateTags.value("closed", SentinelBitSetCoder.of())); private final ExecutableTriggerStateMachine rootTrigger; private final TriggerStateMachineContextFactory contextFactory; diff --git a/runners/core-java/src/test/java/org/apache/beam/runners/core/serialization/SentinelBitSetCoderTest.java b/runners/core-java/src/test/java/org/apache/beam/runners/core/serialization/SentinelBitSetCoderTest.java new file mode 100644 index 000000000000..91d7d369f70b --- /dev/null +++ b/runners/core-java/src/test/java/org/apache/beam/runners/core/serialization/SentinelBitSetCoderTest.java @@ -0,0 +1,150 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.beam.runners.core.serialization; + +import static org.hamcrest.MatcherAssert.assertThat; +import static org.hamcrest.Matchers.equalTo; + +import java.util.Arrays; +import java.util.BitSet; +import java.util.List; +import org.apache.beam.sdk.coders.BitSetCoder; +import org.apache.beam.sdk.coders.Coder; +import org.apache.beam.sdk.coders.Coder.Context; +import org.apache.beam.sdk.coders.CoderException; +import org.apache.beam.sdk.testing.CoderProperties; +import org.apache.beam.sdk.util.CoderUtils; +import org.apache.beam.sdk.values.TypeDescriptor; +import org.junit.Rule; +import org.junit.Test; +import org.junit.rules.ExpectedException; +import org.junit.runner.RunWith; +import org.junit.runners.JUnit4; + +/** Tests for {@link SentinelBitSetCoder}. */ +@RunWith(JUnit4.class) +public class SentinelBitSetCoderTest { + + private static final Coder TEST_CODER = SentinelBitSetCoder.of(); + + private static final List TEST_VALUES = + Arrays.asList( + BitSet.valueOf(new byte[] {0xa, 0xb, 0xc}), + BitSet.valueOf(new byte[] {0xd, 0x3}), + BitSet.valueOf(new byte[] {0xd, 0xe}), + BitSet.valueOf(new byte[] {0}), + BitSet.valueOf(new byte[] {})); + + @Test + public void testDecodeEncodeEquals() throws Exception { + for (BitSet value : TEST_VALUES) { + CoderProperties.coderDecodeEncodeEqual(TEST_CODER, value); + } + } + + @Test + public void testRegisterByteSizeObserver() throws Exception { + CoderProperties.testByteCount( + SentinelBitSetCoder.of(), Coder.Context.OUTER, TEST_VALUES.toArray(new BitSet[] {})); + + CoderProperties.testByteCount( + SentinelBitSetCoder.of(), Coder.Context.NESTED, TEST_VALUES.toArray(new BitSet[] {})); + } + + @Test + public void testStructuralValueConsistentWithEquals() throws Exception { + for (BitSet value1 : TEST_VALUES) { + for (BitSet value2 : TEST_VALUES) { + CoderProperties.structuralValueConsistentWithEquals(TEST_CODER, value1, value2); + } + } + } + + /** + * Generated data to check that the wire format has not changed. "CgsM" is {0xa, 0xb, 0xc} "DQM" + * is {0xd, 0x3} "DQ4" is {0xd, 0xe} "AA==" is {0} (Sentinel for empty BitSet) + */ + private static final List TEST_ENCODINGS = + Arrays.asList("CgsM", "DQM", "DQ4", "AA", "AA"); + + @Test + public void testWireFormatEncode() throws Exception { + CoderProperties.coderEncodesBase64(TEST_CODER, TEST_VALUES, TEST_ENCODINGS); + } + + @Rule public ExpectedException thrown = ExpectedException.none(); + + @Test + public void encodeNullThrowsCoderException() throws Exception { + thrown.expect(CoderException.class); + thrown.expectMessage("cannot encode a null BitSet"); + + CoderUtils.encodeToBase64(TEST_CODER, null); + } + + @Test + public void testEncodedTypeDescriptor() throws Exception { + assertThat(TEST_CODER.getEncodedTypeDescriptor(), equalTo(TypeDescriptor.of(BitSet.class))); + } + + @Test + public void testEmptyBitSetEncoding() throws Exception { + { + byte[] encoded = CoderUtils.encodeToByteArray(TEST_CODER, new BitSet()); + // ByteArrayCoder in OUTER context encodes as is. + assertThat(encoded, equalTo(new byte[] {0})); + } + { + byte[] encoded = CoderUtils.encodeToByteArray(TEST_CODER, new BitSet(), Context.NESTED); + // Varint length = 1, data = 1 + assertThat(encoded, equalTo(new byte[] {1, 0})); + } + } + + @Test + public void testCompatibilityWithBitSetCoder() throws Exception { + BitSetCoder bitSetCoder = BitSetCoder.of(); + SentinelBitSetCoder sentinelCoder = SentinelBitSetCoder.of(); + + for (BitSet bitset : TEST_VALUES) { + for (Coder.Context context : Arrays.asList(Coder.Context.OUTER, Coder.Context.NESTED)) { + // Test SentinelBitSetCoder can decode bytes encoded by BitSetCoder + { + byte[] encodedByBitSet = CoderUtils.encodeToByteArray(bitSetCoder, bitset, context); + BitSet decodedBySentinel = + CoderUtils.decodeFromByteArray(sentinelCoder, encodedByBitSet, context); + assertThat( + "Decoding BitSetCoder encoded value with context " + context, + decodedBySentinel, + equalTo(bitset)); + } + + // Test BitSetCoder can decode bytes encoded by SentinelBitSetCoder + { + byte[] encodedBySentinel = CoderUtils.encodeToByteArray(sentinelCoder, bitset, context); + BitSet decodedByBitSet = + CoderUtils.decodeFromByteArray(bitSetCoder, encodedBySentinel, context); + assertThat( + "Decoding SentinelBitSetCoder encoded value with context " + context, + decodedByBitSet, + equalTo(bitset)); + } + } + } + } +} From a4ff2c6ef010313064620fa9e43e121ec26d2e17 Mon Sep 17 00:00:00 2001 From: Arun Pandian Date: Fri, 10 Apr 2026 05:06:24 +0000 Subject: [PATCH 03/20] fix style --- CHANGES.md | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/CHANGES.md b/CHANGES.md index 58cc8d16d088..acdce4928d72 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -99,6 +99,7 @@ ## Highlights + ## I/Os * DebeziumIO (Java): added `OffsetRetainer` interface and `FileSystemOffsetRetainer` implementation to persist and restore CDC offsets across pipeline restarts, and exposed `withStartOffset` / `withOffsetRetainer` on `DebeziumIO.Read` and the cross-language `ReadBuilder` ([#28248](https://github.com/apache/beam/issues/28248)). @@ -2422,4 +2423,4 @@ Schema Options, it will be removed in version `2.23.0`. ([BEAM-9704](https://iss ## Highlights -- For versions 2.19.0 and older release notes are available on [Apache Beam Blog](https://beam.apache.org/blog/). +- For versions 2.19.0 and older release notes are available on [Apache Beam Blog](https://beam.apache.org/blog/). \ No newline at end of file From c55b4cbb8691e61f88f150bedd048d0e73a03ed1 Mon Sep 17 00:00:00 2001 From: Arun Pandian Date: Fri, 10 Apr 2026 05:29:46 +0000 Subject: [PATCH 04/20] fix style --- CHANGES.md | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/CHANGES.md b/CHANGES.md index acdce4928d72..58cc8d16d088 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -99,7 +99,6 @@ ## Highlights - ## I/Os * DebeziumIO (Java): added `OffsetRetainer` interface and `FileSystemOffsetRetainer` implementation to persist and restore CDC offsets across pipeline restarts, and exposed `withStartOffset` / `withOffsetRetainer` on `DebeziumIO.Read` and the cross-language `ReadBuilder` ([#28248](https://github.com/apache/beam/issues/28248)). @@ -2423,4 +2422,4 @@ Schema Options, it will be removed in version `2.23.0`. ([BEAM-9704](https://iss ## Highlights -- For versions 2.19.0 and older release notes are available on [Apache Beam Blog](https://beam.apache.org/blog/). \ No newline at end of file +- For versions 2.19.0 and older release notes are available on [Apache Beam Blog](https://beam.apache.org/blog/). From cf5050ea348ffb7f8270205d3ea2744ee139c82a Mon Sep 17 00:00:00 2001 From: Arun Pandian Date: Fri, 10 Apr 2026 05:47:33 +0000 Subject: [PATCH 05/20] fix style --- CHANGES.md | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/CHANGES.md b/CHANGES.md index 58cc8d16d088..ca76cc9d6a8b 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -70,9 +70,9 @@ * X feature added (Java/Python) ([#X](https://github.com/apache/beam/issues/X)). * TriggerStateMachineRunner changes from BitSetCoder to SentinelBitSetCoder to - encode finished bitset [#38139](https://github.com/apache/beam/pull/38139). - SentinelBitSetCoder and BitSetCoder are state compatible. Both coders can - decode encoded bytes from the other coder. + encode finished bitset. SentinelBitSetCoder and BitSetCoder are state + compatible. Both coders can decode encoded bytes from the other coder + [#38139](https://github.com/apache/beam/pull/38139). ## Breaking Changes From 848617b5e199734c951bba451a8dcf6790e7d7ee Mon Sep 17 00:00:00 2001 From: Arun Pandian Date: Fri, 10 Apr 2026 07:26:03 +0000 Subject: [PATCH 06/20] fix style --- CHANGES.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/CHANGES.md b/CHANGES.md index ca76cc9d6a8b..44c4620229b7 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -72,7 +72,7 @@ * TriggerStateMachineRunner changes from BitSetCoder to SentinelBitSetCoder to encode finished bitset. SentinelBitSetCoder and BitSetCoder are state compatible. Both coders can decode encoded bytes from the other coder - [#38139](https://github.com/apache/beam/pull/38139). + [#38139](https://github.com/apache/beam/issues/38139). ## Breaking Changes From 493fa8f61a174aea96da079a56a7a156de0acbd5 Mon Sep 17 00:00:00 2001 From: Arun Pandian Date: Fri, 10 Apr 2026 08:00:12 +0000 Subject: [PATCH 07/20] fix style --- CHANGES.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/CHANGES.md b/CHANGES.md index 44c4620229b7..aa9a49a16e68 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -72,7 +72,7 @@ * TriggerStateMachineRunner changes from BitSetCoder to SentinelBitSetCoder to encode finished bitset. SentinelBitSetCoder and BitSetCoder are state compatible. Both coders can decode encoded bytes from the other coder - [#38139](https://github.com/apache/beam/issues/38139). + ([#38139](https://github.com/apache/beam/issues/38139)). ## Breaking Changes From 0d929f146fac101a38adc1582f55463109843817 Mon Sep 17 00:00:00 2001 From: Arun Pandian Date: Sat, 11 Apr 2026 09:01:36 +0000 Subject: [PATCH 08/20] address comment --- .../beam/runners/core/serialization/SentinelBitSetCoder.java | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/serialization/SentinelBitSetCoder.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/serialization/SentinelBitSetCoder.java index e9f0582c5ca9..340816f6e0e1 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/serialization/SentinelBitSetCoder.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/serialization/SentinelBitSetCoder.java @@ -26,9 +26,12 @@ import org.apache.beam.sdk.coders.CoderException; /** - * Coder for {@link BitSet} that stores an empty bit set as a byte array with a single 0 element. + * Coder for {@link BitSet} that stores an empty bit set as a byte array with a single 0 element. In + * general BitSetCoder should be preferred as it encodes an empty bit set as an empty byte array. + * However, there are cases where non-empty values are useful to indicate presence. */ public class SentinelBitSetCoder extends AtomicCoder { + private static final SentinelBitSetCoder INSTANCE = new SentinelBitSetCoder(); private static final ByteArrayCoder BYTE_ARRAY_CODER = ByteArrayCoder.of(); From 55b300c68aac3d2c1bfa48fb47a0fd43ff940ad3 Mon Sep 17 00:00:00 2001 From: Arun Pandian Date: Mon, 27 Apr 2026 21:57:19 +0000 Subject: [PATCH 09/20] add tests and improve setKnownEmpty --- .../beam/runners/core/ReduceFnRunner.java | 20 ++- .../beam/runners/core/ReduceFnRunnerTest.java | 24 ++- .../beam/runners/core/ReduceFnTester.java | 13 ++ .../windmill/state/WindmillWatermarkHold.java | 93 ++++++----- .../state/WindmillStateInternalsTest.java | 152 ++++++++++++++++-- 5 files changed, 241 insertions(+), 61 deletions(-) diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java index ea15ee8eca68..9a9a096130cd 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java @@ -37,6 +37,7 @@ import org.apache.beam.runners.core.triggers.TriggerStateMachineRunner; import org.apache.beam.sdk.metrics.Counter; import org.apache.beam.sdk.metrics.Metrics; +import org.apache.beam.sdk.options.ExperimentalOptions; import org.apache.beam.sdk.options.PipelineOptions; import org.apache.beam.sdk.state.TimeDomain; import org.apache.beam.sdk.transforms.windowing.BoundedWindow; @@ -211,6 +212,9 @@ public class ReduceFnRunner { */ private final NonEmptyPanes nonEmptyPanes; + private final boolean useNewWindowOptimization; + private final boolean disableWatermarkKnownEmptyOptimization; + public ReduceFnRunner( K key, WindowingStrategy windowingStrategy, @@ -237,6 +241,16 @@ public ReduceFnRunner( this.nonEmptyPanes = NonEmptyPanes.create(this.windowingStrategy, this.reduceFn); + this.useNewWindowOptimization = + options != null + && ExperimentalOptions.hasExperiment( + options, "unstable_not_update_compatible_new_window_optimization"); + + this.disableWatermarkKnownEmptyOptimization = + options != null + && ExperimentalOptions.hasExperiment( + options, "unstable_disable_watermark_known_empty_optimization"); + // Note this may incur I/O to load persisted window set data. this.activeWindows = createActiveWindowSet(); @@ -623,11 +637,13 @@ private void processElement(Map windowToMergeResult, WindowedValue StateStyle.RENAMED, value.causedByDrain()); - if (triggerRunner.isNew(directContext.state())) { + if (useNewWindowOptimization && triggerRunner.isNew(directContext.state())) { // Blindly clear state to ensure Windmill doesn't do unnecessary reads. reduceFn.clearState(renamedContext); paneInfoTracker.clear(directContext.state()); - watermarkHold.setKnownEmpty(renamedContext); + if (!disableWatermarkKnownEmptyOptimization) { + watermarkHold.setKnownEmpty(renamedContext); + } nonEmptyPanes.clearPane(renamedContext.state()); } diff --git a/runners/core-java/src/test/java/org/apache/beam/runners/core/ReduceFnRunnerTest.java b/runners/core-java/src/test/java/org/apache/beam/runners/core/ReduceFnRunnerTest.java index 2231a3fa1450..72b379a7c71b 100644 --- a/runners/core-java/src/test/java/org/apache/beam/runners/core/ReduceFnRunnerTest.java +++ b/runners/core-java/src/test/java/org/apache/beam/runners/core/ReduceFnRunnerTest.java @@ -41,6 +41,7 @@ import java.util.ArrayList; import java.util.Arrays; import java.util.BitSet; +import java.util.Collections; import java.util.List; import java.util.Random; import java.util.concurrent.ThreadLocalRandom; @@ -51,6 +52,7 @@ import org.apache.beam.sdk.coders.VarIntCoder; import org.apache.beam.sdk.metrics.MetricName; import org.apache.beam.sdk.metrics.MetricsEnvironment; +import org.apache.beam.sdk.options.ExperimentalOptions; import org.apache.beam.sdk.options.PipelineOptions; import org.apache.beam.sdk.options.PipelineOptionsFactory; import org.apache.beam.sdk.state.TimeDomain; @@ -2353,28 +2355,36 @@ public void testNewWindowOptimization() throws Exception { .withTrigger(AfterPane.elementCountAtLeast(2)) .withMode(AccumulationMode.ACCUMULATING_FIRED_PANES); + PipelineOptions options = PipelineOptionsFactory.create(); + options + .as(ExperimentalOptions.class) + .setExperiments( + Collections.singletonList("unstable_not_update_compatible_new_window_optimization")); ReduceFnTester, IntervalWindow> tester = - ReduceFnTester.nonCombining(strategy); + ReduceFnTester.nonCombining(strategy, options); IntervalWindow window = new IntervalWindow(new Instant(0), new Instant(10)); + assertTrue( + "Window should be new", + tester + .createRunner() + .getTriggerRunner() + .isNew( + tester.createRunner().getContextFactory().base(window, StateStyle.DIRECT).state())); + // 1. First element for a new window. tester.injectElements(TimestampedValue.of(1, new Instant(1))); - // Verify sentinel bit is written. BitSet bitSet = tester .createRunner() .getTriggerRunner() .getFinishedBits( tester.createRunner().getContextFactory().base(window, StateStyle.DIRECT).state()); - - // We expect the bitset to be empty (the sentinel bit is no longer used). assertTrue("Bitset should be empty", bitSet.isEmpty()); - // And trigger not finished. assertFalse("Trigger should not be finished", bitSet.get(0)); - // And verify that it is no longer "new". assertFalse( "Window should no longer be new", tester @@ -2384,11 +2394,11 @@ public void testNewWindowOptimization() throws Exception { tester.createRunner().getContextFactory().base(window, StateStyle.DIRECT).state())); // 2. Second element for the same window. - // We want to verify it doesn't clear the first element. tester.injectElements(TimestampedValue.of(2, new Instant(2))); // Extract output. List>> output = tester.extractOutput(); + // 2 elements fired at end of window. assertThat(output, contains(isSingleWindowedValue(containsInAnyOrder(1, 2), 9, 0, 10))); } } diff --git a/runners/core-java/src/test/java/org/apache/beam/runners/core/ReduceFnTester.java b/runners/core-java/src/test/java/org/apache/beam/runners/core/ReduceFnTester.java index 43b6a3cb0cb0..d1f137bd4936 100644 --- a/runners/core-java/src/test/java/org/apache/beam/runners/core/ReduceFnTester.java +++ b/runners/core-java/src/test/java/org/apache/beam/runners/core/ReduceFnTester.java @@ -129,6 +129,19 @@ ReduceFnTester, W> nonCombining( NullSideInputReader.empty()); } + public static + ReduceFnTester, W> nonCombining( + WindowingStrategy windowingStrategy, PipelineOptions options) throws Exception { + return new ReduceFnTester<>( + windowingStrategy, + TriggerStateMachines.stateMachineForTrigger( + TriggerTranslation.toProto(windowingStrategy.getTrigger())), + SystemReduceFn.buffering(VarIntCoder.of()), + IterableCoder.of(VarIntCoder.of()), + options, + NullSideInputReader.empty()); + } + /** * Creates a {@link ReduceFnTester} for the given {@link WindowingStrategy} and {@link * TriggerStateMachine}, for mocking the interactions between {@link ReduceFnRunner} and the diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillWatermarkHold.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillWatermarkHold.java index 1e9778f0e4e4..2df2699d287a 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillWatermarkHold.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillWatermarkHold.java @@ -17,6 +17,8 @@ */ package org.apache.beam.runners.dataflow.worker.windmill.state; +import static org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkState; + import java.io.Closeable; import java.io.IOException; import java.util.concurrent.ExecutionException; @@ -81,8 +83,12 @@ public void clear() { @Override public void setKnownEmpty() { + checkState(localAdditions == null, "setKnownEmpty called with local additions"); + checkState(!cleared, "setKnownEmpty called after clearing"); + checkState( + cachedValue == null || !cachedValue.isPresent(), + "setKnownEmpty called with a cached value"); cachedValue = Optional.absent(); - localAdditions = null; knownEmpty = true; } @@ -142,43 +148,58 @@ public Future persist( Future result; - if (!knownEmpty && !cleared && localAdditions == null) { - // No changes, so no need to update Windmill and no need to cache any value. - return Futures.immediateFuture(Windmill.WorkItemCommitRequest.newBuilder().buildPartial()); - } - - if (cleared && localAdditions == null) { - // Just clearing the persisted state; blind delete - Windmill.WorkItemCommitRequest.Builder commitBuilder = - Windmill.WorkItemCommitRequest.newBuilder(); - commitBuilder - .addWatermarkHoldsBuilder() - .setTag(stateKey.byteString()) - .setStateFamily(stateFamily) - .setReset(true); - - result = Futures.immediateFuture(commitBuilder.buildPartial()); - } else if (cleared && localAdditions != null) { - // Since we cleared before adding, we can do a blind overwrite of persisted state - Windmill.WorkItemCommitRequest.Builder commitBuilder = - Windmill.WorkItemCommitRequest.newBuilder(); - commitBuilder - .addWatermarkHoldsBuilder() - .setTag(stateKey.byteString()) - .setStateFamily(stateFamily) - .setReset(true) - .addTimestamps(WindmillTimeUtils.harnessToWindmillTimestamp(localAdditions)); - - cachedValue = Optional.of(localAdditions); - - result = Futures.immediateFuture(commitBuilder.buildPartial()); - } else if (!cleared && localAdditions != null) { - // Otherwise, we need to combine the local additions with the already persisted data + if (knownEmpty) { + if (localAdditions != null) { + // 1. We know it's empty, so we can just update with localAdditions + Windmill.WorkItemCommitRequest.Builder commitBuilder = + Windmill.WorkItemCommitRequest.newBuilder(); + commitBuilder + .addWatermarkHoldsBuilder() + .setTag(stateKey.byteString()) + .setStateFamily(stateFamily) + .addTimestamps(WindmillTimeUtils.harnessToWindmillTimestamp(localAdditions)); + + cachedValue = Optional.of(localAdditions); + result = Futures.immediateFuture(commitBuilder.buildPartial()); + } else { + // 2. State is known to be empty and there are no local additions. + // Whether 'cleared' was called or not, the desired state is empty. + // So no need to update Windmill. + result = + Futures.immediateFuture(Windmill.WorkItemCommitRequest.newBuilder().buildPartial()); + } + } else if (cleared) { + if (localAdditions == null) { + // 3. Just clearing the persisted state; blind delete + Windmill.WorkItemCommitRequest.Builder commitBuilder = + Windmill.WorkItemCommitRequest.newBuilder(); + commitBuilder + .addWatermarkHoldsBuilder() + .setTag(stateKey.byteString()) + .setStateFamily(stateFamily) + .setReset(true); + + result = Futures.immediateFuture(commitBuilder.buildPartial()); + } else { + // 4. Since we cleared before adding, we can do an overwrite of persisted state + Windmill.WorkItemCommitRequest.Builder commitBuilder = + Windmill.WorkItemCommitRequest.newBuilder(); + commitBuilder + .addWatermarkHoldsBuilder() + .setTag(stateKey.byteString()) + .setStateFamily(stateFamily) + .setReset(true) + .addTimestamps(WindmillTimeUtils.harnessToWindmillTimestamp(localAdditions)); + + cachedValue = Optional.of(localAdditions); + result = Futures.immediateFuture(commitBuilder.buildPartial()); + } + } else if (localAdditions != null) { + // 5. Otherwise, we need to combine the local additions with the already persisted data result = combineWithPersisted(); - } else if (knownEmpty) { - result = Futures.immediateFuture(Windmill.WorkItemCommitRequest.newBuilder().buildPartial()); } else { - throw new IllegalStateException("Unreachable condition"); + // 6. No changes, so no need to update Windmill. + result = Futures.immediateFuture(Windmill.WorkItemCommitRequest.newBuilder().buildPartial()); } final int estimatedByteSize = ENCODED_SIZE + stateKey.byteString().size(); diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateInternalsTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateInternalsTest.java index 7c4201d2e748..5601041a48f5 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateInternalsTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateInternalsTest.java @@ -3025,13 +3025,13 @@ public void testWatermarkClearBeforeRead() throws Exception { StateTag addr = StateTags.watermarkStateInternal("watermark", TimestampCombiner.EARLIEST); - WatermarkHoldState bag = underTest.state(NAMESPACE, addr); + WatermarkHoldState hold = underTest.state(NAMESPACE, addr); - bag.clear(); - assertThat(bag.read(), Matchers.nullValue()); + hold.clear(); + assertThat(hold.read(), Matchers.nullValue()); - bag.add(new Instant(300)); - assertThat(bag.read(), Matchers.equalTo(new Instant(300))); + hold.add(new Instant(300)); + assertThat(hold.read(), Matchers.equalTo(new Instant(300))); // Shouldn't need to read from windmill because the value is already available. Mockito.verifyNoMoreInteractions(mockReader); @@ -3042,36 +3042,156 @@ public void testWatermarkSetKnownEmptyBeforeRead() throws Exception { StateTag addr = StateTags.watermarkStateInternal("watermark", TimestampCombiner.EARLIEST); - WatermarkHoldState bag = underTest.state(NAMESPACE, addr); + WatermarkHoldState hold = underTest.state(NAMESPACE, addr); - bag.setKnownEmpty(); - assertThat(bag.read(), Matchers.nullValue()); + hold.setKnownEmpty(); + assertThat(hold.read(), Matchers.nullValue()); - bag.add(new Instant(300)); - assertThat(bag.read(), Matchers.equalTo(new Instant(300))); + hold.add(new Instant(300)); + assertThat(hold.read(), Matchers.equalTo(new Instant(300))); // Shouldn't need to read from windmill because the value is already available. Mockito.verifyNoMoreInteractions(mockReader); } @Test - public void testWatermarkSetKnownEmptyPersist() throws Exception { + public void testWatermarkSetKnownEmptyThenAddPersist() throws Exception { StateTag addr = StateTags.watermarkStateInternal("watermark", TimestampCombiner.EARLIEST); - WatermarkHoldState bag = underTest.state(NAMESPACE, addr); + WatermarkHoldState hold = underTest.state(NAMESPACE, addr); - bag.add(new Instant(1000)); - bag.setKnownEmpty(); + hold.setKnownEmpty(); + hold.add(new Instant(1000)); + + Windmill.WorkItemCommitRequest.Builder commitBuilder = + Windmill.WorkItemCommitRequest.newBuilder(); + underTest.persist(commitBuilder); + + assertEquals(1, commitBuilder.getWatermarkHoldsCount()); + + Windmill.WatermarkHold watermarkHold = commitBuilder.getWatermarkHolds(0); + assertEquals(key(NAMESPACE, "watermark"), watermarkHold.getTag()); + assertEquals(TimeUnit.MILLISECONDS.toMicros(1000), watermarkHold.getTimestamps(0)); + + Mockito.verifyNoInteractions(mockReader); + + assertBuildable(commitBuilder); + } + + @Test + public void testNewWatermarkAddPersist() throws Exception { + StateTag addr = + StateTags.watermarkStateInternal("watermark", TimestampCombiner.EARLIEST); + + WatermarkHoldState hold = underTestNewKey.state(NAMESPACE, addr); + + hold.add(new Instant(1000)); + + Windmill.WorkItemCommitRequest.Builder commitBuilder = + Windmill.WorkItemCommitRequest.newBuilder(); + underTestNewKey.persist(commitBuilder); + + assertEquals(1, commitBuilder.getWatermarkHoldsCount()); + + Windmill.WatermarkHold watermarkHold = commitBuilder.getWatermarkHolds(0); + assertEquals(key(NAMESPACE, "watermark"), watermarkHold.getTag()); + assertEquals(TimeUnit.MILLISECONDS.toMicros(1000), watermarkHold.getTimestamps(0)); + + Mockito.verifyNoInteractions(mockReader); + + assertBuildable(commitBuilder); + } + + @Test + public void testNewWatermarkClearPersist() throws Exception { + StateTag addr = + StateTags.watermarkStateInternal("watermark", TimestampCombiner.EARLIEST); + + WatermarkHoldState hold = underTestNewKey.state(NAMESPACE, addr); + + hold.add(new Instant(1000)); + + Windmill.WorkItemCommitRequest.Builder commitBuilder = + Windmill.WorkItemCommitRequest.newBuilder(); + underTestNewKey.persist(commitBuilder); + + assertEquals(1, commitBuilder.getWatermarkHoldsCount()); + + Windmill.WatermarkHold watermarkHold = commitBuilder.getWatermarkHolds(0); + assertEquals(key(NAMESPACE, "watermark"), watermarkHold.getTag()); + assertEquals(TimeUnit.MILLISECONDS.toMicros(1000), watermarkHold.getTimestamps(0)); + + Mockito.verifyNoInteractions(mockReader); + + assertBuildable(commitBuilder); + } + + @Test + public void testWatermarkSetKnownEmptyThenClearPersist() throws Exception { + StateTag addr = + StateTags.watermarkStateInternal("watermark", TimestampCombiner.EARLIEST); + + WatermarkHoldState hold = underTest.state(NAMESPACE, addr); + + hold.setKnownEmpty(); + hold.clear(); Windmill.WorkItemCommitRequest.Builder commitBuilder = Windmill.WorkItemCommitRequest.newBuilder(); underTest.persist(commitBuilder); - // Should be a no-op, no reset, no adds. assertEquals(0, commitBuilder.getWatermarkHoldsCount()); - Mockito.verifyNoMoreInteractions(mockReader); + Mockito.verifyNoInteractions(mockReader); + + assertBuildable(commitBuilder); + } + + @Test + public void testWatermarkSetKnownEmptyThenPersist() throws Exception { + StateTag addr = + StateTags.watermarkStateInternal("watermark", TimestampCombiner.EARLIEST); + + WatermarkHoldState hold = underTest.state(NAMESPACE, addr); + + hold.setKnownEmpty(); + + Windmill.WorkItemCommitRequest.Builder commitBuilder = + Windmill.WorkItemCommitRequest.newBuilder(); + underTest.persist(commitBuilder); + + assertEquals(0, commitBuilder.getWatermarkHoldsCount()); + + Mockito.verifyNoInteractions(mockReader); + + assertBuildable(commitBuilder); + } + + @Test + public void testWatermarkSetKnownEmptyThenClearThenAddPersist() throws Exception { + StateTag addr = + StateTags.watermarkStateInternal("watermark", TimestampCombiner.EARLIEST); + + WatermarkHoldState hold = underTest.state(NAMESPACE, addr); + + hold.setKnownEmpty(); + hold.clear(); + hold.add(new Instant(1000)); + + Windmill.WorkItemCommitRequest.Builder commitBuilder = + Windmill.WorkItemCommitRequest.newBuilder(); + underTest.persist(commitBuilder); + + assertEquals(1, commitBuilder.getWatermarkHoldsCount()); + + Windmill.WatermarkHold watermarkHold = commitBuilder.getWatermarkHolds(0); + assertEquals(key(NAMESPACE, "watermark"), watermarkHold.getTag()); + assertEquals(TimeUnit.MILLISECONDS.toMicros(1000), watermarkHold.getTimestamps(0)); + + Mockito.verifyNoInteractions(mockReader); + + assertBuildable(commitBuilder); } @Test From 8ffe18bd1b08fc2622a09d5438c8f954181873d2 Mon Sep 17 00:00:00 2001 From: Arun Pandian Date: Mon, 27 Apr 2026 22:01:25 +0000 Subject: [PATCH 10/20] doc fix --- .../apache/beam/runners/core/ReduceFnRunner.java | 15 +++++++++------ 1 file changed, 9 insertions(+), 6 deletions(-) diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java index 9a9a096130cd..fad5e71cb828 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java @@ -94,6 +94,11 @@ }) // TODO(https://github.com/apache/beam/issues/20497) public class ReduceFnRunner { + // Experiments guarding optimizations in development. No backward compatibility guarantees. + public static final String UNSTABLE_NOT_UPDATE_COMPATIBLE_NEW_WINDOW_OPTIMIZATION = + "unstable_not_update_compatible_new_window_optimization"; + public static final String UNSTABLE_DISABLE_WATERMARK_KNOWN_EMPTY_OPTIMIZATION = + "unstable_disable_watermark_known_empty_optimization"; /** * The {@link ReduceFnRunner} depends on most aspects of the {@link WindowingStrategy}. * @@ -242,14 +247,12 @@ public ReduceFnRunner( this.nonEmptyPanes = NonEmptyPanes.create(this.windowingStrategy, this.reduceFn); this.useNewWindowOptimization = - options != null - && ExperimentalOptions.hasExperiment( - options, "unstable_not_update_compatible_new_window_optimization"); + ExperimentalOptions.hasExperiment( + options, UNSTABLE_NOT_UPDATE_COMPATIBLE_NEW_WINDOW_OPTIMIZATION); this.disableWatermarkKnownEmptyOptimization = - options != null - && ExperimentalOptions.hasExperiment( - options, "unstable_disable_watermark_known_empty_optimization"); + ExperimentalOptions.hasExperiment( + options, UNSTABLE_DISABLE_WATERMARK_KNOWN_EMPTY_OPTIMIZATION); // Note this may incur I/O to load persisted window set data. this.activeWindows = createActiveWindowSet(); From ae4a5b2004704d53421314e8287447b5c8b28f15 Mon Sep 17 00:00:00 2001 From: Arun Pandian Date: Mon, 27 Apr 2026 22:05:37 +0000 Subject: [PATCH 11/20] improve persist --- .../worker/windmill/state/WindmillWatermarkHold.java | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillWatermarkHold.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillWatermarkHold.java index 2df2699d287a..c0ade373e99f 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillWatermarkHold.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillWatermarkHold.java @@ -148,6 +148,11 @@ public Future persist( Future result; + if (!knownEmpty && !cleared && localAdditions == null) { + // No changes, so no need to update Windmill and no need to cache any value. + return Futures.immediateFuture(Windmill.WorkItemCommitRequest.newBuilder().buildPartial()); + } + if (knownEmpty) { if (localAdditions != null) { // 1. We know it's empty, so we can just update with localAdditions From 64f88995b95adfb32abe453688cee7c48036fa45 Mon Sep 17 00:00:00 2001 From: Arun Pandian Date: Mon, 27 Apr 2026 22:08:02 +0000 Subject: [PATCH 12/20] improve persist --- .../dataflow/worker/windmill/state/WindmillWatermarkHold.java | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillWatermarkHold.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillWatermarkHold.java index c0ade373e99f..be554ce1fb00 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillWatermarkHold.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillWatermarkHold.java @@ -203,8 +203,7 @@ public Future persist( // 5. Otherwise, we need to combine the local additions with the already persisted data result = combineWithPersisted(); } else { - // 6. No changes, so no need to update Windmill. - result = Futures.immediateFuture(Windmill.WorkItemCommitRequest.newBuilder().buildPartial()); + throw new IllegalStateException("Unreachable condition"); } final int estimatedByteSize = ENCODED_SIZE + stateKey.byteString().size(); From de2a32e612e5cbed68f5155142d418aa23dd7418 Mon Sep 17 00:00:00 2001 From: Arun Pandian Date: Mon, 27 Apr 2026 23:05:19 +0000 Subject: [PATCH 13/20] Enable isNewWindowOptimization on non merging windows only --- .../main/java/org/apache/beam/runners/core/ReduceFnRunner.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java index fad5e71cb828..fede633342ce 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java @@ -246,7 +246,7 @@ public ReduceFnRunner( this.nonEmptyPanes = NonEmptyPanes.create(this.windowingStrategy, this.reduceFn); - this.useNewWindowOptimization = + this.useNewWindowOptimization = windowingStrategy.getWindowFn().isNonMerging() && ExperimentalOptions.hasExperiment( options, UNSTABLE_NOT_UPDATE_COMPATIBLE_NEW_WINDOW_OPTIMIZATION); From 3d8f6c583099b38f72dde9c1d55800abec119d45 Mon Sep 17 00:00:00 2001 From: Arun Pandian Date: Mon, 27 Apr 2026 23:06:59 +0000 Subject: [PATCH 14/20] Add todo to prevent global window state growth --- .../main/java/org/apache/beam/runners/core/ReduceFnRunner.java | 2 ++ 1 file changed, 2 insertions(+) diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java index fede633342ce..73a113296bcf 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java @@ -640,6 +640,8 @@ private void processElement(Map windowToMergeResult, WindowedValue StateStyle.RENAMED, value.causedByDrain()); + // TODO: Make sure the NewWindowOptimization does not create unbounded trigger state + // in GlobalWindow if (useNewWindowOptimization && triggerRunner.isNew(directContext.state())) { // Blindly clear state to ensure Windmill doesn't do unnecessary reads. reduceFn.clearState(renamedContext); From 9bb48618b3ed5ce6b2601c7624a6188a8230cc45 Mon Sep 17 00:00:00 2001 From: Arun Pandian Date: Tue, 28 Apr 2026 03:36:47 +0000 Subject: [PATCH 15/20] spotless --- .../java/org/apache/beam/runners/core/ReduceFnRunner.java | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java index 73a113296bcf..b6617ab9afed 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java @@ -246,9 +246,10 @@ public ReduceFnRunner( this.nonEmptyPanes = NonEmptyPanes.create(this.windowingStrategy, this.reduceFn); - this.useNewWindowOptimization = windowingStrategy.getWindowFn().isNonMerging() && - ExperimentalOptions.hasExperiment( - options, UNSTABLE_NOT_UPDATE_COMPATIBLE_NEW_WINDOW_OPTIMIZATION); + this.useNewWindowOptimization = + windowingStrategy.getWindowFn().isNonMerging() + && ExperimentalOptions.hasExperiment( + options, UNSTABLE_NOT_UPDATE_COMPATIBLE_NEW_WINDOW_OPTIMIZATION); this.disableWatermarkKnownEmptyOptimization = ExperimentalOptions.hasExperiment( From 168431d35df50cb82d227300b0eab3fafd8e6ab2 Mon Sep 17 00:00:00 2001 From: Arun Pandian Date: Thu, 14 May 2026 09:46:39 +0000 Subject: [PATCH 16/20] Address comments --- .../beam/runners/core/ReduceFnRunner.java | 7 +- .../triggers/TriggerStateMachineRunner.java | 34 +++-- .../TriggerStateMachineRunnerTest.java | 117 ++++++++++++++++++ .../windmill/state/WindmillWatermarkHold.java | 8 +- 4 files changed, 148 insertions(+), 18 deletions(-) create mode 100644 runners/core-java/src/test/java/org/apache/beam/runners/core/triggers/TriggerStateMachineRunnerTest.java diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java index b6617ab9afed..8573eadca572 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java @@ -274,7 +274,8 @@ public ReduceFnRunner( new TriggerStateMachineRunner<>( triggerStateMachine, new TriggerStateMachineContextFactory<>( - windowingStrategy.getWindowFn(), stateInternals, activeWindows)); + windowingStrategy.getWindowFn(), stateInternals, activeWindows), + this.useNewWindowOptimization); } private ActiveWindowSet createActiveWindowSet() { @@ -787,7 +788,7 @@ public void onTimers(Iterable timers) throws Exception { // Perform prefetching of state to determine if the trigger should fire. if (windowActivation.isGarbageCollection) { - triggerRunner.prefetchIsClosed(directContext.state()); + triggerRunner.prefetchFinishedSet(directContext.state()); } else { triggerRunner.prefetchShouldFire(directContext.window(), directContext.state()); } @@ -966,7 +967,7 @@ private void prefetchEmit( ReduceFn.Context directContext, ReduceFn.Context renamedContext) { triggerRunner.prefetchShouldFire(directContext.window(), directContext.state()); - triggerRunner.prefetchIsClosed(directContext.state()); + triggerRunner.prefetchFinishedSet(directContext.state()); prefetchOnTrigger(directContext, renamedContext); } diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/TriggerStateMachineRunner.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/TriggerStateMachineRunner.java index 2b20087063a0..06ccd76989ca 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/TriggerStateMachineRunner.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/TriggerStateMachineRunner.java @@ -63,13 +63,16 @@ public class TriggerStateMachineRunner { private final ExecutableTriggerStateMachine rootTrigger; private final TriggerStateMachineContextFactory contextFactory; + private final boolean useNewWindowOptimization; public TriggerStateMachineRunner( ExecutableTriggerStateMachine rootTrigger, - TriggerStateMachineContextFactory contextFactory) { + TriggerStateMachineContextFactory contextFactory, + boolean useNewWindowOptimization) { checkState(rootTrigger.getTriggerIndex() == 0); this.rootTrigger = rootTrigger; this.contextFactory = contextFactory; + this.useNewWindowOptimization = useNewWindowOptimization; } private FinishedTriggersBitSet readFinishedBits(ValueState state) { @@ -111,19 +114,19 @@ public BitSet getFinishedBits(StateAccessor state) { return readFinishedBits(state.access(FINISHED_BITS_TAG)).getBitSet(); } - public void prefetchIsClosed(StateAccessor state) { + public void prefetchFinishedSet(StateAccessor state) { if (isFinishedSetNeeded()) { state.access(FINISHED_BITS_TAG).readLater(); } } public void prefetchForValue(W window, StateAccessor state) { - prefetchIsClosed(state); + prefetchFinishedSet(state); rootTrigger.invokePrefetchOnElement(contextFactory.createPrefetchContext(window, rootTrigger)); } public void prefetchShouldFire(W window, StateAccessor state) { - prefetchIsClosed(state); + prefetchFinishedSet(state); rootTrigger.invokePrefetchShouldFire(contextFactory.createPrefetchContext(window, rootTrigger)); } @@ -192,16 +195,29 @@ public void onFire(W window, Timers timers, StateAccessor state) throws Excep persistFinishedSet(state, finishedSet); } - private void persistFinishedSet( - StateAccessor state, FinishedTriggersBitSet modifiedFinishedSet) { + @VisibleForTesting + void persistFinishedSet(StateAccessor state, FinishedTriggersBitSet modifiedFinishedSet) { if (!isFinishedSetNeeded()) { return; } ValueState finishedSetState = state.access(FINISHED_BITS_TAG); - @Nullable BitSet currentBits = finishedSetState.read(); - if (currentBits == null || !currentBits.equals(modifiedFinishedSet.getBitSet())) { - finishedSetState.write(modifiedFinishedSet.getBitSet()); + + if (useNewWindowOptimization) { + @Nullable BitSet bitSet = finishedSetState.read(); + if (bitSet == null || !bitSet.equals(modifiedFinishedSet.getBitSet())) { + // Write a value even if the bitset was empty + finishedSetState.write(modifiedFinishedSet.getBitSet()); + } + return; + } + + if (!readFinishedBits(finishedSetState).equals(modifiedFinishedSet)) { + if (modifiedFinishedSet.getBitSet().isEmpty()) { + finishedSetState.clear(); + } else { + finishedSetState.write(modifiedFinishedSet.getBitSet()); + } } } diff --git a/runners/core-java/src/test/java/org/apache/beam/runners/core/triggers/TriggerStateMachineRunnerTest.java b/runners/core-java/src/test/java/org/apache/beam/runners/core/triggers/TriggerStateMachineRunnerTest.java new file mode 100644 index 000000000000..81c807532756 --- /dev/null +++ b/runners/core-java/src/test/java/org/apache/beam/runners/core/triggers/TriggerStateMachineRunnerTest.java @@ -0,0 +1,117 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.beam.runners.core.triggers; + +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import java.util.BitSet; +import org.apache.beam.runners.core.StateAccessor; +import org.apache.beam.sdk.state.ValueState; +import org.junit.Before; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.junit.runners.JUnit4; +import org.mockito.Mock; +import org.mockito.MockitoAnnotations; + +@RunWith(JUnit4.class) +public class TriggerStateMachineRunnerTest { + + @Mock private StateAccessor mockState; + @Mock private ValueState mockFinishedSetState; + @Mock private TriggerStateMachineContextFactory mockContextFactory; + @Mock private TriggerStateMachine mockTriggerStateMachine; + + private ExecutableTriggerStateMachine rootTrigger; + + @Before + public void setUp() { + MockitoAnnotations.initMocks(this); + when(mockState.access(TriggerStateMachineRunner.FINISHED_BITS_TAG)) + .thenReturn((ValueState) mockFinishedSetState); + rootTrigger = ExecutableTriggerStateMachine.create(mockTriggerStateMachine); + } + + @Test + public void testPersistFinishedSet_emptyAndOptimizationEnabled() throws Exception { + when(mockFinishedSetState.read()).thenReturn(null); + + TriggerStateMachineRunner runner = + new TriggerStateMachineRunner<>( + rootTrigger, + (TriggerStateMachineContextFactory) mockContextFactory, + true /* useNewWindowOptimization */); + + FinishedTriggersBitSet modifiedFinishedSet = FinishedTriggersBitSet.emptyWithCapacity(1); + + runner.persistFinishedSet(mockState, modifiedFinishedSet); + + // Should write empty bitset because optimization is enabled + verify(mockFinishedSetState).write(modifiedFinishedSet.getBitSet()); + } + + @Test + public void testPersistFinishedSet_emptyAndOptimizationDisabled() throws Exception { + when(mockFinishedSetState.read()).thenReturn(null); + + TriggerStateMachineRunner runner = + new TriggerStateMachineRunner<>( + rootTrigger, + (TriggerStateMachineContextFactory) mockContextFactory, + false /* useNewWindowOptimization */); + + FinishedTriggersBitSet modifiedFinishedSet = FinishedTriggersBitSet.emptyWithCapacity(1); + + runner.persistFinishedSet(mockState, modifiedFinishedSet); + + // Should NOT write empty bitset because optimization is disabled and it was already empty (read + // returned null) + verify(mockFinishedSetState, never()).write(modifiedFinishedSet.getBitSet()); + } + + private void runTestPersistFinishedSet_nonEmpty(boolean useNewWindowOptimization) + throws Exception { + when(mockFinishedSetState.read()).thenReturn(null); + + TriggerStateMachineRunner runner = + new TriggerStateMachineRunner<>( + rootTrigger, + (TriggerStateMachineContextFactory) mockContextFactory, + useNewWindowOptimization); + + FinishedTriggersBitSet modifiedFinishedSet = FinishedTriggersBitSet.emptyWithCapacity(1); + modifiedFinishedSet.setFinished(rootTrigger, true); + + runner.persistFinishedSet(mockState, modifiedFinishedSet); + + // Should write non-empty bitset + verify(mockFinishedSetState).write(modifiedFinishedSet.getBitSet()); + } + + @Test + public void testPersistFinishedSet_nonEmpty() throws Exception { + runTestPersistFinishedSet_nonEmpty(false); + } + + @Test + public void testPersistFinishedSet_nonEmptyAndOptimizationEnabled() throws Exception { + runTestPersistFinishedSet_nonEmpty(true); + } +} diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillWatermarkHold.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillWatermarkHold.java index be554ce1fb00..492d3630ba99 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillWatermarkHold.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillWatermarkHold.java @@ -148,11 +148,6 @@ public Future persist( Future result; - if (!knownEmpty && !cleared && localAdditions == null) { - // No changes, so no need to update Windmill and no need to cache any value. - return Futures.immediateFuture(Windmill.WorkItemCommitRequest.newBuilder().buildPartial()); - } - if (knownEmpty) { if (localAdditions != null) { // 1. We know it's empty, so we can just update with localAdditions @@ -203,7 +198,8 @@ public Future persist( // 5. Otherwise, we need to combine the local additions with the already persisted data result = combineWithPersisted(); } else { - throw new IllegalStateException("Unreachable condition"); + // No changes, so no need to update Windmill and no need to cache any value. + return Futures.immediateFuture(Windmill.WorkItemCommitRequest.newBuilder().buildPartial()); } final int estimatedByteSize = ENCODED_SIZE + stateKey.byteString().size(); From ba4766fe9c544b0c371c049f81886b1c7a749094 Mon Sep 17 00:00:00 2001 From: Arun Pandian Date: Thu, 14 May 2026 09:56:59 +0000 Subject: [PATCH 17/20] Refactor --- .../beam/runners/core/triggers/TriggerStateMachineRunner.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/TriggerStateMachineRunner.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/TriggerStateMachineRunner.java index 06ccd76989ca..90daa5ac75b2 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/TriggerStateMachineRunner.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/TriggerStateMachineRunner.java @@ -204,8 +204,8 @@ void persistFinishedSet(StateAccessor state, FinishedTriggersBitSet modifiedF ValueState finishedSetState = state.access(FINISHED_BITS_TAG); if (useNewWindowOptimization) { - @Nullable BitSet bitSet = finishedSetState.read(); - if (bitSet == null || !bitSet.equals(modifiedFinishedSet.getBitSet())) { + if (finishedSetState.read() == null + || !readFinishedBits(finishedSetState).equals(modifiedFinishedSet)) { // Write a value even if the bitset was empty finishedSetState.write(modifiedFinishedSet.getBitSet()); } From ebb5d6f39a8b5df3cd6644fc9f54f7ed1f1c4567 Mon Sep 17 00:00:00 2001 From: Arun Pandian Date: Thu, 14 May 2026 10:00:54 +0000 Subject: [PATCH 18/20] Trigger validate runner --- ...beam_PostCommit_Java_ValidatesRunner_Dataflow_Streaming.json | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.github/trigger_files/beam_PostCommit_Java_ValidatesRunner_Dataflow_Streaming.json b/.github/trigger_files/beam_PostCommit_Java_ValidatesRunner_Dataflow_Streaming.json index 090751435f20..c2110eeca47c 100644 --- a/.github/trigger_files/beam_PostCommit_Java_ValidatesRunner_Dataflow_Streaming.json +++ b/.github/trigger_files/beam_PostCommit_Java_ValidatesRunner_Dataflow_Streaming.json @@ -1,4 +1,4 @@ { "comment": "Modify this file in a trivial way to cause this test suite to run!", - "modification": 7, + "modification": 8, } From 3650cfc3fbbd0eb36efe15cd66bf2ada2b639f86 Mon Sep 17 00:00:00 2001 From: Arun Pandian Date: Mon, 18 May 2026 12:36:38 +0000 Subject: [PATCH 19/20] Fix conflict & add todo --- ...beam_PostCommit_Java_ValidatesRunner_Dataflow_Streaming.json | 2 +- .../main/java/org/apache/beam/runners/core/ReduceFnRunner.java | 2 ++ 2 files changed, 3 insertions(+), 1 deletion(-) diff --git a/.github/trigger_files/beam_PostCommit_Java_ValidatesRunner_Dataflow_Streaming.json b/.github/trigger_files/beam_PostCommit_Java_ValidatesRunner_Dataflow_Streaming.json index c2110eeca47c..e623d3373a93 100644 --- a/.github/trigger_files/beam_PostCommit_Java_ValidatesRunner_Dataflow_Streaming.json +++ b/.github/trigger_files/beam_PostCommit_Java_ValidatesRunner_Dataflow_Streaming.json @@ -1,4 +1,4 @@ { "comment": "Modify this file in a trivial way to cause this test suite to run!", - "modification": 8, + "modification": 1, } diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java index 8573eadca572..4d3bfdbe4b19 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java @@ -646,6 +646,8 @@ private void processElement(Map windowToMergeResult, WindowedValue // in GlobalWindow if (useNewWindowOptimization && triggerRunner.isNew(directContext.state())) { // Blindly clear state to ensure Windmill doesn't do unnecessary reads. + // TODO: Instead of the clears here, we could mark these states as empty locally + // in the state cache and/or explicitly tell that the entries are non-existent reduceFn.clearState(renamedContext); paneInfoTracker.clear(directContext.state()); if (!disableWatermarkKnownEmptyOptimization) { From f3edf1c0a6fc662f3447398923ece5b08d6762d6 Mon Sep 17 00:00:00 2001 From: Arun Pandian Date: Tue, 19 May 2026 00:56:25 +0000 Subject: [PATCH 20/20] Doc fix --- .../main/java/org/apache/beam/runners/core/ReduceFnRunner.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java index 4d3bfdbe4b19..445ce82e56b7 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java @@ -647,7 +647,7 @@ private void processElement(Map windowToMergeResult, WindowedValue if (useNewWindowOptimization && triggerRunner.isNew(directContext.state())) { // Blindly clear state to ensure Windmill doesn't do unnecessary reads. // TODO: Instead of the clears here, we could mark these states as empty locally - // in the state cache and/or explicitly tell that the entries are non-existent + // in the state cache and/or explicitly tell that the entries are non-existent via api reduceFn.clearState(renamedContext); paneInfoTracker.clear(directContext.state()); if (!disableWatermarkKnownEmptyOptimization) {