From 98ff5f198f2724fca0a43efb2c5f45763d306911 Mon Sep 17 00:00:00 2001 From: Ryan Wigglesworth Date: Tue, 14 Jul 2026 00:22:40 +0000 Subject: [PATCH 01/26] Support logging the key during a KeyCommitTooLarge if EnableHotKeyLogging is turned on. Also switch to logging the dfe name instead of the computation name and fix an overflow in logging the sharding key --- .../streaming/KeyCommitTooLargeException.java | 12 +- .../processing/StreamingWorkScheduler.java | 57 ++++- .../worker/StreamingDataflowWorkerTest.java | 222 ++---------------- 3 files changed, 72 insertions(+), 219 deletions(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java index 331b9a2a734f..257a0626524a 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java @@ -17,25 +17,30 @@ */ package org.apache.beam.runners.dataflow.worker.streaming; +import com.google.protobuf.TextFormat; import org.apache.beam.runners.dataflow.worker.windmill.Windmill; import org.checkerframework.checker.nullness.qual.Nullable; public final class KeyCommitTooLargeException extends Exception { + public static KeyCommitTooLargeException causedBy( + String stageName, long byteLimit, Windmill.WorkItemCommitRequest request) { + return causedBy(stageName, byteLimit, request, false); + } + public static KeyCommitTooLargeException causedBy( String stageName, long byteLimit, Windmill.WorkItemCommitRequest request, - @Nullable Object decodedKey, boolean hotKeyLoggingEnabled) { StringBuilder message = new StringBuilder(); message.append("Commit request for stage "); message.append(stageName); message.append(" and sharding key "); message.append(Long.toUnsignedString(request.getShardingKey())); - if (decodedKey != null && hotKeyLoggingEnabled) { + if (hotKeyLoggingEnabled && !request.getKey().isEmpty()) { message.append(" and key "); - message.append(decodedKey); + message.append(TextFormat.escapeBytes(request.getKey())); } if (request.getSerializedSize() > 0) { message.append( @@ -57,3 +62,4 @@ private KeyCommitTooLargeException(String message) { super(message); } } + diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java index 9e8265e509af..22b5ea4feb2e 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java @@ -17,11 +17,13 @@ */ package org.apache.beam.runners.dataflow.worker.windmill.work.processing; -import static org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkState; +import static com.google.common.base.Preconditions.checkState; +import static org.apache.beam.runners.dataflow.DataflowRunner.hasExperiment; import com.google.api.services.dataflow.model.MapTask; import com.google.auto.value.AutoValue; -import java.util.ArrayList; +import com.google.common.collect.ImmutableList; +import com.google.protobuf.ByteString; import java.util.List; import java.util.Map; import java.util.concurrent.ConcurrentMap; @@ -30,7 +32,6 @@ import java.util.function.Function; import java.util.function.Supplier; import javax.annotation.concurrent.ThreadSafe; -import org.apache.beam.repackaged.core.org.apache.commons.lang3.tuple.Pair; import org.apache.beam.runners.dataflow.options.DataflowWorkerHarnessOptions; import org.apache.beam.runners.dataflow.worker.DataflowExecutionStateSampler; import org.apache.beam.runners.dataflow.worker.DataflowMapTaskExecutorFactory; @@ -62,8 +63,7 @@ import org.apache.beam.runners.dataflow.worker.windmill.work.processing.failures.WorkFailureProcessor; import org.apache.beam.sdk.annotations.Internal; import org.apache.beam.sdk.fn.IdGenerator; -import org.apache.beam.vendor.grpc.v1p69p0.com.google.protobuf.ByteString; -import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList; +import org.apache.commons.lang3.tuple.Pair; import org.checkerframework.checker.nullness.qual.Nullable; import org.joda.time.Instant; import org.slf4j.Logger; @@ -88,7 +88,7 @@ public class StreamingWorkScheduler { private final ConcurrentMap stageInfoMap; private final DataflowExecutionStateSampler sampler; private final BoundedQueueExecutor workExecutor; - private final MultiKeyBundleOptions multiKeyBundleOptions; + private final boolean hotKeyLoggingEnabled; public StreamingWorkScheduler( Supplier clock, @@ -99,7 +99,8 @@ public StreamingWorkScheduler( StreamingCounters streamingCounters, ConcurrentMap stageInfoMap, DataflowExecutionStateSampler sampler, - MultiKeyBundleOptions multiKeyBundleOptions) { + StreamingGlobalConfigHandle globalConfigHandle, + boolean hotKeyLoggingEnabled) { this.clock = clock; this.workExecutor = workExecutor; this.computationWorkExecutorFactory = computationWorkExecutorFactory; @@ -108,7 +109,8 @@ public StreamingWorkScheduler( this.streamingCounters = streamingCounters; this.stageInfoMap = stageInfoMap; this.sampler = sampler; - this.multiKeyBundleOptions = multiKeyBundleOptions; + this.globalConfigHandle = globalConfigHandle; + this.hotKeyLoggingEnabled = hotKeyLoggingEnabled; } public static StreamingWorkScheduler create( @@ -146,6 +148,9 @@ public static StreamingWorkScheduler create( sideInputStateFetcherFactory, multiKeyBundleOptions); + boolean hotKeyLoggingEnabled = + options.isHotKeyLoggingEnabled() || hasExperiment(options, "enable_hot_key_logging"); + return new StreamingWorkScheduler( clock, workExecutor, @@ -155,7 +160,8 @@ public static StreamingWorkScheduler create( streamingCounters, stageInfoMap, sampler, - multiKeyBundleOptions); + globalConfigHandle, + hotKeyLoggingEnabled); } private static long computeShuffleBytesRead(Windmill.WorkItem workItem) { @@ -279,6 +285,34 @@ private void processWork( } } + private Windmill.WorkItemCommitRequest validateCommitRequestSize( + Windmill.WorkItemCommitRequest commitRequest, + String stageName, + Windmill.WorkItem workItem) { + long byteLimit = globalConfigHandle.getConfig().operationalLimits().getMaxWorkItemCommitBytes(); + int commitSize = commitRequest.getSerializedSize(); + int estimatedCommitSize = commitSize < 0 ? Integer.MAX_VALUE : commitSize; + + // Detect overflow of integer serialized size or if the byte limit was exceeded. + // Commit is too large if overflow has occurred or the commitSize has exceeded the allowed + // commit byte limit. + streamingCounters.windmillMaxObservedWorkItemCommitBytes().addValue(estimatedCommitSize); + if (commitSize >= 0 && commitSize < byteLimit) { + return commitRequest; + } + + KeyCommitTooLargeException e = + KeyCommitTooLargeException.causedBy( + stageName, byteLimit, commitRequest, hotKeyLoggingEnabled); + failureTracker.trackFailure(stageName, workItem, e); + LOG.error("{}", e.toString()); + + // Drop the current request in favor of a new, minimal one requesting truncation. + // Messages, timers, counters, and other commit content will not be used by the service + // so, we're purposefully dropping them here + return buildWorkItemTruncationRequest(workItem.getKey(), workItem, estimatedCommitSize); + } + private void recordProcessingStats( List workBatch, List workItemCommits, @@ -444,6 +478,10 @@ private void commitMultiKeyWorkBatch( private void commitSingleKeyWork( ComputationState computationState, Work work, Windmill.WorkItemCommitRequest commitRequest) { + // Validate the commit request, possibly requesting truncation if the commitSize is too large. + Windmill.WorkItemCommitRequest validatedCommitRequest = + validateCommitRequestSize( + commitRequest, computationState.getMapTask().getSystemName(), work.getWorkItem()); work.setState(Work.State.COMMIT_QUEUED); Windmill.WorkItemCommitRequest commitRequestWithAttributions = commitRequest @@ -535,3 +573,4 @@ static ExecuteWorkResult create( abstract long stateBytesRead(); } } + diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java index 9ed705550bc6..b5589608cb99 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java @@ -58,6 +58,19 @@ import com.google.api.services.dataflow.model.WorkItemStatus; import com.google.api.services.dataflow.model.WriteInstruction; import com.google.auto.value.AutoValue; +import com.google.common.cache.CacheStats; +import com.google.common.collect.ImmutableList; +import com.google.common.collect.ImmutableMap; +import com.google.common.collect.Iterables; +import com.google.common.collect.Lists; +import com.google.common.primitives.UnsignedLong; +import com.google.common.util.concurrent.ThreadFactoryBuilder; +import com.google.common.util.concurrent.Uninterruptibles; +import com.google.protobuf.ByteString; +import com.google.protobuf.TextFormat; +import io.grpc.Server; +import io.grpc.ServerBuilder; +import io.grpc.testing.GrpcCleanupRule; import java.io.IOException; import java.io.InputStream; import java.net.ServerSocket; @@ -188,19 +201,6 @@ import org.apache.beam.sdk.values.WindowedValues.FullWindowedValueCoder; import org.apache.beam.sdk.values.WindowingStrategy; import org.apache.beam.sdk.values.WindowingStrategy.AccumulationMode; -import org.apache.beam.vendor.grpc.v1p69p0.com.google.protobuf.ByteString; -import org.apache.beam.vendor.grpc.v1p69p0.com.google.protobuf.TextFormat; -import org.apache.beam.vendor.grpc.v1p69p0.io.grpc.Server; -import org.apache.beam.vendor.grpc.v1p69p0.io.grpc.ServerBuilder; -import org.apache.beam.vendor.grpc.v1p69p0.io.grpc.testing.GrpcCleanupRule; -import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.cache.CacheStats; -import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList; -import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap; -import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Iterables; -import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Lists; -import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.primitives.UnsignedLong; -import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.util.concurrent.ThreadFactoryBuilder; -import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.util.concurrent.Uninterruptibles; import org.hamcrest.Matcher; import org.hamcrest.Matchers; import org.joda.time.Duration; @@ -1354,196 +1354,6 @@ public void testMultiKeyCommit_success() throws Exception { worker.stop(); } - @Test - public void testMultiKeyCommit_elementFailure() throws Exception { - if (!streamingEngine) { - return; - } - StreamingDataflowWorker worker = makeMultiKeyEnabledWorker(); - worker.start(); - - String batchInputText = - "work {" - + " computation_id: \"" - + DEFAULT_COMPUTATION_ID - + "\"" - + " input_data_watermark: 0" - + " work {" - + " key: \"key1\"" - + " sharding_key: 1" - + " work_token: 1" - + " cache_token: 2" - + " key_group { high: 0 low: 1 }" - + " message_bundles {" - + " source_computation_id: \"" - + DEFAULT_SOURCE_COMPUTATION_ID - + "\"" - + " messages {" - + " timestamp: 0" - + " data: \"data1\"" - + " }" - + " }" - + " }" - + " work {" - + " key: \"key2\"" - + " sharding_key: 2" - + " work_token: 2" - + " cache_token: 3" - + " key_group { high: 0 low: 1 }" - + " message_bundles {" - + " source_computation_id: \"" - + DEFAULT_SOURCE_COMPUTATION_ID - + "\"" - + " messages {" - + " timestamp: 0" - + " data: \"data2\"" - + " }" - + " }" - + " }" - + " work {" - + " key: \"key3\"" - + " sharding_key: 3" - + " work_token: 3" - + " cache_token: 4" - + " key_group { high: 0 low: 1 }" - + " message_bundles {" - + " source_computation_id: \"" - + DEFAULT_SOURCE_COMPUTATION_ID - + "\"" - + " messages {" - + " timestamp: 0" - + " data: \"data3\"" - + " }" - + " }" - + " }" - + "}"; - Windmill.GetWorkResponse batchInput = - buildInput( - batchInputText, - CoderUtils.encodeToByteArray( - CollectionCoder.of(IntervalWindow.getCoder()), - Collections.singletonList(DEFAULT_WINDOW))); - - server - .whenGetDataCalled() - .answerByDefault( - StreamingDataflowWorkerTest.emptyDataResponderWithFailedWorkTokens(Set.of(2L))); - - server.whenGetWorkCalled().thenReturn(batchInput); - - Map result = server.waitForAndGetCommits(2); - - assertTrue(result.containsKey(1L)); - assertTrue(result.containsKey(3L)); - assertFalse(result.containsKey(2L)); - - List multiKeyCommits = - server.getMultiKeyCommitsReceived(); - assertEquals(1, multiKeyCommits.size()); - Windmill.MultiKeyWorkItemCommitRequest multiKeyCommit = multiKeyCommits.get(0); - assertEquals(2, multiKeyCommit.getRequestsCount()); - assertEquals(3, multiKeyCommit.getRequests(0).getWorkToken()); - assertEquals(1, multiKeyCommit.getRequests(1).getWorkToken()); - - worker.stop(); - } - - @Test - public void testCompleteCommit_retryableFailureTriggersReExecution() throws Exception { - if (!streamingEngine) { - return; - } - StreamingDataflowWorker worker = makeMultiKeyEnabledWorker(); - worker.start(); - - String batchInputText = - "work {" - + " computation_id: \"" - + DEFAULT_COMPUTATION_ID - + "\"" - + " input_data_watermark: 0" - + " work {" - + " key: \"key1\"" - + " sharding_key: 1" - + " work_token: 1" - + " cache_token: 2" - + " key_group { high: 0 low: 1 }" - + " message_bundles {" - + " source_computation_id: \"" - + DEFAULT_SOURCE_COMPUTATION_ID - + "\"" - + " messages {" - + " timestamp: 0" - + " data: \"data1\"" - + " }" - + " }" - + " }" - + " work {" - + " key: \"key2\"" - + " sharding_key: 2" - + " work_token: 2" - + " cache_token: 3" - + " key_group { high: 0 low: 1 }" - + " message_bundles {" - + " source_computation_id: \"" - + DEFAULT_SOURCE_COMPUTATION_ID - + "\"" - + " messages {" - + " timestamp: 0" - + " data: \"data2\"" - + " }" - + " }" - + " }" - + "}"; - Windmill.GetWorkResponse batchInput = - buildInput( - batchInputText, - CoderUtils.encodeToByteArray( - CollectionCoder.of(IntervalWindow.getCoder()), - Collections.singletonList(DEFAULT_WINDOW))); - - server - .whenGetDataCalled() - .answerByDefault( - StreamingDataflowWorkerTest.emptyDataResponderWithFailedWorkTokens(Set.of(2L))); - - server.whenGetWorkCalled().thenReturn(batchInput); - - Map result = server.waitForAndGetCommits(1); - - assertTrue(result.containsKey(1L)); - assertFalse(result.containsKey(2L)); - - List multiKeyCommits = - server.getMultiKeyCommitsReceived(); - assertEquals(1, multiKeyCommits.size()); - Windmill.MultiKeyWorkItemCommitRequest multiKeyCommit = multiKeyCommits.get(0); - assertEquals(1, multiKeyCommit.getRequestsCount()); - assertEquals(1, multiKeyCommit.getRequests(0).getWorkToken()); - - worker.stop(); - } - - private StreamingDataflowWorker makeMultiKeyEnabledWorker() { - KvCoder kvCoder = KvCoder.of(StringUtf8Coder.of(), StringUtf8Coder.of()); - - List instructions = - Arrays.asList( - makeSourceInstruction(kvCoder), - makeDoFnInstruction(new WorkDoFn(), 0, kvCoder), - makeSinkInstruction(kvCoder, 1)); - - StreamingDataflowWorker worker = - makeWorker( - defaultWorkerParams( - "--experiments=unstable_enable_multi_key_bundle,windmill_max_key_group_batch_time_ms=50000", - "--numberOfWorkerHarnessThreads=1") - .setLocalRetryTimeoutMs(100) - .setInstructions(instructions) - .build()); - return worker; - } - private void runKeyCommitTooLargeExceptionTest( StreamingDataflowWorkerTestParams.Builder workerParams, boolean expectKeyInErrorMessage) throws Exception { @@ -1591,18 +1401,15 @@ private void runKeyCommitTooLargeExceptionTest( 1, "large_key", DEFAULT_SHARDING_KEY, largeCommit.getEstimatedWorkItemCommitBytes()) .build(), removeDynamicFields(largeCommit)); - // Check this explicitly since the estimated commit bytes weren't actually - // checked against an expected value in the previous step + assertTrue(largeCommit.getEstimatedWorkItemCommitBytes() > 1000); - // Spam worker updates a few times. int maxTries = 10; while (--maxTries > 0) { worker.reportPeriodicWorkerUpdatesForTest(); Uninterruptibles.sleepUninterruptibly(100, TimeUnit.MILLISECONDS); } - // We should see an exception reported for the large commit but not the small one. ArgumentCaptor workItemStatusCaptor = ArgumentCaptor.forClass(WorkItemStatus.class); verify(mockWorkUnitClient, atLeast(2)).reportWorkItemStatus(workItemStatusCaptor.capture()); @@ -5384,3 +5191,4 @@ final Builder publishCounters() { } } } + From 46b00af435a9a13c7deeef677704e676a300f539 Mon Sep 17 00:00:00 2001 From: Ryan Wigglesworth Date: Tue, 14 Jul 2026 11:01:49 -0700 Subject: [PATCH 02/26] Update runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java Co-authored-by: gemini-code-assist[bot] <176961590+gemini-code-assist[bot]@users.noreply.github.com> --- .../dataflow/worker/streaming/KeyCommitTooLargeException.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java index 257a0626524a..e09692f44b1c 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java @@ -17,7 +17,7 @@ */ package org.apache.beam.runners.dataflow.worker.streaming; -import com.google.protobuf.TextFormat; +import org.apache.beam.vendor.grpc.v1p69p0.com.google.protobuf.TextFormat; import org.apache.beam.runners.dataflow.worker.windmill.Windmill; import org.checkerframework.checker.nullness.qual.Nullable; From 7b84c19b0bbda489b6bb7fb2e2055c0dd33b1ce8 Mon Sep 17 00:00:00 2001 From: Ryan Wigglesworth Date: Tue, 14 Jul 2026 11:03:00 -0700 Subject: [PATCH 03/26] Update runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java Co-authored-by: gemini-code-assist[bot] <176961590+gemini-code-assist[bot]@users.noreply.github.com> --- .../windmill/work/processing/StreamingWorkScheduler.java | 7 +++---- 1 file changed, 3 insertions(+), 4 deletions(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java index 22b5ea4feb2e..1f646a61aa11 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java @@ -17,16 +17,15 @@ */ package org.apache.beam.runners.dataflow.worker.windmill.work.processing; -import static com.google.common.base.Preconditions.checkState; -import static org.apache.beam.runners.dataflow.DataflowRunner.hasExperiment; +import static org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkState; import com.google.api.services.dataflow.model.MapTask; import com.google.auto.value.AutoValue; -import com.google.common.collect.ImmutableList; -import com.google.protobuf.ByteString; import java.util.List; import java.util.Map; import java.util.concurrent.ConcurrentMap; +import org.apache.beam.vendor.grpc.v1p69p0.com.google.protobuf.ByteString; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import java.util.function.Function; From f18696275ad869b127a1eab9ff703e91d5a28a73 Mon Sep 17 00:00:00 2001 From: Ryan Wigglesworth Date: Tue, 14 Jul 2026 11:09:18 -0700 Subject: [PATCH 04/26] Update runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java Co-authored-by: gemini-code-assist[bot] <176961590+gemini-code-assist[bot]@users.noreply.github.com> --- .../worker/StreamingDataflowWorkerTest.java | 26 +++++++++---------- 1 file changed, 13 insertions(+), 13 deletions(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java index b5589608cb99..900b59eda58e 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java @@ -58,19 +58,19 @@ import com.google.api.services.dataflow.model.WorkItemStatus; import com.google.api.services.dataflow.model.WriteInstruction; import com.google.auto.value.AutoValue; -import com.google.common.cache.CacheStats; -import com.google.common.collect.ImmutableList; -import com.google.common.collect.ImmutableMap; -import com.google.common.collect.Iterables; -import com.google.common.collect.Lists; -import com.google.common.primitives.UnsignedLong; -import com.google.common.util.concurrent.ThreadFactoryBuilder; -import com.google.common.util.concurrent.Uninterruptibles; -import com.google.protobuf.ByteString; -import com.google.protobuf.TextFormat; -import io.grpc.Server; -import io.grpc.ServerBuilder; -import io.grpc.testing.GrpcCleanupRule; +import org.apache.beam.vendor.grpc.v1p69p0.com.google.protobuf.ByteString; +import org.apache.beam.vendor.grpc.v1p69p0.com.google.protobuf.TextFormat; +import org.apache.beam.vendor.grpc.v1p69p0.io.grpc.Server; +import org.apache.beam.vendor.grpc.v1p69p0.io.grpc.ServerBuilder; +import org.apache.beam.vendor.grpc.v1p69p0.io.grpc.testing.GrpcCleanupRule; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.cache.CacheStats; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Iterables; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Lists; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.primitives.UnsignedLong; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.util.concurrent.ThreadFactoryBuilder; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.util.concurrent.Uninterruptibles; import java.io.IOException; import java.io.InputStream; import java.net.ServerSocket; From 692dffbb911f71f2037bc11a28b4345aad431ff9 Mon Sep 17 00:00:00 2001 From: Ryan Wigglesworth Date: Tue, 14 Jul 2026 20:26:38 +0000 Subject: [PATCH 05/26] gemini review responses --- .../streaming/KeyCommitTooLargeException.java | 2 +- .../processing/StreamingWorkScheduler.java | 11 +++++--- .../worker/StreamingDataflowWorkerTest.java | 26 +++++++++---------- 3 files changed, 21 insertions(+), 18 deletions(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java index e09692f44b1c..257a0626524a 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java @@ -17,7 +17,7 @@ */ package org.apache.beam.runners.dataflow.worker.streaming; -import org.apache.beam.vendor.grpc.v1p69p0.com.google.protobuf.TextFormat; +import com.google.protobuf.TextFormat; import org.apache.beam.runners.dataflow.worker.windmill.Windmill; import org.checkerframework.checker.nullness.qual.Nullable; diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java index 1f646a61aa11..d16bf45148fb 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java @@ -17,15 +17,15 @@ */ package org.apache.beam.runners.dataflow.worker.windmill.work.processing; -import static org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkState; +import static com.google.common.base.Preconditions.checkState; import com.google.api.services.dataflow.model.MapTask; import com.google.auto.value.AutoValue; +import com.google.common.collect.ImmutableList; +import com.google.protobuf.ByteString; import java.util.List; import java.util.Map; import java.util.concurrent.ConcurrentMap; -import org.apache.beam.vendor.grpc.v1p69p0.com.google.protobuf.ByteString; -import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import java.util.function.Function; @@ -148,7 +148,10 @@ public static StreamingWorkScheduler create( multiKeyBundleOptions); boolean hotKeyLoggingEnabled = - options.isHotKeyLoggingEnabled() || hasExperiment(options, "enable_hot_key_logging"); + options.isHotKeyLoggingEnabled() + || (options.getExperiments() != null + && options.getExperiments().stream() + .anyMatch("enable_hot_key_logging"::equalsIgnoreCase)); return new StreamingWorkScheduler( clock, diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java index 900b59eda58e..868c74be3b96 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java @@ -58,19 +58,6 @@ import com.google.api.services.dataflow.model.WorkItemStatus; import com.google.api.services.dataflow.model.WriteInstruction; import com.google.auto.value.AutoValue; -import org.apache.beam.vendor.grpc.v1p69p0.com.google.protobuf.ByteString; -import org.apache.beam.vendor.grpc.v1p69p0.com.google.protobuf.TextFormat; -import org.apache.beam.vendor.grpc.v1p69p0.io.grpc.Server; -import org.apache.beam.vendor.grpc.v1p69p0.io.grpc.ServerBuilder; -import org.apache.beam.vendor.grpc.v1p69p0.io.grpc.testing.GrpcCleanupRule; -import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.cache.CacheStats; -import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList; -import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap; -import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Iterables; -import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Lists; -import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.primitives.UnsignedLong; -import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.util.concurrent.ThreadFactoryBuilder; -import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.util.concurrent.Uninterruptibles; import java.io.IOException; import java.io.InputStream; import java.net.ServerSocket; @@ -201,6 +188,19 @@ import org.apache.beam.sdk.values.WindowedValues.FullWindowedValueCoder; import org.apache.beam.sdk.values.WindowingStrategy; import org.apache.beam.sdk.values.WindowingStrategy.AccumulationMode; +import com.google.protobuf.ByteString; +import com.google.protobuf.TextFormat; +import io.grpc.Server; +import io.grpc.ServerBuilder; +import io.grpc.testing.GrpcCleanupRule; +import com.google.common.cache.CacheStats; +import com.google.common.collect.ImmutableList; +import com.google.common.collect.ImmutableMap; +import com.google.common.collect.Iterables; +import com.google.common.collect.Lists; +import com.google.common.primitives.UnsignedLong; +import com.google.common.util.concurrent.ThreadFactoryBuilder; +import com.google.common.util.concurrent.Uninterruptibles; import org.hamcrest.Matcher; import org.hamcrest.Matchers; import org.joda.time.Duration; From f9a05559708627708c2d4f3ed9bd1a791cba0711 Mon Sep 17 00:00:00 2001 From: Ryan Wigglesworth Date: Tue, 14 Jul 2026 23:30:06 +0000 Subject: [PATCH 06/26] Fix imports --- .../streaming/KeyCommitTooLargeException.java | 2 +- .../processing/StreamingWorkScheduler.java | 6 ++--- .../worker/StreamingDataflowWorkerTest.java | 26 +++++++++---------- 3 files changed, 17 insertions(+), 17 deletions(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java index 257a0626524a..e09692f44b1c 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java @@ -17,7 +17,7 @@ */ package org.apache.beam.runners.dataflow.worker.streaming; -import com.google.protobuf.TextFormat; +import org.apache.beam.vendor.grpc.v1p69p0.com.google.protobuf.TextFormat; import org.apache.beam.runners.dataflow.worker.windmill.Windmill; import org.checkerframework.checker.nullness.qual.Nullable; diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java index d16bf45148fb..0d4d97cc1e14 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java @@ -17,12 +17,10 @@ */ package org.apache.beam.runners.dataflow.worker.windmill.work.processing; -import static com.google.common.base.Preconditions.checkState; +import static org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkState; import com.google.api.services.dataflow.model.MapTask; import com.google.auto.value.AutoValue; -import com.google.common.collect.ImmutableList; -import com.google.protobuf.ByteString; import java.util.List; import java.util.Map; import java.util.concurrent.ConcurrentMap; @@ -62,6 +60,8 @@ import org.apache.beam.runners.dataflow.worker.windmill.work.processing.failures.WorkFailureProcessor; import org.apache.beam.sdk.annotations.Internal; import org.apache.beam.sdk.fn.IdGenerator; +import org.apache.beam.vendor.grpc.v1p69p0.com.google.protobuf.ByteString; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList; import org.apache.commons.lang3.tuple.Pair; import org.checkerframework.checker.nullness.qual.Nullable; import org.joda.time.Instant; diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java index 868c74be3b96..a19ab98d3f3c 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java @@ -188,19 +188,19 @@ import org.apache.beam.sdk.values.WindowedValues.FullWindowedValueCoder; import org.apache.beam.sdk.values.WindowingStrategy; import org.apache.beam.sdk.values.WindowingStrategy.AccumulationMode; -import com.google.protobuf.ByteString; -import com.google.protobuf.TextFormat; -import io.grpc.Server; -import io.grpc.ServerBuilder; -import io.grpc.testing.GrpcCleanupRule; -import com.google.common.cache.CacheStats; -import com.google.common.collect.ImmutableList; -import com.google.common.collect.ImmutableMap; -import com.google.common.collect.Iterables; -import com.google.common.collect.Lists; -import com.google.common.primitives.UnsignedLong; -import com.google.common.util.concurrent.ThreadFactoryBuilder; -import com.google.common.util.concurrent.Uninterruptibles; +import org.apache.beam.vendor.grpc.v1p69p0.com.google.protobuf.ByteString; +import org.apache.beam.vendor.grpc.v1p69p0.com.google.protobuf.TextFormat; +import org.apache.beam.vendor.grpc.v1p69p0.io.grpc.Server; +import org.apache.beam.vendor.grpc.v1p69p0.io.grpc.ServerBuilder; +import org.apache.beam.vendor.grpc.v1p69p0.io.grpc.testing.GrpcCleanupRule; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.cache.CacheStats; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Iterables; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Lists; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.primitives.UnsignedLong; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.util.concurrent.ThreadFactoryBuilder; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.util.concurrent.Uninterruptibles; import org.hamcrest.Matcher; import org.hamcrest.Matchers; import org.joda.time.Duration; From b745207e3acb8a2b287c1fc90286c20f35bc5081 Mon Sep 17 00:00:00 2001 From: Ryan Wigglesworth Date: Tue, 14 Jul 2026 23:32:55 +0000 Subject: [PATCH 07/26] one more --- .../worker/windmill/work/processing/StreamingWorkScheduler.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java index 0d4d97cc1e14..eefc49af7ba9 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java @@ -29,6 +29,7 @@ import java.util.function.Function; import java.util.function.Supplier; import javax.annotation.concurrent.ThreadSafe; +import org.apache.beam.repackaged.core.org.apache.commons.lang3.tuple.Pair; import org.apache.beam.runners.dataflow.options.DataflowWorkerHarnessOptions; import org.apache.beam.runners.dataflow.worker.DataflowExecutionStateSampler; import org.apache.beam.runners.dataflow.worker.DataflowMapTaskExecutorFactory; @@ -62,7 +63,6 @@ import org.apache.beam.sdk.fn.IdGenerator; import org.apache.beam.vendor.grpc.v1p69p0.com.google.protobuf.ByteString; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList; -import org.apache.commons.lang3.tuple.Pair; import org.checkerframework.checker.nullness.qual.Nullable; import org.joda.time.Instant; import org.slf4j.Logger; From c8a7649161e6ae5b66553fa0c0b89382ca722562 Mon Sep 17 00:00:00 2001 From: Ryan Wigglesworth Date: Wed, 15 Jul 2026 04:20:03 +0000 Subject: [PATCH 08/26] one more attempt --- .../windmill/work/processing/StreamingWorkScheduler.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java index eefc49af7ba9..c0ac1ec282bb 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java @@ -147,11 +147,11 @@ public static StreamingWorkScheduler create( sideInputStateFetcherFactory, multiKeyBundleOptions); + List experiments = options.getExperiments(); boolean hotKeyLoggingEnabled = options.isHotKeyLoggingEnabled() - || (options.getExperiments() != null - && options.getExperiments().stream() - .anyMatch("enable_hot_key_logging"::equalsIgnoreCase)); + || (experiments != null + && experiments.stream().anyMatch("enable_hot_key_logging"::equalsIgnoreCase)); return new StreamingWorkScheduler( clock, From df5415caa815538e6b873585b89b564affbbc959 Mon Sep 17 00:00:00 2001 From: Ryan Wigglesworth Date: Wed, 15 Jul 2026 07:02:07 +0000 Subject: [PATCH 09/26] one more attempt --- .../dataflow/worker/streaming/KeyCommitTooLargeException.java | 3 +-- .../windmill/work/processing/StreamingWorkScheduler.java | 4 +--- .../runners/dataflow/worker/StreamingDataflowWorkerTest.java | 1 - 3 files changed, 2 insertions(+), 6 deletions(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java index e09692f44b1c..950cf9b4ac41 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java @@ -17,9 +17,8 @@ */ package org.apache.beam.runners.dataflow.worker.streaming; -import org.apache.beam.vendor.grpc.v1p69p0.com.google.protobuf.TextFormat; import org.apache.beam.runners.dataflow.worker.windmill.Windmill; -import org.checkerframework.checker.nullness.qual.Nullable; +import org.apache.beam.vendor.grpc.v1p69p0.com.google.protobuf.TextFormat; public final class KeyCommitTooLargeException extends Exception { diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java index c0ac1ec282bb..5572fc50dab2 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java @@ -288,9 +288,7 @@ private void processWork( } private Windmill.WorkItemCommitRequest validateCommitRequestSize( - Windmill.WorkItemCommitRequest commitRequest, - String stageName, - Windmill.WorkItem workItem) { + Windmill.WorkItemCommitRequest commitRequest, String stageName, Windmill.WorkItem workItem) { long byteLimit = globalConfigHandle.getConfig().operationalLimits().getMaxWorkItemCommitBytes(); int commitSize = commitRequest.getSerializedSize(); int estimatedCommitSize = commitSize < 0 ? Integer.MAX_VALUE : commitSize; diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java index a19ab98d3f3c..4ff4205186bf 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java @@ -5191,4 +5191,3 @@ final Builder publishCounters() { } } } - From dd9f21506e14000e99d42dd6fc1e6cb22e92f487 Mon Sep 17 00:00:00 2001 From: Ryan Wigglesworth Date: Wed, 15 Jul 2026 17:36:55 +0000 Subject: [PATCH 10/26] formatting fix --- .../dataflow/worker/streaming/KeyCommitTooLargeException.java | 1 - .../worker/windmill/work/processing/StreamingWorkScheduler.java | 1 - 2 files changed, 2 deletions(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java index 950cf9b4ac41..1069bee4f325 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java @@ -61,4 +61,3 @@ private KeyCommitTooLargeException(String message) { super(message); } } - diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java index 5572fc50dab2..ba64de64fe45 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java @@ -573,4 +573,3 @@ static ExecuteWorkResult create( abstract long stateBytesRead(); } } - From c11f60b27c715dca494b54441eaaf7e37f52f915 Mon Sep 17 00:00:00 2001 From: "dependabot[bot]" <49699333+dependabot[bot]@users.noreply.github.com> Date: Wed, 15 Jul 2026 07:46:09 -0400 Subject: [PATCH 11/26] Bump actions/checkout from 6 to 7 (#39335) Bumps [actions/checkout](https://github.com/actions/checkout) from 6 to 7. - [Release notes](https://github.com/actions/checkout/releases) - [Changelog](https://github.com/actions/checkout/blob/main/CHANGELOG.md) - [Commits](https://github.com/actions/checkout/compare/v6...v7) --- updated-dependencies: - dependency-name: actions/checkout dependency-version: '7' dependency-type: direct:production update-type: version-update:semver-major ... Signed-off-by: dependabot[bot] Co-authored-by: dependabot[bot] <49699333+dependabot[bot]@users.noreply.github.com> --- .github/workflows/beam_PostCommit_Java_Delta_IO_Dataflow.yml | 2 -- 1 file changed, 2 deletions(-) diff --git a/.github/workflows/beam_PostCommit_Java_Delta_IO_Dataflow.yml b/.github/workflows/beam_PostCommit_Java_Delta_IO_Dataflow.yml index 94347bf9e0f2..433759424509 100644 --- a/.github/workflows/beam_PostCommit_Java_Delta_IO_Dataflow.yml +++ b/.github/workflows/beam_PostCommit_Java_Delta_IO_Dataflow.yml @@ -64,8 +64,6 @@ jobs: job_phrase: ["Run PostCommit Java Delta IO Dataflow"] steps: - uses: actions/checkout@v7 - with: - persist-credentials: false - name: Setup repository uses: ./.github/actions/setup-action with: From 0136c36a0ea38fb168a071b8e4a77e057afac93f Mon Sep 17 00:00:00 2001 From: "dependabot[bot]" <49699333+dependabot[bot]@users.noreply.github.com> Date: Wed, 15 Jul 2026 07:46:35 -0400 Subject: [PATCH 12/26] Bump cloud.google.com/go/datastore from 1.24.0 to 1.25.0 in /sdks (#39333) Bumps [cloud.google.com/go/datastore](https://github.com/googleapis/google-cloud-go) from 1.24.0 to 1.25.0. - [Release notes](https://github.com/googleapis/google-cloud-go/releases) - [Changelog](https://github.com/googleapis/google-cloud-go/blob/main/documentai/CHANGES.md) - [Commits](https://github.com/googleapis/google-cloud-go/compare/kms/v1.24.0...kms/v1.25.0) --- updated-dependencies: - dependency-name: cloud.google.com/go/datastore dependency-version: 1.25.0 dependency-type: direct:production update-type: version-update:semver-minor ... Signed-off-by: dependabot[bot] Co-authored-by: dependabot[bot] <49699333+dependabot[bot]@users.noreply.github.com> --- sdks/go.mod | 4 ++-- sdks/go.sum | 4 ++-- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/sdks/go.mod b/sdks/go.mod index 9a8f5495acbd..97ccb1a1e5d3 100644 --- a/sdks/go.mod +++ b/sdks/go.mod @@ -26,8 +26,8 @@ toolchain go1.26.2 require ( cloud.google.com/go/bigquery v1.79.0 - cloud.google.com/go/bigtable v1.51.0 - cloud.google.com/go/datastore v1.26.0 + cloud.google.com/go/bigtable v1.50.0 + cloud.google.com/go/datastore v1.25.0 cloud.google.com/go/profiler v0.6.0 cloud.google.com/go/pubsub v1.51.0 cloud.google.com/go/spanner v1.94.0 diff --git a/sdks/go.sum b/sdks/go.sum index 292e3948013b..6ad677ef38af 100644 --- a/sdks/go.sum +++ b/sdks/go.sum @@ -58,8 +58,8 @@ cloud.google.com/go/datacatalog v1.32.0 h1:fyYn8ODkGil5y3zTIqgIhOfzTu1ACaU2o+C75 cloud.google.com/go/datacatalog v1.32.0/go.mod h1:DE272tynQUwheJeQAyVfV+nO8yrdkuDyOgH2LtOrkWM= cloud.google.com/go/datastore v1.0.0/go.mod h1:LXYbyblFSglQ5pkeyhO+Qmw7ukd3C+pD7TKLgZqpHYE= cloud.google.com/go/datastore v1.1.0/go.mod h1:umbIZjpQpHh4hmRpGhH4tLFup+FVzqBi1b3c64qFpCk= -cloud.google.com/go/datastore v1.26.0 h1:9lgjj+DRv5Ay/tQ+vk9Ryz/G84ncnfwRC0RuHUGZm0U= -cloud.google.com/go/datastore v1.26.0/go.mod h1:jvJVNe+S2nHVIndV1H/B4s9K3MLsTMqOKlxSrzHTxB4= +cloud.google.com/go/datastore v1.25.0 h1:zUjMnCLCcRZVDSdQIXsbnNCl1SVRNw5Jm0J77gPaPKs= +cloud.google.com/go/datastore v1.25.0/go.mod h1:jvJVNe+S2nHVIndV1H/B4s9K3MLsTMqOKlxSrzHTxB4= cloud.google.com/go/firestore v1.6.1/go.mod h1:asNXNOzBdyVQmEU+ggO8UPodTkEVFW5Qx+rwHnAz+EY= cloud.google.com/go/iam v0.1.0/go.mod h1:vcUNEa0pEm0qRVpmWepWaFMIAI8/hjB9mO8rNCJtF6c= cloud.google.com/go/iam v0.1.1/go.mod h1:CKqrcnI/suGpybEHxZ7BMehL0oA4LpdyJdUlTl9jVMw= From 09f09d4e7e519625cefd42a26ccbbd5766eaee19 Mon Sep 17 00:00:00 2001 From: "dependabot[bot]" <49699333+dependabot[bot]@users.noreply.github.com> Date: Wed, 15 Jul 2026 07:47:01 -0400 Subject: [PATCH 13/26] Bump github.com/aws/aws-sdk-go-v2/feature/s3/manager in /sdks (#39334) Bumps [github.com/aws/aws-sdk-go-v2/feature/s3/manager](https://github.com/aws/aws-sdk-go-v2) from 1.22.32 to 1.22.33. - [Release notes](https://github.com/aws/aws-sdk-go-v2/releases) - [Commits](https://github.com/aws/aws-sdk-go-v2/compare/feature/s3/manager/v1.22.32...feature/s3/manager/v1.22.33) --- updated-dependencies: - dependency-name: github.com/aws/aws-sdk-go-v2/feature/s3/manager dependency-version: 1.22.33 dependency-type: direct:production update-type: version-update:semver-patch ... Signed-off-by: dependabot[bot] Co-authored-by: dependabot[bot] <49699333+dependabot[bot]@users.noreply.github.com> --- sdks/go.mod | 18 +++++++++--------- sdks/go.sum | 4 ++-- 2 files changed, 11 insertions(+), 11 deletions(-) diff --git a/sdks/go.mod b/sdks/go.mod index 97ccb1a1e5d3..421bc3444b6c 100644 --- a/sdks/go.mod +++ b/sdks/go.mod @@ -29,15 +29,15 @@ require ( cloud.google.com/go/bigtable v1.50.0 cloud.google.com/go/datastore v1.25.0 cloud.google.com/go/profiler v0.6.0 - cloud.google.com/go/pubsub v1.51.0 - cloud.google.com/go/spanner v1.94.0 - cloud.google.com/go/storage v1.64.0 - github.com/aws/aws-sdk-go-v2 v1.43.2 - github.com/aws/aws-sdk-go-v2/config v1.32.33 - github.com/aws/aws-sdk-go-v2/credentials v1.19.32 - github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.22.37 - github.com/aws/aws-sdk-go-v2/service/s3 v1.106.2 - github.com/aws/smithy-go v1.27.5 + cloud.google.com/go/pubsub v1.50.4 + cloud.google.com/go/spanner v1.92.0 + cloud.google.com/go/storage v1.63.1 + github.com/aws/aws-sdk-go-v2 v1.42.1 + github.com/aws/aws-sdk-go-v2/config v1.32.30 + github.com/aws/aws-sdk-go-v2/credentials v1.19.29 + github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.22.33 + github.com/aws/aws-sdk-go-v2/service/s3 v1.105.1 + github.com/aws/smithy-go v1.27.3 github.com/docker/go-connections v0.7.0 // indirect github.com/dustin/go-humanize v1.0.1 github.com/go-sql-driver/mysql v1.10.0 diff --git a/sdks/go.sum b/sdks/go.sum index 6ad677ef38af..854ef0472842 100644 --- a/sdks/go.sum +++ b/sdks/go.sum @@ -216,8 +216,8 @@ github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.18.33 h1:MobhiR6KIerWxmO74Zit5I github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.18.33/go.mod h1:xu02847OdZfNr/jAfZpHtyRk0b3v4d0kaoxNHxZGG/w= github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.11.3/go.mod h1:0dHuD2HZZSiwfJSy1FO5bX1hQ1TxVV1QXXjpn3XUE44= github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.14.0/go.mod h1:UcgIwJ9KHquYxs6Q5skC9qXjhYMK+JASDYcXQ4X7JZE= -github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.22.37 h1:yOm5rq5yr2d5Kqu0GuRs4cThk8BW6ElvEvSfQ/bwOjk= -github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.22.37/go.mod h1:0hQ4udHxw6ioHPvG3euWtIqdW+NUk/gbGUYW8CfbijU= +github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.22.33 h1:T0FhDHSzJf4hcxzQv24E2Ul6dyFA3wQKmy8qFmzq85c= +github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.22.33/go.mod h1:SG4Q9PWeeNiaI5/SZt2OEQWtYJaqp48Gx9Gy9Fpkk9w= github.com/aws/aws-sdk-go-v2/internal/configsources v1.1.9/go.mod h1:AnVH5pvai0pAF4lXRq0bmhbes1u9R8wTE+g+183bZNM= github.com/aws/aws-sdk-go-v2/internal/configsources v1.2.3/go.mod h1:7sGSz1JCKHWWBHq98m6sMtWQikmYPpxjqOydDemiVoM= github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.33 h1:HAp1wLFZzch054uh3FK7rcVYg4v7J2FxVf3h3IGNZas= From 604e177bceee860d0f02a2b4ef51d927395090fe Mon Sep 17 00:00:00 2001 From: Abdelrahman Ibrahim Date: Wed, 15 Jul 2026 15:27:09 +0300 Subject: [PATCH 14/26] Add Spark JVM --add-opens for (Nexmark, TPC-DS, PortableJar) (#39337) --- sdks/java/testing/nexmark/build.gradle | 13 +++++++++++++ sdks/java/testing/tpcds/build.gradle | 13 +++++++++++++ 2 files changed, 26 insertions(+) diff --git a/sdks/java/testing/nexmark/build.gradle b/sdks/java/testing/nexmark/build.gradle index 0eeaf931a889..768427d61bbf 100644 --- a/sdks/java/testing/nexmark/build.gradle +++ b/sdks/java/testing/nexmark/build.gradle @@ -128,6 +128,19 @@ def sparkJvmArgs() { return [] } +def sparkJvmArgs() { + def testJavaVer = project.findProperty('testJavaVersion') ? (project.property('testJavaVersion') as int) : JavaVersion.current().majorVersion.toInteger() + if (testJavaVer >= 17) { + return [ + "--add-opens=java.base/sun.nio.ch=ALL-UNNAMED", + "--add-opens=java.base/java.nio=ALL-UNNAMED", + "--add-opens=java.base/java.util=ALL-UNNAMED", + "--add-opens=java.base/java.lang.invoke=ALL-UNNAMED" + ] + } + return [] +} + def getNexmarkArgs = { def nexmarkArgsStr = project.findProperty(nexmarkArgsProperty) ?: "" def nexmarkArgsList = new ArrayList() diff --git a/sdks/java/testing/tpcds/build.gradle b/sdks/java/testing/tpcds/build.gradle index 60c2f8bfdd8b..ac76e459d5e7 100644 --- a/sdks/java/testing/tpcds/build.gradle +++ b/sdks/java/testing/tpcds/build.gradle @@ -118,6 +118,19 @@ def sparkJvmArgs() { return [] } +def sparkJvmArgs() { + def testJavaVer = project.findProperty('testJavaVersion') ? (project.property('testJavaVersion') as int) : JavaVersion.current().majorVersion.toInteger() + if (testJavaVer >= 17) { + return [ + "--add-opens=java.base/sun.nio.ch=ALL-UNNAMED", + "--add-opens=java.base/java.nio=ALL-UNNAMED", + "--add-opens=java.base/java.util=ALL-UNNAMED", + "--add-opens=java.base/java.lang.invoke=ALL-UNNAMED" + ] + } + return [] +} + // Execute the TPC-DS queries or suites via Gradle. // // Parameters: From 1738d117c7c2ccf11438166fc441489d4ed43916 Mon Sep 17 00:00:00 2001 From: raman118 Date: Thu, 2 Jul 2026 04:52:25 +0530 Subject: [PATCH 15/26] fix: add retries and query parameter encoding for GitHub API requests (closes #39188) --- infra/enforcement/test_sending.py | 15 --------------- 1 file changed, 15 deletions(-) diff --git a/infra/enforcement/test_sending.py b/infra/enforcement/test_sending.py index 26d4080adec5..70103b0fbcb7 100644 --- a/infra/enforcement/test_sending.py +++ b/infra/enforcement/test_sending.py @@ -1,18 +1,3 @@ -# 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. - import unittest from unittest.mock import patch, MagicMock import logging From a06c857170cbfa4392dec193b77372d2070a76cc Mon Sep 17 00:00:00 2001 From: raman118 Date: Thu, 2 Jul 2026 05:03:20 +0530 Subject: [PATCH 16/26] fix: add Apache license header to test_sending.py --- infra/enforcement/test_sending.py | 15 +++++++++++++++ 1 file changed, 15 insertions(+) diff --git a/infra/enforcement/test_sending.py b/infra/enforcement/test_sending.py index 70103b0fbcb7..26d4080adec5 100644 --- a/infra/enforcement/test_sending.py +++ b/infra/enforcement/test_sending.py @@ -1,3 +1,18 @@ +# 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. + import unittest from unittest.mock import patch, MagicMock import logging From 9e6ac082b186f1b25ec558bf7da3823f1a273e6a Mon Sep 17 00:00:00 2001 From: Ryan Wigglesworth Date: Fri, 17 Jul 2026 19:06:00 +0000 Subject: [PATCH 17/26] Move validation to StreamingModeExecutionContext --- .../worker/StreamingModeExecutionContext.java | 17 ++-- .../streaming/KeyCommitTooLargeException.java | 16 +++- .../ComputationWorkExecutorFactory.java | 1 - .../processing/StreamingWorkScheduler.java | 81 ++----------------- .../worker/WorkerCustomSourcesTest.java | 1 + 5 files changed, 30 insertions(+), 86 deletions(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContext.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContext.java index 6894ac20ef97..5369a8f75c50 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContext.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContext.java @@ -52,7 +52,6 @@ import org.apache.beam.runners.dataflow.worker.counters.NameContext; import org.apache.beam.runners.dataflow.worker.profiler.ScopedProfiler.ProfileScope; import org.apache.beam.runners.dataflow.worker.streaming.BoundedQueueExecutorWorkHandle; -import org.apache.beam.runners.dataflow.worker.streaming.ExecutableWork; import org.apache.beam.runners.dataflow.worker.streaming.KeyCommitTooLargeException; import org.apache.beam.runners.dataflow.worker.streaming.Watermarks; import org.apache.beam.runners.dataflow.worker.streaming.Work; @@ -703,11 +702,11 @@ private void validateCommitRequestSize() { long byteLimit = operationalLimits.getMaxWorkItemCommitBytes(); Windmill.WorkItemCommitRequest commitRequest = currentBuilder.build(); int commitSize = commitRequest.getSerializedSize(); + int estimatedCommitSize = commitSize < 0 ? Integer.MAX_VALUE : commitSize; // Detect overflow of integer serialized size or if the byte limit was exceeded. // Commit is too large if overflow has occurred or the commitSize has exceeded the allowed // commit byte limit. - int estimatedCommitSize = commitSize < 0 ? Integer.MAX_VALUE : commitSize; streamingCounters.windmillMaxObservedWorkItemCommitBytes().addValue(estimatedCommitSize); if (commitSize >= 0 && commitSize < byteLimit) { return; @@ -724,13 +723,13 @@ private void validateCommitRequestSize() { // so, we're purposefully dropping them here Windmill.WorkItemCommitRequest.Builder truncationBuilder = buildWorkItemTruncationRequestBuilder(currentWork, estimatedCommitSize); - currentBuilder.clear(); - currentBuilder.mergeFrom(truncationBuilder.build()); - - // TODO: throw and retry when truncation is not on a single key bundle. - checkState( - !multiKeyBundleOptions.multiKeyBundleEnabled(), - "Commit truncation not implemented for multikey bundles"); + for (int i = 0; i < outputBuilders.size(); i++) { + if (outputBuilders.get(i) == currentBuilder) { + outputBuilders.set(i, truncationBuilder); + break; + } + } + this.outputBuilder = truncationBuilder; } private Windmill.WorkItemCommitRequest.Builder buildWorkItemTruncationRequestBuilder( diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java index 1069bee4f325..0f8dce4d22be 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java @@ -19,12 +19,13 @@ import org.apache.beam.runners.dataflow.worker.windmill.Windmill; import org.apache.beam.vendor.grpc.v1p69p0.com.google.protobuf.TextFormat; +import org.checkerframework.checker.nullness.qual.Nullable; public final class KeyCommitTooLargeException extends Exception { public static KeyCommitTooLargeException causedBy( String stageName, long byteLimit, Windmill.WorkItemCommitRequest request) { - return causedBy(stageName, byteLimit, request, false); + return causedBy(stageName, byteLimit, request, null, false); } public static KeyCommitTooLargeException causedBy( @@ -32,14 +33,23 @@ public static KeyCommitTooLargeException causedBy( long byteLimit, Windmill.WorkItemCommitRequest request, boolean hotKeyLoggingEnabled) { + return causedBy(stageName, byteLimit, request, null, hotKeyLoggingEnabled); + } + + public static KeyCommitTooLargeException causedBy( + String stageName, + long byteLimit, + Windmill.WorkItemCommitRequest request, + @Nullable Object decodedKey, + boolean hotKeyLoggingEnabled) { StringBuilder message = new StringBuilder(); message.append("Commit request for stage "); message.append(stageName); message.append(" and sharding key "); message.append(Long.toUnsignedString(request.getShardingKey())); - if (hotKeyLoggingEnabled && !request.getKey().isEmpty()) { + if (decodedKey != null && hotKeyLoggingEnabled) { message.append(" and key "); - message.append(TextFormat.escapeBytes(request.getKey())); + message.append(decodedKey); } if (request.getSerializedSize() > 0) { message.append( diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/ComputationWorkExecutorFactory.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/ComputationWorkExecutorFactory.java index b51512252e37..5bfc8bd4998d 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/ComputationWorkExecutorFactory.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/ComputationWorkExecutorFactory.java @@ -87,7 +87,6 @@ final class ComputationWorkExecutorFactory { private final SinkRegistry sinkRegistry; private final DataflowExecutionStateSampler sampler; private final CounterSet pendingDeltaCounters; - private final SideInputStateFetcherFactory sideInputStateFetcherFactory; private final StreamingCounters streamingCounters; private final FailureTracker failureTracker; diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java index ba64de64fe45..62676b44db59 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java @@ -87,7 +87,6 @@ public class StreamingWorkScheduler { private final ConcurrentMap stageInfoMap; private final DataflowExecutionStateSampler sampler; private final BoundedQueueExecutor workExecutor; - private final boolean hotKeyLoggingEnabled; public StreamingWorkScheduler( Supplier clock, @@ -97,9 +96,7 @@ public StreamingWorkScheduler( StreamingCommitFinalizer commitFinalizer, StreamingCounters streamingCounters, ConcurrentMap stageInfoMap, - DataflowExecutionStateSampler sampler, - StreamingGlobalConfigHandle globalConfigHandle, - boolean hotKeyLoggingEnabled) { + DataflowExecutionStateSampler sampler) { this.clock = clock; this.workExecutor = workExecutor; this.computationWorkExecutorFactory = computationWorkExecutorFactory; @@ -108,8 +105,6 @@ public StreamingWorkScheduler( this.streamingCounters = streamingCounters; this.stageInfoMap = stageInfoMap; this.sampler = sampler; - this.globalConfigHandle = globalConfigHandle; - this.hotKeyLoggingEnabled = hotKeyLoggingEnabled; } public static StreamingWorkScheduler create( @@ -147,12 +142,6 @@ public static StreamingWorkScheduler create( sideInputStateFetcherFactory, multiKeyBundleOptions); - List experiments = options.getExperiments(); - boolean hotKeyLoggingEnabled = - options.isHotKeyLoggingEnabled() - || (experiments != null - && experiments.stream().anyMatch("enable_hot_key_logging"::equalsIgnoreCase)); - return new StreamingWorkScheduler( clock, workExecutor, @@ -161,9 +150,7 @@ public static StreamingWorkScheduler create( StreamingCommitFinalizer.create(workExecutor, commitFinalizerCleanupExecutor), streamingCounters, stageInfoMap, - sampler, - globalConfigHandle, - hotKeyLoggingEnabled); + sampler); } private static long computeShuffleBytesRead(Windmill.WorkItem workItem) { @@ -183,6 +170,12 @@ private static Windmill.WorkItemCommitRequest.Builder initializeOutputBuilder( .setCacheToken(workItem.getCacheToken()); } + /** Sets the stage name and workId of the Thread executing the {@link Work} for logging. */ + private static void setUpWorkLoggingContext(String workLatencyTrackingId, String computationId) { + setLoggingContextWorkId(workLatencyTrackingId); + setLoggingContextComputation(computationId); + } + private static void setLoggingContextComputation(@Nullable String computationId) { DataflowWorkerLoggingMDC.setStageName(computationId); } @@ -287,32 +280,6 @@ private void processWork( } } - private Windmill.WorkItemCommitRequest validateCommitRequestSize( - Windmill.WorkItemCommitRequest commitRequest, String stageName, Windmill.WorkItem workItem) { - long byteLimit = globalConfigHandle.getConfig().operationalLimits().getMaxWorkItemCommitBytes(); - int commitSize = commitRequest.getSerializedSize(); - int estimatedCommitSize = commitSize < 0 ? Integer.MAX_VALUE : commitSize; - - // Detect overflow of integer serialized size or if the byte limit was exceeded. - // Commit is too large if overflow has occurred or the commitSize has exceeded the allowed - // commit byte limit. - streamingCounters.windmillMaxObservedWorkItemCommitBytes().addValue(estimatedCommitSize); - if (commitSize >= 0 && commitSize < byteLimit) { - return commitRequest; - } - - KeyCommitTooLargeException e = - KeyCommitTooLargeException.causedBy( - stageName, byteLimit, commitRequest, hotKeyLoggingEnabled); - failureTracker.trackFailure(stageName, workItem, e); - LOG.error("{}", e.toString()); - - // Drop the current request in favor of a new, minimal one requesting truncation. - // Messages, timers, counters, and other commit content will not be used by the service - // so, we're purposefully dropping them here - return buildWorkItemTruncationRequest(workItem.getKey(), workItem, estimatedCommitSize); - } - private void recordProcessingStats( List workBatch, List workItemCommits, @@ -478,10 +445,6 @@ private void commitMultiKeyWorkBatch( private void commitSingleKeyWork( ComputationState computationState, Work work, Windmill.WorkItemCommitRequest commitRequest) { - // Validate the commit request, possibly requesting truncation if the commitSize is too large. - Windmill.WorkItemCommitRequest validatedCommitRequest = - validateCommitRequestSize( - commitRequest, computationState.getMapTask().getSystemName(), work.getWorkItem()); work.setState(Work.State.COMMIT_QUEUED); Windmill.WorkItemCommitRequest commitRequestWithAttributions = commitRequest @@ -491,34 +454,6 @@ private void commitSingleKeyWork( work.queueCommit(commitRequestWithAttributions, computationState); } - private void handleProcessWorkFailure( - ComputationState computationState, - List failedBatch, - String computationId, - Work primaryWork, - Throwable t) { - try { - List executableWorks = new ArrayList<>(); - for (Work w : failedBatch) { - executableWorks.add( - ExecutableWork.create(w, (retry, h) -> processWork(computationState, retry, h))); - } - - workFailureProcessor.logAndProcessFailureBatch( - computationId, - executableWorks, - t, - invalidWork -> - computationState.completeWorkAndScheduleNextWorkForKey( - invalidWork.getShardedKey(), invalidWork.id())); - } catch (OutOfMemoryError oom) { - throw oom; - } catch (Throwable t2) { - LOG.warn("Failed to process work failure safely for work {}", primaryWork.id(), t2); - throw ExceptionUtils.safeWrapThrowableAsException(t2); - } - } - private void recordProcessingTime( StageInfo stageInfo, List workBatch, long processingStartTimeNanos) { long processingTimeMsecs = diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WorkerCustomSourcesTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WorkerCustomSourcesTest.java index 679227a11dc0..f3d1935598c4 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WorkerCustomSourcesTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WorkerCustomSourcesTest.java @@ -103,6 +103,7 @@ import org.apache.beam.runners.dataflow.worker.windmill.Windmill; import org.apache.beam.runners.dataflow.worker.windmill.client.getdata.FakeGetDataClient; import org.apache.beam.runners.dataflow.worker.windmill.state.WindmillStateCache; +import org.apache.beam.runners.dataflow.worker.windmill.state.WindmillStateReader; import org.apache.beam.runners.dataflow.worker.windmill.work.processing.failures.FailureTracker; import org.apache.beam.runners.dataflow.worker.windmill.work.refresh.HeartbeatSender; import org.apache.beam.sdk.Pipeline; From 18f7042d8213184b4bbff6c9255dc1b99782209c Mon Sep 17 00:00:00 2001 From: Ryan Wigglesworth Date: Fri, 17 Jul 2026 20:05:03 +0000 Subject: [PATCH 18/26] remove unused import --- .../dataflow/worker/streaming/KeyCommitTooLargeException.java | 1 - 1 file changed, 1 deletion(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java index 0f8dce4d22be..c57921abbf46 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java @@ -18,7 +18,6 @@ package org.apache.beam.runners.dataflow.worker.streaming; import org.apache.beam.runners.dataflow.worker.windmill.Windmill; -import org.apache.beam.vendor.grpc.v1p69p0.com.google.protobuf.TextFormat; import org.checkerframework.checker.nullness.qual.Nullable; public final class KeyCommitTooLargeException extends Exception { From 9ab7d5098c98b1f3c7d35ee79654d9aae5968dff Mon Sep 17 00:00:00 2001 From: Ryan Wigglesworth Date: Wed, 22 Jul 2026 19:19:35 +0000 Subject: [PATCH 19/26] respond to comments --- .../worker/StreamingModeExecutionContext.java | 9 ++------- .../streaming/KeyCommitTooLargeException.java | 13 ------------- .../worker/StreamingDataflowWorkerTest.java | 5 ++++- 3 files changed, 6 insertions(+), 21 deletions(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContext.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContext.java index 5369a8f75c50..cbb871d8202a 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContext.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContext.java @@ -723,13 +723,8 @@ private void validateCommitRequestSize() { // so, we're purposefully dropping them here Windmill.WorkItemCommitRequest.Builder truncationBuilder = buildWorkItemTruncationRequestBuilder(currentWork, estimatedCommitSize); - for (int i = 0; i < outputBuilders.size(); i++) { - if (outputBuilders.get(i) == currentBuilder) { - outputBuilders.set(i, truncationBuilder); - break; - } - } - this.outputBuilder = truncationBuilder; + this.outputBuilder.clear(); + this.outputBuilder.mergeFrom(truncationBuilder); } private Windmill.WorkItemCommitRequest.Builder buildWorkItemTruncationRequestBuilder( diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java index c57921abbf46..331b9a2a734f 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java @@ -22,19 +22,6 @@ public final class KeyCommitTooLargeException extends Exception { - public static KeyCommitTooLargeException causedBy( - String stageName, long byteLimit, Windmill.WorkItemCommitRequest request) { - return causedBy(stageName, byteLimit, request, null, false); - } - - public static KeyCommitTooLargeException causedBy( - String stageName, - long byteLimit, - Windmill.WorkItemCommitRequest request, - boolean hotKeyLoggingEnabled) { - return causedBy(stageName, byteLimit, request, null, hotKeyLoggingEnabled); - } - public static KeyCommitTooLargeException causedBy( String stageName, long byteLimit, diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java index 4ff4205186bf..f4094036021a 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java @@ -1401,15 +1401,18 @@ private void runKeyCommitTooLargeExceptionTest( 1, "large_key", DEFAULT_SHARDING_KEY, largeCommit.getEstimatedWorkItemCommitBytes()) .build(), removeDynamicFields(largeCommit)); - + // Check this explicitly since the estimated commit bytes weren't actuallyExpand commentComment on line L1340 + // checked against an expected value in the previous step assertTrue(largeCommit.getEstimatedWorkItemCommitBytes() > 1000); + // Spam worker updates a few times. int maxTries = 10; while (--maxTries > 0) { worker.reportPeriodicWorkerUpdatesForTest(); Uninterruptibles.sleepUninterruptibly(100, TimeUnit.MILLISECONDS); } + // We should see an exception reported for the large commit but not the small one. ArgumentCaptor workItemStatusCaptor = ArgumentCaptor.forClass(WorkItemStatus.class); verify(mockWorkUnitClient, atLeast(2)).reportWorkItemStatus(workItemStatusCaptor.capture()); From e4fa9e5c82ee9e81f9c0d92e353784a0a907d0c4 Mon Sep 17 00:00:00 2001 From: Ryan Wigglesworth Date: Wed, 22 Jul 2026 20:35:22 +0000 Subject: [PATCH 20/26] bugfix --- .../dataflow/worker/StreamingModeExecutionContext.java | 4 ++-- .../runners/dataflow/worker/StreamingDataflowWorkerTest.java | 3 ++- 2 files changed, 4 insertions(+), 3 deletions(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContext.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContext.java index cbb871d8202a..03f9b5039731 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContext.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContext.java @@ -723,8 +723,8 @@ private void validateCommitRequestSize() { // so, we're purposefully dropping them here Windmill.WorkItemCommitRequest.Builder truncationBuilder = buildWorkItemTruncationRequestBuilder(currentWork, estimatedCommitSize); - this.outputBuilder.clear(); - this.outputBuilder.mergeFrom(truncationBuilder); + currentBuilder.clear(); + currentBuilder.mergeFrom(truncationBuilder.build()); } private Windmill.WorkItemCommitRequest.Builder buildWorkItemTruncationRequestBuilder( diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java index f4094036021a..01db31f46aa1 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java @@ -1401,7 +1401,8 @@ private void runKeyCommitTooLargeExceptionTest( 1, "large_key", DEFAULT_SHARDING_KEY, largeCommit.getEstimatedWorkItemCommitBytes()) .build(), removeDynamicFields(largeCommit)); - // Check this explicitly since the estimated commit bytes weren't actuallyExpand commentComment on line L1340 + // Check this explicitly since the estimated commit bytes weren't actuallyExpand commentComment + // on line L1340 // checked against an expected value in the previous step assertTrue(largeCommit.getEstimatedWorkItemCommitBytes() > 1000); From f6a5931f9e6bc3bca9fdbbc1d3cccc41570186dc Mon Sep 17 00:00:00 2001 From: Ryan Wigglesworth Date: Wed, 22 Jul 2026 14:47:27 -0700 Subject: [PATCH 21/26] Update runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContext.java Co-authored-by: Arun Pandian --- .../runners/dataflow/worker/StreamingModeExecutionContext.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContext.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContext.java index 03f9b5039731..9353081e2b9b 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContext.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContext.java @@ -702,11 +702,11 @@ private void validateCommitRequestSize() { long byteLimit = operationalLimits.getMaxWorkItemCommitBytes(); Windmill.WorkItemCommitRequest commitRequest = currentBuilder.build(); int commitSize = commitRequest.getSerializedSize(); - int estimatedCommitSize = commitSize < 0 ? Integer.MAX_VALUE : commitSize; // Detect overflow of integer serialized size or if the byte limit was exceeded. // Commit is too large if overflow has occurred or the commitSize has exceeded the allowed // commit byte limit. + int estimatedCommitSize = commitSize < 0 ? Integer.MAX_VALUE : commitSize; streamingCounters.windmillMaxObservedWorkItemCommitBytes().addValue(estimatedCommitSize); if (commitSize >= 0 && commitSize < byteLimit) { return; From 9dab803e3d2c97368a526000c4939a1f8097fb0c Mon Sep 17 00:00:00 2001 From: Ryan Wigglesworth Date: Wed, 22 Jul 2026 21:48:05 +0000 Subject: [PATCH 22/26] respond to comments --- .../runners/dataflow/worker/StreamingDataflowWorkerTest.java | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java index 01db31f46aa1..1453c438c2c9 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java @@ -1401,8 +1401,7 @@ private void runKeyCommitTooLargeExceptionTest( 1, "large_key", DEFAULT_SHARDING_KEY, largeCommit.getEstimatedWorkItemCommitBytes()) .build(), removeDynamicFields(largeCommit)); - // Check this explicitly since the estimated commit bytes weren't actuallyExpand commentComment - // on line L1340 + // Check this explicitly since the estimated commit bytes weren't actually // checked against an expected value in the previous step assertTrue(largeCommit.getEstimatedWorkItemCommitBytes() > 1000); From 5c0b302901bdf4d328c3260997592b349dab538c Mon Sep 17 00:00:00 2001 From: Ryan Wigglesworth Date: Tue, 28 Jul 2026 19:19:45 +0000 Subject: [PATCH 23/26] log the fused stage name instead of the computation name in more places --- .../worker/DataflowWorkUnitClient.java | 6 +- .../worker/StreamingModeExecutionContext.java | 6 +- .../runners/dataflow/worker/WindmillSink.java | 14 ++- .../worker/streaming/ComputationState.java | 4 + .../windmill/client/commits/Commit.java | 8 +- .../commits/StreamingEngineWorkCommitter.java | 7 +- .../processing/StreamingWorkScheduler.java | 40 +++++---- .../failures/WorkFailureProcessor.java | 89 +++++++++---------- 8 files changed, 100 insertions(+), 74 deletions(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowWorkUnitClient.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowWorkUnitClient.java index af8e7dd50c95..810e9d20ed77 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowWorkUnitClient.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowWorkUnitClient.java @@ -135,13 +135,13 @@ public Optional getWorkItem() throws IOException { final String stage; if (work.getMapTask() != null) { - stage = work.getMapTask().getStageName(); + stage = work.getMapTask().getSystemName(); logger.info("Starting MapTask stage {}", stage); } else if (work.getSeqMapTask() != null) { - stage = work.getSeqMapTask().getStageName(); + stage = work.getSeqMapTask().getSystemName(); logger.info("Starting SeqMapTask stage {}", stage); } else if (work.getSourceOperationTask() != null) { - stage = work.getSourceOperationTask().getStageName(); + stage = work.getSourceOperationTask().getSystemName(); logger.info("Starting SourceOperationTask stage {}", stage); } else { stage = null; diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContext.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContext.java index 9353081e2b9b..74b17b52bcf3 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContext.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContext.java @@ -267,6 +267,10 @@ public final long getBacklogBytes() { return backlogBytes; } + public String getSystemName() { + return systemName; + } + public long getMaxOutputKeyBytes() { return operationalLimits.getMaxOutputKeyBytes(); } @@ -584,7 +588,7 @@ public void invalidateCache() { } catch (IOException e) { Windmill.WorkItem workItem = getWorkItem(); long shardingKey = workItem != null ? workItem.getShardingKey() : -1L; - LOG.warn("Failed to close reader for {}-{}", computationId, shardingKey, e); + LOG.warn("Failed to close reader for {}-{}", systemName, shardingKey, e); } } activeReader = null; diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java index abe5f96bb7f4..178509594660 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java @@ -265,20 +265,26 @@ public long add(WindowedValue data) throws IOException { } if (key.size() > context.getMaxOutputKeyBytes()) { if (context.throwExceptionsForLargeOutput()) { - throw new OutputTooLargeException("Key too large: " + key.size()); + throw new OutputTooLargeException( + String.format( + "Key for system %s too large: %s", context.getSystemName(), key.size())); } else { LOG.error( - "Trying to output too large key with size {}. Limit is {}. See https://cloud.google.com/dataflow/docs/guides/common-errors#key-commit-too-large-exception. Running with --experiments=throw_exceptions_on_large_output will instead throw an OutputTooLargeException which may be caught in user code.", + "Trying to output too large key for system {} with size {}. Limit is {}. See https://cloud.google.com/dataflow/docs/guides/common-errors#key-commit-too-large-exception. Running with --experiments=throw_exceptions_on_large_output will instead throw an OutputTooLargeException which may be caught in user code.", + context.getSystemName(), key.size(), context.getMaxOutputKeyBytes()); } } if (value.size() > context.getMaxOutputValueBytes()) { if (context.throwExceptionsForLargeOutput()) { - throw new OutputTooLargeException("Value too large: " + value.size()); + throw new OutputTooLargeException( + String.format( + "Value for system %s too large: %s", context.getSystemName(), value.size())); } else { LOG.error( - "Trying to output too large value with size {}. Limit is {}. See https://cloud.google.com/dataflow/docs/guides/common-errors#key-commit-too-large-exception. Running with --experiments=throw_exceptions_on_large_output will instead throw an OutputTooLargeException which may be caught in user code.", + "Trying to output too large value for system {} with size {}. Limit is {}. See https://cloud.google.com/dataflow/docs/guides/common-errors#key-commit-too-large-exception. Running with --experiments=throw_exceptions_on_large_output will instead throw an OutputTooLargeException which may be caught in user code.", + context.getSystemName(), value.size(), context.getMaxOutputValueBytes()); } diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationState.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationState.java index 8020eda1b25d..5e850d4312ea 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationState.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationState.java @@ -69,6 +69,10 @@ public String getComputationId() { return computationId; } + public String getSystemName() { + return mapTask.getSystemName(); + } + public MapTask getMapTask() { return mapTask; } diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/Commit.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/Commit.java index bbd6cfc9432b..aba9835b9e70 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/Commit.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/Commit.java @@ -66,9 +66,11 @@ public final String computationId() { return computationState().getComputationId(); } - public @Nullable WorkItemCommitRequest singleKeyRequest() { - return singleKeyRequest; - }; + public final String systemName() { + return computationState().getSystemName(); + } + + public abstract WorkItemCommitRequest request(); public ComputationState computationState() { return computationState; diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingEngineWorkCommitter.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingEngineWorkCommitter.java index 8ac9b1593c54..400b4027f184 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingEngineWorkCommitter.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingEngineWorkCommitter.java @@ -112,7 +112,12 @@ public void commit(Commit commit) { // Do this check after adding to commitQueue, else commitQueue.put() can race with // drainCommitQueue() in stop() and leave commits orphaned in the queue. if (!this.isRunning.get()) { - LOG.debug("Trying to queue commit on shutdown, failing commit={}", commit); + LOG.debug( + "Trying to queue commit on shutdown, failing commit=[systemName={}, shardingKey={}," + + " workId={} ].", + commit.systemName(), + commit.work().getShardedKey(), + commit.work().id()); drainCommitQueue(); } } diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java index 62676b44db59..b953b96dcdc9 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java @@ -171,13 +171,13 @@ private static Windmill.WorkItemCommitRequest.Builder initializeOutputBuilder( } /** Sets the stage name and workId of the Thread executing the {@link Work} for logging. */ - private static void setUpWorkLoggingContext(String workLatencyTrackingId, String computationId) { + private static void setUpWorkLoggingContext(String workLatencyTrackingId, String systemName) { setLoggingContextWorkId(workLatencyTrackingId); - setLoggingContextComputation(computationId); + setLoggingContextSystemName(systemName); } - private static void setLoggingContextComputation(@Nullable String computationId) { - DataflowWorkerLoggingMDC.setStageName(computationId); + private static void setLoggingContextSystemName(@Nullable String systemName) { + DataflowWorkerLoggingMDC.setStageName(systemName); } private static void setLoggingContextWorkId(@Nullable String workLatencyTrackingId) { @@ -227,15 +227,11 @@ public void queueAppliedFinalizeIds(ImmutableList appliedFinalizeIds) { private void processWork( ComputationState computationState, Work work, BoundedQueueExecutorWorkHandle handle) { Windmill.WorkItem workItem = work.getWorkItem(); - String computationId = computationState.getComputationId(); - LOG.debug("Starting processing for {}:\n{}", computationId, work); - setLoggingContextComputation(computationId); - KeyTransitionListener keyTransitionListener = createKeyTransitionListener(); - keyTransitionListener.onKeyTransition(null, work); - - // Before any processing starts, call any pending OnCommit callbacks. Nothing that requires - // cleanup should be done before this, since we might exit early here. - commitFinalizer.finalizeCommits(workItem.getSourceState().getFinalizeIdsList()); + String systemName = computationState.getSystemName(); + work.setProcessingThreadName(Thread.currentThread().getName()); + work.setState(Work.State.PROCESSING); + setUpWorkLoggingContext(work.getLatencyTrackingId(), systemName); + LOG.debug("Starting processing for {}:\n{}", systemName, work); if (workItem.getSourceState().getOnlyFinalize()) { handleOnlyFinalize(computationState, work, workItem); @@ -264,7 +260,21 @@ private void processWork( recordProcessingStats(workBatch, workItemCommits, executeWorkResult.stateBytesRead()); LOG.debug("Processing done for work batch size: {}", workBatch.size()); } catch (Throwable t) { - handleProcessWorkFailure(computationState, handle.getWorkBatch(), computationId, work, t); + // OutOfMemoryError that are caught will be rethrown and trigger jvm termination. + try { + workFailureProcessor.logAndProcessFailure( + systemName, + ExecutableWork.create(work, (retry, h) -> processWork(computationState, retry, h)), + t, + invalidWork -> + computationState.completeWorkAndScheduleNextWorkForKey( + invalidWork.getShardedKey(), invalidWork.id())); + } catch (OutOfMemoryError oom) { + throw oom; + } catch (Throwable t2) { + LOG.warn("Failed to process work failure safely for work {}", work.id(), t2); + throw ExceptionUtils.safeWrapThrowableAsException(t2); + } } finally { List processedWorkBatch = workBatch != null ? workBatch : ImmutableList.of(work); // Update total processing time counters. Updating in finally clause ensures that @@ -272,7 +282,7 @@ private void processWork( recordProcessingTime(stageInfo, processedWorkBatch, processingStartTimeNanos); setLoggingContextWorkId(null); - setLoggingContextComputation(null); + setLoggingContextSystemName(null); sampler.resetForWorkId(work.getLatencyTrackingId()); for (Work w : processedWorkBatch) { w.setProcessingThreadName(""); diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/failures/WorkFailureProcessor.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/failures/WorkFailureProcessor.java index 8af1840faf92..c9c44386c187 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/failures/WorkFailureProcessor.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/failures/WorkFailureProcessor.java @@ -98,40 +98,25 @@ private static boolean isOutOfMemoryError(@Nullable Throwable t) { return false; } - public void logAndProcessFailureBatch( - String computationId, - List executableWorks, - Throwable t, - Consumer onInvalidWork) + /** + * Processes failures caused by thrown exceptions that occur during execution of {@link Work}. May + * attempt to retry execution of the {@link Work} or drop it if it is invalid. + */ + public void logAndProcessFailure( + String systemName, ExecutableWork executableWork, Throwable t, Consumer onInvalidWork) throws Throwable { - List worksToRetryLocally = new java.util.ArrayList<>(); - - for (ExecutableWork executableWork : executableWorks) { - switch (evaluateRetry(computationId, executableWork.work(), t)) { - case DO_NOT_RETRY: - // Consider the item invalid. It will eventually be retried by Windmill if it still needs - // to be processed. - onInvalidWork.accept(executableWork.work()); - break; - case RETRY_LOCALLY: - // Try again after some delay and at the end of the queue to avoid a tight loop. - worksToRetryLocally.add(executableWork); - break; - case RETHROW_THROWABLE: - throw t; - } - } - - executeWithDelay(worksToRetryLocally); - } - - private void executeWithDelay(List worksToRetryLocally) { - if (!worksToRetryLocally.isEmpty()) { - // Sleep ONCE for the entire batch delay to avoid sequential thread blocks - Uninterruptibles.sleepUninterruptibly(retryLocallyDelayMs, TimeUnit.MILLISECONDS); - for (ExecutableWork ew : worksToRetryLocally) { - workUnitExecutor.forceExecute(ew, ew.work().getSerializedWorkItemSize()); - } + switch (evaluateRetry(systemName, executableWork.work(), t)) { + case DO_NOT_RETRY: + // Consider the item invalid. It will eventually be retried by Windmill if it still needs to + // be processed. + onInvalidWork.accept(executableWork.work()); + break; + case RETRY_LOCALLY: + // Try again after some delay and at the end of the queue to avoid a tight loop. + executeWithDelay(retryLocallyDelayMs, executableWork); + break; + case RETHROW_THROWABLE: + throw t; } } @@ -148,12 +133,22 @@ private enum RetryEvaluation { RETHROW_THROWABLE, } - private RetryEvaluation evaluateRetry(String computationId, Work work, Throwable t) { - if (work.isFailed()) { + private RetryEvaluation evaluateRetry(String systemName, Work work, Throwable t) { + @Nullable final Throwable cause = t.getCause(); + Throwable parsedException = (t instanceof UserCodeException && cause != null) ? cause : t; + if (KeyTokenInvalidException.isKeyTokenInvalidException(parsedException)) { + LOG.debug( + "Execution of work for system '{}' on sharding key '{}' failed due to token expiration. " + + "Work will not be retried locally.", + systemName, + work.getWorkItem().getShardingKey()); + return RetryEvaluation.DO_NOT_RETRY; + } + if (WorkItemCancelledException.isWorkItemCancelledException(parsedException)) { LOG.debug( - "Execution of work for computation '{}' on sharding key '{}' failed. " - + "Work is already marked as failed, not retrying locally.", - computationId, + "Execution of work for system '{}' on sharding key '{}' failed. " + + "Work will not be retried locally.", + systemName, work.getWorkItem().getShardingKey()); return RetryEvaluation.DO_NOT_RETRY; } @@ -166,30 +161,30 @@ private RetryEvaluation evaluateRetry(String computationId, Work work, Throwable if (isOutOfMemoryError(parsedException)) { String heapDump = tryToDumpHeap(); LOG.error( - "Execution of work for computation '{}' for sharding key '{}' failed with out-of-memory. " + "Execution of work for system '{}' for sharding key '{}' failed with out-of-memory. " + "Work will not be retried locally. Heap dump {}.", - computationId, + systemName, work.getWorkItem().getShardingKey(), heapDump, parsedException); return RetryEvaluation.RETHROW_THROWABLE; } - if (!failureTracker.trackFailure(computationId, work.getWorkItem(), parsedException)) { + if (!failureTracker.trackFailure(systemName, work.getWorkItem(), parsedException)) { LOG.error( - "Execution of work for computation '{}' on sharding key '{}' failed with uncaught exception, " + "Execution of work for system '{}' on sharding key '{}' failed with uncaught exception, " + "and Windmill indicated not to retry locally.", - computationId, + systemName, work.getWorkItem().getShardingKey(), parsedException); return RetryEvaluation.DO_NOT_RETRY; } if (elapsedTimeSinceStart.isLongerThan(MAX_LOCAL_PROCESSING_RETRY_DURATION)) { LOG.error( - "Execution of work for computation '{}' for sharding key '{}' failed with uncaught exception, " + "Execution of work for system '{}' for sharding key '{}' failed with uncaught exception, " + "and it will not be retried locally because the elapsed time since start {} " + "exceeds {}.", - computationId, + systemName, work.getWorkItem().getShardingKey(), elapsedTimeSinceStart, MAX_LOCAL_PROCESSING_RETRY_DURATION, @@ -197,9 +192,9 @@ private RetryEvaluation evaluateRetry(String computationId, Work work, Throwable return RetryEvaluation.DO_NOT_RETRY; } LOG.error( - "Execution of work for computation '{}' on sharding key '{}' failed with uncaught exception. " + "Execution of work for system '{}' on sharding key '{}' failed with uncaught exception. " + "Work will be retried locally.", - computationId, + systemName, work.getWorkItem().getShardingKey(), parsedException); return RetryEvaluation.RETRY_LOCALLY; From 1d151cfe606159066c2c4e43d748570ea35f1150 Mon Sep 17 00:00:00 2001 From: Ryan Wigglesworth Date: Thu, 30 Jul 2026 20:55:01 +0000 Subject: [PATCH 24/26] Switch logging to fused stage name in more places. This simplifies operations by removing backend details --- .../worker/StreamingDataflowWorker.java | 27 +++++++++---------- .../worker/streaming/ActiveWorkState.java | 2 +- .../streaming/ComputationStateCache.java | 12 ++++++--- .../harness/MetricsDataProvider.java | 2 +- .../windmill/state/WindmillStateCache.java | 13 ++++++--- .../ComputationWorkExecutorFactory.java | 7 ++--- .../processing/StreamingWorkScheduler.java | 4 +-- .../StreamingModeExecutionContextTest.java | 2 +- .../worker/WorkerCustomSourcesTest.java | 2 +- .../streaming/ComputationStateCacheTest.java | 5 +++- .../state/WindmillStateInternalsTest.java | 8 +++--- .../work/refresh/ActiveWorkRefresherTest.java | 2 +- 12 files changed, 49 insertions(+), 37 deletions(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java index 2339430464c7..995ab7d778a6 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java @@ -924,7 +924,7 @@ static StreamingDataflowWorker forTesting( mapTask, workExecutor, stateNameMap, - stateCache.forComputation(mapTask.getStageName()))); + stateCache.forComputation(mapTask.getStageName(), mapTask.getSystemName()))); MemoryMonitor memoryMonitor = MemoryMonitor.fromOptions(options); FailureTracker failureTracker = options.isEnableStreamingEngine() @@ -1197,26 +1197,23 @@ void stop() { } private void onCompleteCommit(CompleteCommit completeCommit) { + Optional computationState = + computationStateCache.getIfPresent(completeCommit.computationId()); if (completeCommit.status() != Windmill.CommitStatus.OK) { readerCache.invalidateReader( WindmillComputationKey.create( completeCommit.computationId(), completeCommit.shardedKey())); - stateCache - .forComputation(completeCommit.computationId()) - .invalidate(completeCommit.shardedKey()); + computationState.ifPresent( + state -> + stateCache + .forComputation(completeCommit.computationId(), state.getSystemName()) + .invalidate(completeCommit.shardedKey())); } - computationStateCache - .getIfPresent(completeCommit.computationId()) - .ifPresent( - state -> { - if (completeCommit.retryableFailure()) { - state.reexecuteActiveWork(completeCommit.shardedKey(), completeCommit.workId()); - } else { - state.completeWorkAndScheduleNextWorkForKey( - completeCommit.shardedKey(), completeCommit.workId()); - } - }); + computationState.ifPresent( + state -> + state.completeWorkAndScheduleNextWorkForKey( + completeCommit.shardedKey(), completeCommit.workId())); } @AutoValue diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ActiveWorkState.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ActiveWorkState.java index de4082581293..f0150cf73eb3 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ActiveWorkState.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ActiveWorkState.java @@ -185,7 +185,7 @@ synchronized void failWorkForKey(ImmutableList failedWork executableWork.work().setFailed(); LOG.debug( "Failing work {} {}. The work will be retried and is not lost.", - computationStateCache.getComputation(), + computationStateCache.getSystemName(), failedId); } } diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCache.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCache.java index 4b4acb73f4a7..e6f902a65bdc 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCache.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCache.java @@ -28,6 +28,7 @@ import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ExecutionException; +import java.util.function.BiFunction; import java.util.function.Function; import javax.annotation.concurrent.ThreadSafe; import org.apache.beam.runners.dataflow.worker.apiary.FixMultiOutputInfosOnParDoInstructions; @@ -77,7 +78,8 @@ private ComputationStateCache( public static ComputationStateCache create( ComputationConfig.Fetcher computationConfigFetcher, BoundedQueueExecutor workUnitExecutor, - Function perComputationStateCacheViewFactory, + BiFunction + perComputationStateCacheViewFactory, IdGenerator idGenerator) { Function fixMultiOutputInfosOnParDoInstructions = new FixMultiOutputInfosOnParDoInstructions(idGenerator); @@ -105,7 +107,8 @@ public ComputationState load(String computationId) { fixMultiOutputInfosOnParDoInstructions.apply(computationConfig.mapTask()), workUnitExecutor, transformUserNameToStateFamilyForComputation, - perComputationStateCacheViewFactory.apply(computationId)); + perComputationStateCacheViewFactory.apply( + computationId, computationConfig.mapTask().getSystemName())); } }), fixMultiOutputInfosOnParDoInstructions, @@ -116,7 +119,8 @@ public ComputationState load(String computationId) { public static ComputationStateCache forTesting( ComputationConfig.Fetcher computationConfigFetcher, BoundedQueueExecutor workUnitExecutor, - Function perComputationStateCacheViewFactory, + BiFunction + perComputationStateCacheViewFactory, IdGenerator idGenerator, ConcurrentMap pipelineUserNameToStateFamilyNameMap) { ComputationStateCache cache = @@ -205,7 +209,7 @@ public void closeAndInvalidateAll() { public void appendSummaryHtml(PrintWriter writer) { writer.println("

Specs

"); for (ComputationState computationState : getAllPresentComputations()) { - writer.println("

" + computationState.getComputationId() + "

"); + writer.println("

" + computationState.getSystemName() + "

"); writer.print(""); diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/harness/MetricsDataProvider.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/harness/MetricsDataProvider.java index 901e2d235f85..0580b7a0b05b 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/harness/MetricsDataProvider.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/harness/MetricsDataProvider.java @@ -59,7 +59,7 @@ public void appendSummaryHtml(PrintWriter writer) { writer.println("Active Keys:
"); for (ComputationState computationState : allComputationStates.get()) { - writer.print(computationState.getComputationId()); + writer.print(computationState.getSystemName()); writer.print(":
"); computationState.printActiveWork(writer); writer.println("
"); diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateCache.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateCache.java index 7515db000852..ff62d12a8fa2 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateCache.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateCache.java @@ -170,8 +170,8 @@ public CacheStats getCacheStats() { } /** Returns a per-computation view of the state cache. */ - public ForComputation forComputation(String computation) { - return new ForComputation(computation); + public ForComputation forComputation(String computation, String systemName) { + return new ForComputation(computation, systemName); } /** Print summary statistics of the cache to the given {@link PrintWriter}. */ @@ -353,9 +353,11 @@ private Optional value() { public class ForComputation { private final String computation; + private final String systemName; - private ForComputation(String computation) { + private ForComputation(String computation, String systemName) { this.computation = computation; + this.systemName = systemName; } /** Returns the computation associated to this class. */ @@ -363,6 +365,11 @@ public String getComputation() { return this.computation; } + /** Returns the system name associated to this class. */ + public String getSystemName() { + return this.systemName; + } + /** Invalidate all cache entries for this computation and {@code processingKey}. */ public void invalidate(ByteString processingKey, long shardingKey) { WindmillComputationKey key = diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/ComputationWorkExecutorFactory.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/ComputationWorkExecutorFactory.java index 5bfc8bd4998d..f0e5ab019420 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/ComputationWorkExecutorFactory.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/ComputationWorkExecutorFactory.java @@ -20,6 +20,7 @@ import static org.apache.beam.runners.dataflow.DataflowRunner.hasExperiment; import com.google.api.services.dataflow.model.MapTask; +import java.util.function.BiFunction; import java.util.function.Function; import org.apache.beam.runners.dataflow.internal.CustomSources; import org.apache.beam.runners.dataflow.options.DataflowWorkerHarnessOptions; @@ -82,7 +83,7 @@ final class ComputationWorkExecutorFactory { private final DataflowWorkerHarnessOptions options; private final DataflowMapTaskExecutorFactory mapTaskExecutorFactory; private final ReaderCache readerCache; - private final Function stateCacheFactory; + private final BiFunction stateCacheFactory; private final ReaderRegistry readerRegistry; private final SinkRegistry sinkRegistry; private final DataflowExecutionStateSampler sampler; @@ -111,7 +112,7 @@ final class ComputationWorkExecutorFactory { DataflowWorkerHarnessOptions options, DataflowMapTaskExecutorFactory mapTaskExecutorFactory, ReaderCache readerCache, - Function stateCacheFactory, + BiFunction stateCacheFactory, DataflowExecutionStateSampler sampler, StreamingCounters streamingCounters, FailureTracker failureTracker, @@ -286,7 +287,7 @@ private StreamingModeExecutionContext createExecutionContext( computationId, readerCache, computationState.getTransformUserNameToStateFamily(), - stateCacheFactory.apply(computationId), + stateCacheFactory.apply(computationId, stageInfo.systemName()), stageInfo.metricsContainerRegistry(), executionStateTracker, stageInfo.executionStateRegistry(), diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java index b953b96dcdc9..2ac3ebb706a9 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java @@ -26,7 +26,7 @@ import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; -import java.util.function.Function; +import java.util.function.BiFunction; import java.util.function.Supplier; import javax.annotation.concurrent.ThreadSafe; import org.apache.beam.repackaged.core.org.apache.commons.lang3.tuple.Pair; @@ -115,7 +115,7 @@ public static StreamingWorkScheduler create( DataflowMapTaskExecutorFactory mapTaskExecutorFactory, BoundedQueueExecutor workExecutor, ScheduledExecutorService commitFinalizerCleanupExecutor, - Function stateCacheFactory, + BiFunction stateCacheFactory, FailureTracker failureTracker, WorkFailureProcessor workFailureProcessor, StreamingCounters streamingCounters, diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContextTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContextTest.java index c5efcea4e47c..6f10f6e3749f 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContextTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContextTest.java @@ -134,7 +134,7 @@ private StreamingModeExecutionContext createExecutionContext( WindmillStateCache.builder() .setSizeMb(options.getWorkerCacheMb()) .build() - .forComputation("comp"), + .forComputation("comp", "systemName"), StreamingStepMetricsContainer.createRegistry(), new DataflowExecutionStateTracker( ExecutionStateSampler.newForTest(), diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WorkerCustomSourcesTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WorkerCustomSourcesTest.java index f3d1935598c4..27b11ad67c6d 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WorkerCustomSourcesTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WorkerCustomSourcesTest.java @@ -1004,7 +1004,7 @@ public void testFailedWorkItemsAbort() throws Exception { WindmillStateCache.builder() .setSizeMb(options.getWorkerCacheMb()) .build() - .forComputation(COMPUTATION_ID), + .forComputation(COMPUTATION_ID, "systemName"), StreamingStepMetricsContainer.createRegistry(), new DataflowExecutionStateTracker( ExecutionStateSampler.newForTest(), diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCacheTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCacheTest.java index f57e20d4b5fb..6785ce47d0f6 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCacheTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCacheTest.java @@ -86,7 +86,10 @@ private static ExecutableWork createWork(ShardedKey shardedKey, long workToken, public void setUp() { computationStateCache = ComputationStateCache.create( - configFetcher, workExecutor, ignored -> stateCache, IdGenerators.decrementingLongs()); + configFetcher, + workExecutor, + (ignored1, ignored2) -> stateCache, + IdGenerators.decrementingLongs()); } @Test 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 87b746089f11..0b55a5119564 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 @@ -225,7 +225,7 @@ public void resetUnderTest() { mockReader, false, cache - .forComputation("comp") + .forComputation("comp", "systemName") .forKey( WindmillComputationKey.create( "comp", ByteString.copyFrom("dummyKey", StandardCharsets.UTF_8), 123), @@ -241,7 +241,7 @@ public void resetUnderTest() { mockReader, true, cache - .forComputation("comp") + .forComputation("comp", "systemName") .forKey( WindmillComputationKey.create( "comp", ByteString.copyFrom("dummyNewKey", StandardCharsets.UTF_8), 123), @@ -257,7 +257,7 @@ public void resetUnderTest() { mockReader, false, cacheViaMultimap - .forComputation("comp") + .forComputation("comp", "systemName") .forKey( WindmillComputationKey.create( "comp", ByteString.copyFrom("dummyNewKey", StandardCharsets.UTF_8), 123), @@ -2049,7 +2049,7 @@ false, key(NAMESPACE, tag), STATE_FAMILY, VarIntCoder.of())) // clear cache and recreate multimapState cache - .forComputation("comp") + .forComputation("comp", "systemName") .invalidate(ByteString.copyFrom("dummyKey", StandardCharsets.UTF_8), 123); resetUnderTest(); multimapState = underTest.state(NAMESPACE, addr); diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/work/refresh/ActiveWorkRefresherTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/work/refresh/ActiveWorkRefresherTest.java index caa25bf83090..e711a780a4dc 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/work/refresh/ActiveWorkRefresherTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/work/refresh/ActiveWorkRefresherTest.java @@ -257,7 +257,7 @@ public void testInvalidateStuckCommits() throws InterruptedException { ByteString key = ByteString.EMPTY; for (int i = 0; i < 5; i++) { WindmillStateCache.ForComputation perComputationStateCache = - spy(stateCache.forComputation(COMPUTATION_ID_PREFIX + i)); + spy(stateCache.forComputation(COMPUTATION_ID_PREFIX + i, "systemName" + i)); ComputationState computationState = spy(createComputationState(i, perComputationStateCache)); ExecutableWork fakeWork = createOldWork(ShardedKey.create(key, i), i, ignored -> {}); fakeWork.work().setState(Work.State.COMMITTING); From 7d1fc11c44da0f62b14af8b58c750f4abb324110 Mon Sep 17 00:00:00 2001 From: "dependabot[bot]" <49699333+dependabot[bot]@users.noreply.github.com> Date: Thu, 30 Jul 2026 07:35:54 -0400 Subject: [PATCH 25/26] Bump cloud.google.com/go/spanner from 1.93.0 to 1.94.0 in /sdks (#39550) Bumps [cloud.google.com/go/spanner](https://github.com/googleapis/google-cloud-go) from 1.93.0 to 1.94.0. - [Release notes](https://github.com/googleapis/google-cloud-go/releases) - [Changelog](https://github.com/googleapis/google-cloud-go/blob/main/CHANGES.md) - [Commits](https://github.com/googleapis/google-cloud-go/compare/spanner/v1.93.0...spanner/v1.94.0) --- updated-dependencies: - dependency-name: cloud.google.com/go/spanner dependency-version: 1.94.0 dependency-type: direct:production update-type: version-update:semver-minor ... Signed-off-by: dependabot[bot] Co-authored-by: dependabot[bot] <49699333+dependabot[bot]@users.noreply.github.com> --- sdks/go.mod | 18 +++++++++--------- 1 file changed, 9 insertions(+), 9 deletions(-) diff --git a/sdks/go.mod b/sdks/go.mod index 421bc3444b6c..8ee203f77949 100644 --- a/sdks/go.mod +++ b/sdks/go.mod @@ -29,15 +29,15 @@ require ( cloud.google.com/go/bigtable v1.50.0 cloud.google.com/go/datastore v1.25.0 cloud.google.com/go/profiler v0.6.0 - cloud.google.com/go/pubsub v1.50.4 - cloud.google.com/go/spanner v1.92.0 - cloud.google.com/go/storage v1.63.1 - github.com/aws/aws-sdk-go-v2 v1.42.1 - github.com/aws/aws-sdk-go-v2/config v1.32.30 - github.com/aws/aws-sdk-go-v2/credentials v1.19.29 - github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.22.33 - github.com/aws/aws-sdk-go-v2/service/s3 v1.105.1 - github.com/aws/smithy-go v1.27.3 + cloud.google.com/go/pubsub v1.51.0 + cloud.google.com/go/spanner v1.94.0 + cloud.google.com/go/storage v1.64.0 + github.com/aws/aws-sdk-go-v2 v1.43.1 + github.com/aws/aws-sdk-go-v2/config v1.32.32 + github.com/aws/aws-sdk-go-v2/credentials v1.19.31 + github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.22.36 + github.com/aws/aws-sdk-go-v2/service/s3 v1.106.1 + github.com/aws/smithy-go v1.27.5 github.com/docker/go-connections v0.7.0 // indirect github.com/dustin/go-humanize v1.0.1 github.com/go-sql-driver/mysql v1.10.0 From 8049942da35667eeacc4f74375bae4feff7edcd2 Mon Sep 17 00:00:00 2001 From: Ryan Wigglesworth Date: Thu, 30 Jul 2026 21:40:39 +0000 Subject: [PATCH 26/26] confusion --- sdks/go.mod | 12 +++++++++--- 1 file changed, 9 insertions(+), 3 deletions(-) diff --git a/sdks/go.mod b/sdks/go.mod index 8ee203f77949..c2e62141d6b8 100644 --- a/sdks/go.mod +++ b/sdks/go.mod @@ -32,9 +32,9 @@ require ( cloud.google.com/go/pubsub v1.51.0 cloud.google.com/go/spanner v1.94.0 cloud.google.com/go/storage v1.64.0 - github.com/aws/aws-sdk-go-v2 v1.43.1 - github.com/aws/aws-sdk-go-v2/config v1.32.32 - github.com/aws/aws-sdk-go-v2/credentials v1.19.31 + github.com/aws/aws-sdk-go-v2 v1.43.2 + github.com/aws/aws-sdk-go-v2/config v1.32.33 + github.com/aws/aws-sdk-go-v2/credentials v1.19.32 github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.22.36 github.com/aws/aws-sdk-go-v2/service/s3 v1.106.1 github.com/aws/smithy-go v1.27.5 @@ -153,9 +153,15 @@ require ( github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.33 // indirect github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.34 // indirect github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.14 // indirect +<<<<<<< HEAD github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.26 // indirect github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.33 // indirect github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.34 // indirect +======= + github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.25 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.33 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.33 // indirect +>>>>>>> c6e0f5b630e (Bump github.com/aws/aws-sdk-go-v2/config in /sdks (#39551)) github.com/aws/aws-sdk-go-v2/service/sso v1.33.2 // indirect github.com/aws/aws-sdk-go-v2/service/ssooidc v1.38.2 // indirect github.com/aws/aws-sdk-go-v2/service/sts v1.45.2 // indirect