From c831c3363077e48942a661f02faa3893fe09cd3e Mon Sep 17 00:00:00 2001 From: Ryan Wigglesworth Date: Tue, 14 Jul 2026 00:22:40 +0000 Subject: [PATCH 01/21] 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 | 25 ++++++-- .../processing/StreamingWorkScheduler.java | 30 ++++++--- .../worker/StreamingDataflowWorkerTest.java | 61 +++++++++++++------ 3 files changed, 80 insertions(+), 36 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 4d6ae8a208c1..227dfa235218 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,17 +17,30 @@ */ package org.apache.beam.runners.dataflow.worker.streaming; +import com.google.protobuf.TextFormat; import org.apache.beam.runners.dataflow.worker.windmill.Windmill; public final class KeyCommitTooLargeException extends Exception { public static KeyCommitTooLargeException causedBy( - String computationId, long byteLimit, Windmill.WorkItemCommitRequest request) { + String stageName, long byteLimit, Windmill.WorkItemCommitRequest request) { + return causedBy(stageName, byteLimit, request, false); + } + + public static KeyCommitTooLargeException causedBy( + String stageName, + long byteLimit, + Windmill.WorkItemCommitRequest request, + boolean hotKeyLoggingEnabled) { StringBuilder message = new StringBuilder(); message.append("Commit request for stage "); - message.append(computationId); + message.append(stageName); message.append(" and sharding key "); - message.append(request.getShardingKey()); + message.append(Long.toUnsignedString(request.getShardingKey())); + if (hotKeyLoggingEnabled && !request.getKey().isEmpty()) { + message.append(" and key "); + message.append(TextFormat.escapeBytes(request.getKey())); + } if (request.getSerializedSize() > 0) { message.append( " has size " @@ -38,9 +51,8 @@ public static KeyCommitTooLargeException causedBy( message.append(" is larger than 2GB and cannot be processed"); } message.append( - ". This may be caused by grouping a very " - + "large amount of data in a single window without using Combine," - + " or by producing a large amount of data from a single input element." + ". This may be caused by grouping a very large amount of data in a single window without" + + " using Combine, or by producing a large amount of data from a single input element." + " See https://cloud.google.com/dataflow/docs/guides/common-errors#key-commit-too-large-exception."); return new KeyCommitTooLargeException(message.toString()); } @@ -49,3 +61,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 e9b85d720d2b..eb713a35d150 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,10 +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 com.google.common.collect.ImmutableList; +import com.google.protobuf.ByteString; import java.util.List; import java.util.Map; import java.util.concurrent.ConcurrentMap; @@ -29,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; @@ -61,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; @@ -89,6 +90,7 @@ public class StreamingWorkScheduler { private final DataflowExecutionStateSampler sampler; private final StreamingGlobalConfigHandle globalConfigHandle; private final BoundedQueueExecutor workExecutor; + private final boolean hotKeyLoggingEnabled; public StreamingWorkScheduler( Supplier clock, @@ -100,7 +102,8 @@ public StreamingWorkScheduler( StreamingCounters streamingCounters, ConcurrentMap stageInfoMap, DataflowExecutionStateSampler sampler, - StreamingGlobalConfigHandle globalConfigHandle) { + StreamingGlobalConfigHandle globalConfigHandle, + boolean hotKeyLoggingEnabled) { this.clock = clock; this.workExecutor = workExecutor; this.computationWorkExecutorFactory = computationWorkExecutorFactory; @@ -111,6 +114,7 @@ public StreamingWorkScheduler( this.stageInfoMap = stageInfoMap; this.sampler = sampler; this.globalConfigHandle = globalConfigHandle; + this.hotKeyLoggingEnabled = hotKeyLoggingEnabled; } public static StreamingWorkScheduler create( @@ -145,6 +149,9 @@ public static StreamingWorkScheduler create( hotKeyLogger, sideInputStateFetcherFactory); + boolean hotKeyLoggingEnabled = + options.isHotKeyLoggingEnabled() || hasExperiment(options, "enable_hot_key_logging"); + return new StreamingWorkScheduler( clock, workExecutor, @@ -155,7 +162,8 @@ public static StreamingWorkScheduler create( streamingCounters, stageInfoMap, sampler, - globalConfigHandle); + globalConfigHandle, + hotKeyLoggingEnabled); } private static long computeShuffleBytesRead(Windmill.WorkItem workItem) { @@ -307,7 +315,7 @@ private void processWork( private Windmill.WorkItemCommitRequest validateCommitRequestSize( Windmill.WorkItemCommitRequest commitRequest, - String computationId, + String stageName, Windmill.WorkItem workItem) { long byteLimit = globalConfigHandle.getConfig().operationalLimits().getMaxWorkItemCommitBytes(); int commitSize = commitRequest.getSerializedSize(); @@ -322,8 +330,9 @@ private Windmill.WorkItemCommitRequest validateCommitRequestSize( } KeyCommitTooLargeException e = - KeyCommitTooLargeException.causedBy(computationId, byteLimit, commitRequest); - failureTracker.trackFailure(computationId, workItem, 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. @@ -451,7 +460,7 @@ private void commitSingleKeyWork( // Validate the commit request, possibly requesting truncation if the commitSize is too large. Windmill.WorkItemCommitRequest validatedCommitRequest = validateCommitRequestSize( - commitRequest, computationState.getComputationId(), work.getWorkItem()); + commitRequest, computationState.getMapTask().getSystemName(), work.getWorkItem()); work.setState(Work.State.COMMIT_QUEUED); validatedCommitRequest = validatedCommitRequest @@ -509,3 +518,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 23730bc57705..43eda2ae8106 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; @@ -187,19 +200,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; @@ -1290,8 +1290,9 @@ public void testKeyTokenInvalidException() throws Exception { worker.stop(); } - @Test - public void testKeyCommitTooLargeException() throws Exception { + private void runKeyCommitTooLargeExceptionTest( + StreamingDataflowWorkerTestParams.Builder workerParams, boolean expectKeyInErrorMessage) + throws Exception { KvCoder kvCoder = KvCoder.of(StringUtf8Coder.of(), StringUtf8Coder.of()); List instructions = @@ -1304,7 +1305,7 @@ public void testKeyCommitTooLargeException() throws Exception { StreamingDataflowWorker worker = makeWorker( - defaultWorkerParams() + workerParams .setInstructions(instructions) .setStreamingGlobalConfig( StreamingGlobalConfig.builder() @@ -1337,18 +1338,14 @@ public void testKeyCommitTooLargeException() throws Exception { .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()); @@ -1360,12 +1357,35 @@ public void testKeyCommitTooLargeException() throws Exception { foundErrors = true; String errorMessage = status.getErrors().get(0).getMessage(); assertThat(errorMessage, Matchers.containsString("KeyCommitTooLargeException")); + assertThat(errorMessage, Matchers.containsString("Commit request for stage computation")); + if (expectKeyInErrorMessage) { + assertThat(errorMessage, Matchers.containsString("and key large_key")); + } else { + assertThat(errorMessage, Matchers.not(Matchers.containsString("and key large_key"))); + } } } assertTrue(foundErrors); worker.stop(); } + @Test + public void testKeyCommitTooLargeException() throws Exception { + runKeyCommitTooLargeExceptionTest(defaultWorkerParams(), /* expectKeyInErrorMessage= */ false); + } + + @Test + public void testKeyCommitTooLargeException_withHotKeyLoggingEnabled() throws Exception { + runKeyCommitTooLargeExceptionTest( + defaultWorkerParams("--hotKeyLoggingEnabled=true"), /* expectKeyInErrorMessage= */ true); + } + + @Test + public void testKeyCommitTooLargeException_withHotKeyLoggingDisabled() throws Exception { + runKeyCommitTooLargeExceptionTest( + defaultWorkerParams("--hotKeyLoggingEnabled=false"), /* expectKeyInErrorMessage= */ false); + } + @Test public void testOutputKeyTooLargeException() throws Exception { KvCoder kvCoder = KvCoder.of(StringUtf8Coder.of(), StringUtf8Coder.of()); @@ -5105,3 +5125,4 @@ final Builder publishCounters() { } } } + From 25d7bbc34189611eabc781bab6f98a82cec81cbd Mon Sep 17 00:00:00 2001 From: Ryan Wigglesworth Date: Tue, 14 Jul 2026 11:01:49 -0700 Subject: [PATCH 02/21] 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 227dfa235218..59b2d050d957 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; public final class KeyCommitTooLargeException extends Exception { From 2513ffebcb01177da6006a4e56c6de9cb5919c34 Mon Sep 17 00:00:00 2001 From: Ryan Wigglesworth Date: Tue, 14 Jul 2026 11:03:00 -0700 Subject: [PATCH 03/21] 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 eb713a35d150..7ad95ee7958c 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 647a81542c8065daca4a9b43dfebacc95de06779 Mon Sep 17 00:00:00 2001 From: Ryan Wigglesworth Date: Tue, 14 Jul 2026 11:09:18 -0700 Subject: [PATCH 04/21] 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 43eda2ae8106..e73bbbfc6dc3 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 369aff8a9d0fe3cd6680eabab56ecac5ad7af846 Mon Sep 17 00:00:00 2001 From: Ryan Wigglesworth Date: Tue, 14 Jul 2026 20:26:38 +0000 Subject: [PATCH 05/21] 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 59b2d050d957..227dfa235218 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; 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 7ad95ee7958c..84457942de06 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; @@ -149,7 +149,10 @@ public static StreamingWorkScheduler create( sideInputStateFetcherFactory); 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 e73bbbfc6dc3..31aceb012eac 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; @@ -200,6 +187,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 15751e2636bd7a00070115e193f6e4729cfb7b5c Mon Sep 17 00:00:00 2001 From: Ryan Wigglesworth Date: Tue, 14 Jul 2026 23:30:06 +0000 Subject: [PATCH 06/21] 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 227dfa235218..59b2d050d957 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; 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 84457942de06..0a73c70843d9 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 31aceb012eac..645a3447562f 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 @@ -187,19 +187,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 d3a4d1c6e7d2de41556f70f391fa28093c707eec Mon Sep 17 00:00:00 2001 From: Ryan Wigglesworth Date: Tue, 14 Jul 2026 23:32:55 +0000 Subject: [PATCH 07/21] 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 0a73c70843d9..7fb37edf6162 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 d5babb02d6051888474aa1bd6efe28529a0192b2 Mon Sep 17 00:00:00 2001 From: Ryan Wigglesworth Date: Wed, 15 Jul 2026 04:20:03 +0000 Subject: [PATCH 08/21] 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 7fb37edf6162..dbb5bc391cb9 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 @@ -148,11 +148,11 @@ public static StreamingWorkScheduler create( hotKeyLogger, sideInputStateFetcherFactory); + 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 9f611b500893f95f725b951ccbce0a5f1fba6b58 Mon Sep 17 00:00:00 2001 From: Ryan Wigglesworth Date: Wed, 15 Jul 2026 07:02:07 +0000 Subject: [PATCH 09/21] one more attempt --- .../dataflow/worker/streaming/KeyCommitTooLargeException.java | 2 +- .../windmill/work/processing/StreamingWorkScheduler.java | 4 +--- .../runners/dataflow/worker/StreamingDataflowWorkerTest.java | 1 - 3 files changed, 2 insertions(+), 5 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 59b2d050d957..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,8 +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.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 dbb5bc391cb9..4947d3769e88 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 @@ -316,9 +316,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 645a3447562f..fc8f12c45cb3 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 @@ -5125,4 +5125,3 @@ final Builder publishCounters() { } } } - From a5bd7d0d38ad441d3d1e3bc94eae9263a7e44db6 Mon Sep 17 00:00:00 2001 From: Ryan Wigglesworth Date: Wed, 15 Jul 2026 17:36:55 +0000 Subject: [PATCH 10/21] 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 4947d3769e88..45b9fbe93396 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 @@ -518,4 +518,3 @@ static ExecuteWorkResult create( abstract long stateBytesRead(); } } - From 2655ca53ab1c350a8d5d26c0fdc9c59fa52d85b9 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/21] 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, 1 insertion(+), 1 deletion(-) diff --git a/.github/workflows/beam_PostCommit_Java_Delta_IO_Dataflow.yml b/.github/workflows/beam_PostCommit_Java_Delta_IO_Dataflow.yml index 604be67d90fa..433759424509 100644 --- a/.github/workflows/beam_PostCommit_Java_Delta_IO_Dataflow.yml +++ b/.github/workflows/beam_PostCommit_Java_Delta_IO_Dataflow.yml @@ -63,7 +63,7 @@ jobs: job_name: ["beam_PostCommit_Java_Delta_IO_Dataflow"] job_phrase: ["Run PostCommit Java Delta IO Dataflow"] steps: - - uses: actions/checkout@v6 + - uses: actions/checkout@v7 - name: Setup repository uses: ./.github/actions/setup-action with: From c3dfad90fbf11c476169dde4075613addfc211f9 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/21] 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 | 2 +- sdks/go.sum | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/sdks/go.mod b/sdks/go.mod index ceb23f1ae046..6ae70cab524d 100644 --- a/sdks/go.mod +++ b/sdks/go.mod @@ -27,7 +27,7 @@ toolchain go1.26.2 require ( cloud.google.com/go/bigquery v1.79.0 cloud.google.com/go/bigtable v1.50.0 - cloud.google.com/go/datastore v1.24.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 diff --git a/sdks/go.sum b/sdks/go.sum index 7360fcd17dad..1434dab988e9 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.24.0 h1:auNUPJTT9gFcHNj2iKOEeE23nrjf7dE7VA6TO3jw8h0= -cloud.google.com/go/datastore v1.24.0/go.mod h1:cEkLhU6Ti/gauQ7DFrUrG8bQjiMIxi++b5ePiThi5So= +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 66224976972ffe4ae68797b48eee370e8be79688 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/21] 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 | 2 +- sdks/go.sum | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/sdks/go.mod b/sdks/go.mod index 6ae70cab524d..fcaafb670e41 100644 --- a/sdks/go.mod +++ b/sdks/go.mod @@ -35,7 +35,7 @@ require ( 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.32 + 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 diff --git a/sdks/go.sum b/sdks/go.sum index 1434dab988e9..784d9c54dd94 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.30 h1:/hi1JADLEW9YYryEz1w4GQ github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.18.30/go.mod h1:/3AOgy4K17Dm4ucMZVC/MJkzy5kmfKUcINRHZyo0koQ= 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.32 h1:kUb0wd0/NfYv2RDoDfogxBy/Hevby5yLIL12iCcs0hY= -github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.22.32/go.mod h1:+rUx79uJZfEavbKROY+U1ez7bUMie7oOAdrjgL0C1FQ= +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.30 h1:xM/Is9cKMHa8Jj8zkvWhvrFkZsXJV9E+BB4g0HW0duQ= From 4e1732bfabc34b0483bbabec67dd2bf3da1a60eb Mon Sep 17 00:00:00 2001 From: Yi Hu Date: Wed, 15 Jul 2026 08:02:11 -0400 Subject: [PATCH 14/21] Fix DataflowV1 test failure by fixing getSimpleName access (#39330) --- .../beam_PostCommit_Java_DataflowV1.json | 2 +- .../sdk/transforms/display/DisplayData.java | 7 +++--- .../transforms/reflect/DoFnSignatures.java | 3 ++- .../beam/sdk/util/common/ReflectHelpers.java | 22 +++++++++++++++++-- 4 files changed, 27 insertions(+), 7 deletions(-) diff --git a/.github/trigger_files/beam_PostCommit_Java_DataflowV1.json b/.github/trigger_files/beam_PostCommit_Java_DataflowV1.json index ae6cb268ff68..aff502e19620 100644 --- a/.github/trigger_files/beam_PostCommit_Java_DataflowV1.json +++ b/.github/trigger_files/beam_PostCommit_Java_DataflowV1.json @@ -1,5 +1,5 @@ { - "https://github.com/apache/beam/pull/36138": "Cleanly separating v1 worker and v2 sdk harness container image handling", + "https://github.com/apache/beam/pull/39330": "Fix DataflowV1 test failure by fixing getSimpleName access", "https://github.com/apache/beam/pull/34902": "Introducing OutputBuilder", "https://github.com/apache/beam/pull/35177": "Introducing WindowedValueReceiver to runners", "comment": "Modify this file in a trivial way to cause this test suite to run", diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/display/DisplayData.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/display/DisplayData.java index 3be59e556158..dcc84ad8e1a0 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/display/DisplayData.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/display/DisplayData.java @@ -33,6 +33,7 @@ import java.util.Set; import org.apache.beam.sdk.options.ValueProvider; import org.apache.beam.sdk.transforms.PTransform; +import org.apache.beam.sdk.util.common.ReflectHelpers; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Joiner; 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; @@ -617,7 +618,7 @@ FormattedItemValue format(Object value) { @Override FormattedItemValue format(Object value) { Class clazz = checkType(value, Class.class, JAVA_CLASS); - return new FormattedItemValue(clazz.getName(), clazz.getSimpleName()); + return new FormattedItemValue(clazz.getName(), ReflectHelpers.getSimpleName(clazz)); } }; @@ -753,10 +754,10 @@ private Builder include(Path path, HasDisplayData subComponent) { // Common case: AutoValue classes such as AutoValue_FooIO_Read. It's more useful // to show the user the FooIO.Read class, which is the direct superclass of the AutoValue // generated class. - if (namespace.getSimpleName().startsWith("AutoValue_")) { + if (ReflectHelpers.getSimpleName(namespace).startsWith("AutoValue_")) { namespace = namespace.getSuperclass(); } - if (namespace.isSynthetic() && namespace.getSimpleName().contains("$$Lambda")) { + if (namespace.isSynthetic() && ReflectHelpers.getSimpleName(namespace).contains("$$Lambda")) { try { String className = namespace.getCanonicalName(); // A local class, local interface, or anonymous class does not have a canonical name. diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/reflect/DoFnSignatures.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/reflect/DoFnSignatures.java index 9f3491bca7b9..6d28f9343866 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/reflect/DoFnSignatures.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/reflect/DoFnSignatures.java @@ -2421,7 +2421,8 @@ private static String format(TypeDescriptor t) { } private static String format(Class kls) { - return kls.getSimpleName().isEmpty() ? kls.getName() : kls.getSimpleName(); + String simpleName = ReflectHelpers.getSimpleName(kls); + return simpleName.isEmpty() ? kls.getName() : simpleName; } static class ErrorReporter { diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/common/ReflectHelpers.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/common/ReflectHelpers.java index c2d945bbaac1..7d5964cb83ca 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/common/ReflectHelpers.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/common/ReflectHelpers.java @@ -52,11 +52,29 @@ public class ReflectHelpers { private static final Joiner COMMA_SEPARATOR = Joiner.on(", "); + /** + * Returns the simple name of the given class, or a fallback simple name derived from the class + * name if {@link Class#getSimpleName()} throws an exception (such as an {@link + * IllegalAccessError} on Java 17+ for non-public inner classes). + */ + public static String getSimpleName(Class clazz) { + try { + return clazz.getSimpleName(); + } catch (Throwable t) { + if (clazz.isArray()) { + return getSimpleName(clazz.getComponentType()) + "[]"; + } + String name = clazz.getName(); + int idx = Math.max(name.lastIndexOf('.'), name.lastIndexOf('$')); + return idx != -1 ? name.substring(idx + 1) : name; + } + } + /** Returns a string representation of the signature of a {@link Method}. */ public static String formatMethod(Method input) { String parameterTypes = FluentIterable.from(asList(input.getParameterTypes())) - .transform(Class::getSimpleName) + .transform(ReflectHelpers::getSimpleName) .join(COMMA_SEPARATOR); return String.format("%s(%s)", input.getName(), parameterTypes); } @@ -100,7 +118,7 @@ private static void format(StringBuilder builder, Type t) { } private static void formatClass(StringBuilder builder, Class clazz) { - builder.append(clazz.getSimpleName()); + builder.append(getSimpleName(clazz)); } private static void formatTypeVariable(StringBuilder builder, TypeVariable t) { From 2bea70be7750c174aa12f97a605c1faf5486a4f3 Mon Sep 17 00:00:00 2001 From: Abdelrahman Ibrahim Date: Wed, 15 Jul 2026 15:04:21 +0300 Subject: [PATCH 15/21] add jdkAddOpenModules for Dataflow V2 (#39306) --- .../trigger_files/beam_PostCommit_Java_DataflowV2.json | 2 +- sdks/java/container/common.gradle | 1 + sdks/java/container/java17/option-arrow.json | 9 +++++++++ sdks/java/container/java21/option-arrow.json | 9 +++++++++ sdks/java/container/java25/option-arrow.json | 9 +++++++++ 5 files changed, 29 insertions(+), 1 deletion(-) create mode 100644 sdks/java/container/java17/option-arrow.json create mode 100644 sdks/java/container/java21/option-arrow.json create mode 100644 sdks/java/container/java25/option-arrow.json diff --git a/.github/trigger_files/beam_PostCommit_Java_DataflowV2.json b/.github/trigger_files/beam_PostCommit_Java_DataflowV2.json index d0d502569834..c4a6954cecb7 100644 --- a/.github/trigger_files/beam_PostCommit_Java_DataflowV2.json +++ b/.github/trigger_files/beam_PostCommit_Java_DataflowV2.json @@ -1,4 +1,4 @@ { - "modification": 7, + "modification": 9, "https://github.com/apache/beam/pull/35159": "moving WindowedValue and making an interface" } diff --git a/sdks/java/container/common.gradle b/sdks/java/container/common.gradle index c81a33827bef..5b2322f61a35 100644 --- a/sdks/java/container/common.gradle +++ b/sdks/java/container/common.gradle @@ -83,6 +83,7 @@ task copyGolangLicenses(type: Copy) { task copyJdkOptions(type: Copy) { from "option-jamm.json" + from "option-arrow.json" from "java${imageJavaVersion}-security.properties" from "option-java${imageJavaVersion}-security.json" into "build/target/options" diff --git a/sdks/java/container/java17/option-arrow.json b/sdks/java/container/java17/option-arrow.json new file mode 100644 index 000000000000..a22b46ac1f78 --- /dev/null +++ b/sdks/java/container/java17/option-arrow.json @@ -0,0 +1,9 @@ +{ + "name": "arrow", + "enabled": true, + "options": { + "java_arguments": [ + "--add-opens=java.base/java.nio=ALL-UNNAMED" + ] + } +} diff --git a/sdks/java/container/java21/option-arrow.json b/sdks/java/container/java21/option-arrow.json new file mode 100644 index 000000000000..a22b46ac1f78 --- /dev/null +++ b/sdks/java/container/java21/option-arrow.json @@ -0,0 +1,9 @@ +{ + "name": "arrow", + "enabled": true, + "options": { + "java_arguments": [ + "--add-opens=java.base/java.nio=ALL-UNNAMED" + ] + } +} diff --git a/sdks/java/container/java25/option-arrow.json b/sdks/java/container/java25/option-arrow.json new file mode 100644 index 000000000000..a22b46ac1f78 --- /dev/null +++ b/sdks/java/container/java25/option-arrow.json @@ -0,0 +1,9 @@ +{ + "name": "arrow", + "enabled": true, + "options": { + "java_arguments": [ + "--add-opens=java.base/java.nio=ALL-UNNAMED" + ] + } +} From 1e6b61b3c8008154dcbf6cb0f21901318f913b82 Mon Sep 17 00:00:00 2001 From: Abdelrahman Ibrahim Date: Wed, 15 Jul 2026 15:27:09 +0300 Subject: [PATCH 16/21] Add Spark JVM --add-opens for (Nexmark, TPC-DS, PortableJar) (#39337) --- .../beam_PostCommit_Java_Nexmark_Spark.json | 4 ++++ .../beam_PostCommit_Java_Tpcds_Spark.json | 4 ++++ .../beam_PostCommit_PortableJar_Spark.json | 4 ++++ runners/portability/test_pipeline_jar.sh | 11 ++++++++++- sdks/java/testing/nexmark/build.gradle | 14 ++++++++++++++ sdks/java/testing/tpcds/build.gradle | 14 ++++++++++++++ 6 files changed, 50 insertions(+), 1 deletion(-) create mode 100644 .github/trigger_files/beam_PostCommit_Java_Nexmark_Spark.json create mode 100644 .github/trigger_files/beam_PostCommit_Java_Tpcds_Spark.json create mode 100644 .github/trigger_files/beam_PostCommit_PortableJar_Spark.json diff --git a/.github/trigger_files/beam_PostCommit_Java_Nexmark_Spark.json b/.github/trigger_files/beam_PostCommit_Java_Nexmark_Spark.json new file mode 100644 index 000000000000..e3d6056a5de9 --- /dev/null +++ b/.github/trigger_files/beam_PostCommit_Java_Nexmark_Spark.json @@ -0,0 +1,4 @@ +{ + "comment": "Modify this file in a trivial way to cause this test suite to run", + "modification": 1 +} diff --git a/.github/trigger_files/beam_PostCommit_Java_Tpcds_Spark.json b/.github/trigger_files/beam_PostCommit_Java_Tpcds_Spark.json new file mode 100644 index 000000000000..e3d6056a5de9 --- /dev/null +++ b/.github/trigger_files/beam_PostCommit_Java_Tpcds_Spark.json @@ -0,0 +1,4 @@ +{ + "comment": "Modify this file in a trivial way to cause this test suite to run", + "modification": 1 +} diff --git a/.github/trigger_files/beam_PostCommit_PortableJar_Spark.json b/.github/trigger_files/beam_PostCommit_PortableJar_Spark.json new file mode 100644 index 000000000000..e3d6056a5de9 --- /dev/null +++ b/.github/trigger_files/beam_PostCommit_PortableJar_Spark.json @@ -0,0 +1,4 @@ +{ + "comment": "Modify this file in a trivial way to cause this test suite to run", + "modification": 1 +} diff --git a/runners/portability/test_pipeline_jar.sh b/runners/portability/test_pipeline_jar.sh index f01d67b6580c..871a604e319f 100755 --- a/runners/portability/test_pipeline_jar.sh +++ b/runners/portability/test_pipeline_jar.sh @@ -123,7 +123,16 @@ OUTPUT_JAR="test-pipeline-${RUNNER}-$(date +%Y%m%d-%H%M%S).jar" if [[ "$TEST_EXIT_CODE" -eq 0 ]]; then # Execute the jar - java -jar $OUTPUT_JAR || TEST_EXIT_CODE=$? + JAVA_ARGS=() + if [[ "$RUNNER" = "SparkRunner" ]]; then + JAVA_ARGS+=( + "--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" + ) + fi + java "${JAVA_ARGS[@]}" -jar $OUTPUT_JAR || TEST_EXIT_CODE=$? fi rm -rf $ENV_DIR diff --git a/sdks/java/testing/nexmark/build.gradle b/sdks/java/testing/nexmark/build.gradle index bd917a3935ad..b554e9d9297f 100644 --- a/sdks/java/testing/nexmark/build.gradle +++ b/sdks/java/testing/nexmark/build.gradle @@ -114,6 +114,19 @@ if (isSparkRunner) { } } +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() @@ -179,6 +192,7 @@ task run(type: JavaExec) { systemProperty "spark.ui.showConsoleProgress", "false" // Dataset runner only systemProperty "spark.sql.shuffle.partitions", "4" + jvmArgs += sparkJvmArgs() } mainClass = "org.apache.beam.sdk.nexmark.Main" diff --git a/sdks/java/testing/tpcds/build.gradle b/sdks/java/testing/tpcds/build.gradle index 6a05df96ab20..15fdd480f076 100644 --- a/sdks/java/testing/tpcds/build.gradle +++ b/sdks/java/testing/tpcds/build.gradle @@ -104,6 +104,19 @@ if (isSpark) { } } +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: @@ -141,6 +154,7 @@ task run(type: JavaExec) { // Dataset runner only systemProperty "spark.sql.shuffle.partitions", "4" systemProperty "spark.sql.adaptive.enabled", "false" // high overhead for complex queries + jvmArgs += sparkJvmArgs() } mainClass = "org.apache.beam.sdk.tpcds.BeamTpcds" From 163cd12937bcc2b707b27325f3c8d55d11b28900 Mon Sep 17 00:00:00 2001 From: Abdelrahman Ibrahim Date: Wed, 15 Jul 2026 15:28:53 +0300 Subject: [PATCH 17/21] add Spark job-server --add-opens (#39339) --- .../trigger_files/beam_PostCommit_Go_VR_Spark.json | 4 ++-- .../beam_PostCommit_Python_Examples_Spark.json | 4 ++++ .../beam_PostCommit_Python_ValidatesRunner_Spark.json | 3 ++- sdks/go/test/run_validatesrunner_tests.sh | 4 ++++ .../apache_beam/runners/portability/job_server.py | 4 ++-- .../runners/portability/job_server_test.py | 2 +- .../apache_beam/runners/portability/spark_runner.py | 11 +++++++++++ .../runners/portability/spark_runner_test.py | 2 ++ 8 files changed, 28 insertions(+), 6 deletions(-) create mode 100644 .github/trigger_files/beam_PostCommit_Python_Examples_Spark.json diff --git a/.github/trigger_files/beam_PostCommit_Go_VR_Spark.json b/.github/trigger_files/beam_PostCommit_Go_VR_Spark.json index 72b690e649d3..0dcf05b4215b 100644 --- a/.github/trigger_files/beam_PostCommit_Go_VR_Spark.json +++ b/.github/trigger_files/beam_PostCommit_Go_VR_Spark.json @@ -1,5 +1,5 @@ { "comment": "Modify this file in a trivial way to cause this test suite to run", - "modification": 1, - "https://github.com/apache/beam/pull/36527": "skip a processing time timer test in spark", + "modification": 2, + "https://github.com/apache/beam/pull/36527": "skip a processing time timer test in spark" } diff --git a/.github/trigger_files/beam_PostCommit_Python_Examples_Spark.json b/.github/trigger_files/beam_PostCommit_Python_Examples_Spark.json new file mode 100644 index 000000000000..e3d6056a5de9 --- /dev/null +++ b/.github/trigger_files/beam_PostCommit_Python_Examples_Spark.json @@ -0,0 +1,4 @@ +{ + "comment": "Modify this file in a trivial way to cause this test suite to run", + "modification": 1 +} diff --git a/.github/trigger_files/beam_PostCommit_Python_ValidatesRunner_Spark.json b/.github/trigger_files/beam_PostCommit_Python_ValidatesRunner_Spark.json index f4288649d8ab..f4ec72dc416b 100644 --- a/.github/trigger_files/beam_PostCommit_Python_ValidatesRunner_Spark.json +++ b/.github/trigger_files/beam_PostCommit_Python_ValidatesRunner_Spark.json @@ -2,5 +2,6 @@ "https://github.com/apache/beam/pull/34830": "testing", "https://github.com/apache/beam/issues/35429": "testing", "trigger-2026-04-04": "portable_runner expand_sdf opt-in", - "https://github.com/apache/beam/pull/38892": "UnboundedSource portable VR test" + "https://github.com/apache/beam/pull/38892": "UnboundedSource portable VR test", + "modification": 1 } diff --git a/sdks/go/test/run_validatesrunner_tests.sh b/sdks/go/test/run_validatesrunner_tests.sh index 9ea89fd22cfc..f1d681267691 100755 --- a/sdks/go/test/run_validatesrunner_tests.sh +++ b/sdks/go/test/run_validatesrunner_tests.sh @@ -286,6 +286,10 @@ if [[ "$RUNNER" == "flink" || "$RUNNER" == "spark" || "$RUNNER" == "portable" || --artifact-port 0 & elif [[ "$RUNNER" == "spark" ]]; then "$JAVA_CMD" \ + --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 \ -jar $SPARK_JOB_SERVER_JAR \ --spark-master-url local \ --job-port $JOB_PORT \ diff --git a/sdks/python/apache_beam/runners/portability/job_server.py b/sdks/python/apache_beam/runners/portability/job_server.py index 8bdad6b70fea..6c55f642c615 100644 --- a/sdks/python/apache_beam/runners/portability/job_server.py +++ b/sdks/python/apache_beam/runners/portability/job_server.py @@ -175,8 +175,8 @@ def subprocess_cmd_and_endpoint(self): self._artifacts_dir if self._artifacts_dir else self.local_temp_dir( prefix='artifacts')) job_port, = subprocess_server.pick_port(self._job_port) - subprocess_cmd = [self._java_launcher, '-jar'] + self._jvm_properties + [ - jar_path + subprocess_cmd = [self._java_launcher] + list(self._jvm_properties) + [ + '-jar', jar_path ] + list( self.java_arguments( job_port, self._artifact_port, self._expansion_port, artifacts_dir)) diff --git a/sdks/python/apache_beam/runners/portability/job_server_test.py b/sdks/python/apache_beam/runners/portability/job_server_test.py index 13b3629b24bf..fbf668557794 100644 --- a/sdks/python/apache_beam/runners/portability/job_server_test.py +++ b/sdks/python/apache_beam/runners/portability/job_server_test.py @@ -62,8 +62,8 @@ def test_subprocess_cmd_and_endpoint(self): subprocess_cmd, [ '/path/to/java', - '-jar', '-Dsome.property=value', + '-jar', '/path/to/jar', '--artifacts-dir', '/path/to/artifacts/', diff --git a/sdks/python/apache_beam/runners/portability/spark_runner.py b/sdks/python/apache_beam/runners/portability/spark_runner.py index 480fbdecdce3..f7eeb6509f6a 100644 --- a/sdks/python/apache_beam/runners/portability/spark_runner.py +++ b/sdks/python/apache_beam/runners/portability/spark_runner.py @@ -31,6 +31,13 @@ # https://spark.apache.org/docs/latest/submitting-applications.html#master-urls LOCAL_MASTER_PATTERN = r'^local(\[.+\])?$' +SPARK_JAR_JOB_SERVER_JVM_ARGS = [ + '--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', +] + # Since Java job servers are heavyweight external processes, cache them. # This applies only to SparkJarJobServer, not SparkUberJarJobServer. JOB_SERVER_CACHE = {} @@ -84,6 +91,10 @@ def __init__(self, options): self._jar = options.spark_job_server_jar self._master_url = options.spark_master_url self._spark_version = options.spark_version + self._jvm_properties = list(self._jvm_properties) + for arg in SPARK_JAR_JOB_SERVER_JVM_ARGS: + if arg not in self._jvm_properties: + self._jvm_properties.append(arg) def path_to_jar(self): if self._jar: diff --git a/sdks/python/apache_beam/runners/portability/spark_runner_test.py b/sdks/python/apache_beam/runners/portability/spark_runner_test.py index 54238bbf1928..4152b8d09f4f 100644 --- a/sdks/python/apache_beam/runners/portability/spark_runner_test.py +++ b/sdks/python/apache_beam/runners/portability/spark_runner_test.py @@ -29,6 +29,7 @@ from apache_beam.runners.portability import job_server from apache_beam.runners.portability import portable_runner from apache_beam.runners.portability import portable_runner_test +from apache_beam.runners.portability.spark_runner import SPARK_JAR_JOB_SERVER_JVM_ARGS from apache_beam.utils import subprocess_server # Run as @@ -100,6 +101,7 @@ def _subprocess_command(cls, job_port, expansion_port): try: return [ subprocess_server.JavaHelper.get_java(), + *SPARK_JAR_JOB_SERVER_JVM_ARGS, '-Dbeam.spark.test.reuseSparkContext=true', '-jar', cls.spark_job_server_jar, From dca9622face7f121ad53e187f63082b1c13dbd3f Mon Sep 17 00:00:00 2001 From: raman118 Date: Thu, 2 Jul 2026 04:52:25 +0530 Subject: [PATCH 18/21] fix: add retries and query parameter encoding for GitHub API requests (closes #39188) --- infra/enforcement/account_keys.py | 7 +-- infra/enforcement/sending.py | 66 +++++++++++++++++++---- infra/enforcement/test_sending.py | 88 +++++++++++++++++++++++++++++++ 3 files changed, 148 insertions(+), 13 deletions(-) create mode 100644 infra/enforcement/test_sending.py diff --git a/infra/enforcement/account_keys.py b/infra/enforcement/account_keys.py index 31c1354319d6..56ccbf654b03 100644 --- a/infra/enforcement/account_keys.py +++ b/infra/enforcement/account_keys.py @@ -164,14 +164,15 @@ def _get_all_live_service_accounts(self) -> List[str]: request.name = f"projects/{self.project_id}" try: - accounts = self.service_account_client.list_service_accounts(request=request) - self.logger.debug(f"Retrieved {len(accounts.accounts)} service accounts for project {self.project_id}") + accounts_pager = self.service_account_client.list_service_accounts(request=request) + accounts = list(accounts_pager) + self.logger.debug(f"Retrieved {len(accounts)} service accounts for project {self.project_id}") if not accounts: self.logger.warning(f"No service accounts found in project {self.project_id}.") return [] - return [self._normalize_account_email(account.email) for account in accounts.accounts if not account.disabled] + return [self._normalize_account_email(account.email) for account in accounts if not account.disabled] except Exception as e: self.logger.error(f"Failed to retrieve service accounts for project {self.project_id}: {e}") raise diff --git a/infra/enforcement/sending.py b/infra/enforcement/sending.py index 37de025a207f..bd9787b6ce87 100644 --- a/infra/enforcement/sending.py +++ b/infra/enforcement/sending.py @@ -59,25 +59,71 @@ def __init__(self, logger: logging.Logger, github_token: str, github_repo: str, self.logger = logger self.github_api_url = "https://api.github.com" - def _make_github_request(self, method: str, endpoint: str, json: Optional[dict] = None) -> requests.Response: + def _make_github_request(self, method: str, endpoint: str, json: Optional[dict] = None, params: Optional[dict] = None) -> requests.Response: """ - Makes a request to the GitHub API. + Makes a request to the GitHub API with retry logic for transient errors and rate limiting. Args: method (str): The HTTP method to use (e.g., "GET", "POST", "PATCH"). endpoint (str): The API endpoint to call. json (Optional[dict]): The JSON payload to send with the request. + params (Optional[dict]): The URL parameters to send with the request. Returns: requests.Response: The response from the API. """ + import time url = f"{self.github_api_url}/{endpoint}" - response = requests.request(method, url, headers=self.headers, json=json) + max_retries = 5 + backoff = 2 - if not response.ok: - self.logger.error(f"Failed GitHub API request to {endpoint}: {response.status_code} - {response.text}") - response.raise_for_status() - + for attempt in range(max_retries): + try: + response = requests.request(method, url, headers=self.headers, json=json, params=params) + + # Check for rate limiting / secondary rate limit (403, 429) + if response.status_code in [403, 429]: + retry_after = response.headers.get("Retry-After") + if retry_after: + sleep_seconds = int(retry_after) + else: + sleep_seconds = backoff + + self.logger.warning( + f"GitHub API rate limit hit ({response.status_code}) on {endpoint}. " + f"Retrying in {sleep_seconds} seconds... (Attempt {attempt + 1}/{max_retries})" + ) + time.sleep(sleep_seconds) + backoff *= 2 + continue + + # Check for transient server errors (500, 502, 503, 504) + if response.status_code in [500, 502, 503, 504]: + self.logger.warning( + f"GitHub API server error ({response.status_code}) on {endpoint}. " + f"Retrying in {backoff} seconds... (Attempt {attempt + 1}/{max_retries})" + ) + time.sleep(backoff) + backoff *= 2 + continue + + if not response.ok: + self.logger.error(f"Failed GitHub API request to {endpoint}: {response.status_code} - {response.text}") + response.raise_for_status() + + return response + except requests.exceptions.RequestException as e: + if attempt == max_retries - 1: + raise + self.logger.warning( + f"GitHub API request exception on {endpoint}: {e}. " + f"Retrying in {backoff} seconds... (Attempt {attempt + 1}/{max_retries})" + ) + time.sleep(backoff) + backoff *= 2 + + self.logger.error(f"Failed GitHub API request to {endpoint} after {max_retries} attempts.") + response.raise_for_status() return response def _send_email(self, title: str, body: str, recipient: str) -> None: @@ -97,13 +143,13 @@ def _send_email(self, title: str, body: str, recipient: str) -> None: def _get_open_issues(self, title: str) -> List[GitHubIssue]: """ - Retrieves the number of open GitHub issues with a given title. + Retrieves the open GitHub issues with a given title. Args: title (str): The title of the GitHub issue. """ - endpoint = f"search/issues?q=is:issue+repo:{self.github_repo}+in:title+{title}+is:open" - response = self._make_github_request("GET", endpoint) + q = f'is:issue repo:{self.github_repo} in:title "{title}" is:open' + response = self._make_github_request("GET", "search/issues", params={"q": q}) issues = response.json().get('items', []) parsed_issues = [] for issue in issues: diff --git a/infra/enforcement/test_sending.py b/infra/enforcement/test_sending.py new file mode 100644 index 000000000000..70103b0fbcb7 --- /dev/null +++ b/infra/enforcement/test_sending.py @@ -0,0 +1,88 @@ +import unittest +from unittest.mock import patch, MagicMock +import logging +import requests +from sending import SendingClient + +class TestSendingClient(unittest.TestCase): + def setUp(self): + self.logger = logging.getLogger("TestLogger") + self.logger.setLevel(logging.DEBUG) + # Create a SendingClient with dummy values + self.client = SendingClient( + logger=self.logger, + github_token="dummy-token", + github_repo="apache/beam", + smtp_server="smtp.example.com", + smtp_port=587, + email="sender@example.com", + password="password" + ) + + @patch("requests.request") + def test_get_open_issues_flaky_retry(self, mock_request): + # We simulate a 403 secondary rate limit error on the first request, + # and a successful 200 response on the second request. + mock_response_fail = MagicMock() + mock_response_fail.status_code = 403 + mock_response_fail.ok = False + mock_response_fail.headers = {"Retry-After": "1"} + mock_response_fail.raise_for_status.side_effect = requests.exceptions.HTTPError("403 Client Error") + + mock_response_success = MagicMock() + mock_response_success.status_code = 200 + mock_response_success.ok = True + mock_response_success.json.return_value = { + "items": [ + { + "number": 1234, + "title": "[SECURITY] Action Required: Unmanaged Service Account Keys Detected", + "body": "Test body", + "state": "open", + "html_url": "https://github.com/apache/beam/issues/1234", + "created_at": "2026-07-02T00:00:00Z", + "updated_at": "2026-07-02T00:00:00Z" + } + ] + } + + mock_request.side_effect = [mock_response_fail, mock_response_success] + + # Call get_open_issues + issues = self.client._get_open_issues("[SECURITY] Action Required: Unmanaged Service Account Keys Detected") + + # Verify that two requests were made (one retry) + self.assertEqual(mock_request.call_count, 2) + self.assertEqual(len(issues), 1) + self.assertEqual(issues[0].number, 1234) + + @patch("requests.request") + def test_get_open_issues_query_format(self, mock_request): + mock_response = MagicMock() + mock_response.status_code = 200 + mock_response.ok = True + mock_response.json.return_value = {"items": []} + mock_request.return_value = mock_response + + title = "[SECURITY] Action Required: Unmanaged Service Account Keys Detected" + self.client._get_open_issues(title) + + # Verify that the query parameter was passed correctly to requests + mock_request.assert_called_once() + args, kwargs = mock_request.call_args + + # Check that we passed params dict containing 'q' + self.assertIn("params", kwargs) + self.assertIn("q", kwargs["params"]) + + q = kwargs["params"]["q"] + # The query should specify the repo, title (quoted), state and type + self.assertIn('repo:apache/beam', q) + self.assertIn('is:issue', q) + self.assertIn('is:open', q) + self.assertIn('in:title', q) + self.assertIn(f'"{title}"', q) + + +if __name__ == "__main__": + unittest.main() From 7b88d16f58b8213384c372fb3b892fecf26f2492 Mon Sep 17 00:00:00 2001 From: raman118 Date: Thu, 2 Jul 2026 05:03:20 +0530 Subject: [PATCH 19/21] 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 5939496fd151aaeeb8386cde1e512f9746f7dc68 Mon Sep 17 00:00:00 2001 From: Ryan Wigglesworth Date: Fri, 17 Jul 2026 19:06:00 +0000 Subject: [PATCH 20/21] Move validation to StreamingModeExecutionContext --- .../worker/StreamingModeExecutionContext.java | 58 +++++++++++++++ .../streaming/KeyCommitTooLargeException.java | 16 ++++- .../ComputationWorkExecutorFactory.java | 14 +++- .../processing/StreamingWorkScheduler.java | 70 ++----------------- .../StreamingModeExecutionContextTest.java | 5 ++ .../worker/WorkerCustomSourcesTest.java | 8 +++ 6 files changed, 103 insertions(+), 68 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 9401cc5f8ed9..41119742b15b 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,10 +52,12 @@ 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.KeyCommitTooLargeException; import org.apache.beam.runners.dataflow.worker.streaming.Watermarks; import org.apache.beam.runners.dataflow.worker.streaming.Work; import org.apache.beam.runners.dataflow.worker.streaming.config.StreamingGlobalConfig; import org.apache.beam.runners.dataflow.worker.streaming.config.StreamingGlobalConfigHandle; +import org.apache.beam.runners.dataflow.worker.streaming.harness.StreamingCounters; import org.apache.beam.runners.dataflow.worker.streaming.sideinput.SideInput; import org.apache.beam.runners.dataflow.worker.streaming.sideinput.SideInputState; import org.apache.beam.runners.dataflow.worker.streaming.sideinput.SideInputStateFetcher; @@ -77,6 +79,7 @@ import org.apache.beam.runners.dataflow.worker.windmill.state.WindmillTagEncodingV1; import org.apache.beam.runners.dataflow.worker.windmill.state.WindmillTagEncodingV2; import org.apache.beam.runners.dataflow.worker.windmill.state.WindmillTimerData; +import org.apache.beam.runners.dataflow.worker.windmill.work.processing.failures.FailureTracker; import org.apache.beam.sdk.annotations.Internal; import org.apache.beam.sdk.coders.Coder; import org.apache.beam.sdk.coders.CoderException; @@ -177,6 +180,9 @@ public class StreamingModeExecutionContext private final HotKeyLogger hotKeyLogger; private final boolean hotKeyLoggingEnabled; private final String stepName; + private final String systemName; + private final StreamingCounters streamingCounters; + private final FailureTracker failureTracker; private @Nullable Coder keyCoder; // Key switch listener to delegate MDC logging context and thread name updates @@ -212,6 +218,9 @@ public StreamingModeExecutionContext( HotKeyLogger hotKeyLogger, boolean hotKeyLoggingEnabled, String stepName, + String systemName, + StreamingCounters streamingCounters, + FailureTracker failureTracker, String sourceBytesProcessCounterName, SideInputStateFetcherFactory sideInputStateFetcherFactory) { super( @@ -231,6 +240,9 @@ public StreamingModeExecutionContext( this.hotKeyLogger = checkNotNull(hotKeyLogger); this.hotKeyLoggingEnabled = hotKeyLoggingEnabled; this.stepName = checkNotNull(stepName); + this.systemName = checkNotNull(systemName); + this.streamingCounters = checkNotNull(streamingCounters); + this.failureTracker = checkNotNull(failureTracker); this.sourceBytesProcessCounterName = checkNotNull(sourceBytesProcessCounterName); this.sideInputStateFetcherFactory = sideInputStateFetcherFactory; StreamingGlobalConfig config = globalConfigHandle.getConfig(); @@ -669,6 +681,52 @@ private void flushStateInternal() { getOutputBuilder() .setSourceBytesProcessed(computeSourceBytesProcessed(sourceBytesProcessCounterName)); + + validateCommitRequestSize(); + } + + private void validateCommitRequestSize() { + Windmill.WorkItemCommitRequest.Builder currentBuilder = getOutputBuilder(); + Work currentWork = getWork(); + 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. + streamingCounters.windmillMaxObservedWorkItemCommitBytes().addValue(estimatedCommitSize); + if (commitSize >= 0 && commitSize < byteLimit) { + return; + } + + KeyCommitTooLargeException e = + KeyCommitTooLargeException.causedBy( + systemName, byteLimit, commitRequest, key, hotKeyLoggingEnabled); + failureTracker.trackFailure(systemName, currentWork.getWorkItem(), 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 + 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; + } + + private Windmill.WorkItemCommitRequest.Builder buildWorkItemTruncationRequestBuilder( + Work work, int estimatedCommitSize) { + Windmill.WorkItemCommitRequest.Builder outputBuilder = createOutputBuilder(work); + outputBuilder.setExceedsMaxWorkItemCommitBytes(true); + outputBuilder.setEstimatedWorkItemCommitBytes(estimatedCommitSize); + return outputBuilder; } private final long computeSourceBytesProcessed(String sourceBytesCounterName) { 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 4a52d9fde771..0c3591102f9c 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 @@ -49,11 +49,13 @@ import org.apache.beam.runners.dataflow.worker.streaming.ComputationWorkExecutor; import org.apache.beam.runners.dataflow.worker.streaming.StageInfo; import org.apache.beam.runners.dataflow.worker.streaming.config.StreamingGlobalConfigHandle; +import org.apache.beam.runners.dataflow.worker.streaming.harness.StreamingCounters; import org.apache.beam.runners.dataflow.worker.streaming.sideinput.SideInputStateFetcherFactory; import org.apache.beam.runners.dataflow.worker.util.common.worker.MapTaskExecutor; import org.apache.beam.runners.dataflow.worker.util.common.worker.OutputObjectAndByteCounter; import org.apache.beam.runners.dataflow.worker.util.common.worker.ReadOperation; import org.apache.beam.runners.dataflow.worker.windmill.state.WindmillStateCache; +import org.apache.beam.runners.dataflow.worker.windmill.work.processing.failures.FailureTracker; import org.apache.beam.sdk.coders.Coder; import org.apache.beam.sdk.coders.KvCoder; import org.apache.beam.sdk.fn.IdGenerator; @@ -84,6 +86,8 @@ final class ComputationWorkExecutorFactory { private final SinkRegistry sinkRegistry; private final DataflowExecutionStateSampler sampler; private final CounterSet pendingDeltaCounters; + private final StreamingCounters streamingCounters; + private final FailureTracker failureTracker; /** * Function which converts map tasks to their network representation for execution. @@ -108,7 +112,8 @@ final class ComputationWorkExecutorFactory { ReaderCache readerCache, Function stateCacheFactory, DataflowExecutionStateSampler sampler, - CounterSet pendingDeltaCounters, + StreamingCounters streamingCounters, + FailureTracker failureTracker, IdGenerator idGenerator, StreamingGlobalConfigHandle globalConfigHandle, HotKeyLogger hotKeyLogger, @@ -122,7 +127,9 @@ final class ComputationWorkExecutorFactory { this.readerRegistry = ReaderRegistry.defaultRegistry(); this.sinkRegistry = SinkRegistry.defaultRegistry(); this.sampler = sampler; - this.pendingDeltaCounters = pendingDeltaCounters; + this.streamingCounters = streamingCounters; + this.failureTracker = failureTracker; + this.pendingDeltaCounters = streamingCounters.pendingDeltaCounters(); this.mapTaskToNetwork = new MapTaskToNetworkFunction(idGenerator); this.maxSinkBytes = hasExperiment(options, DISABLE_SINK_BYTE_LIMIT_EXPERIMENT) @@ -286,6 +293,9 @@ private StreamingModeExecutionContext createExecutionContext( hotKeyLogger, hotKeyLoggingEnabled, stepName, + stageInfo.systemName(), + streamingCounters, + failureTracker, computationState.sourceBytesProcessCounterName(), sideInputStateFetcherFactory); } 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 45b9fbe93396..990393f51f8f 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 @@ -43,7 +43,6 @@ import org.apache.beam.runners.dataflow.worker.streaming.ComputationState; import org.apache.beam.runners.dataflow.worker.streaming.ComputationWorkExecutor; 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.StageInfo; import org.apache.beam.runners.dataflow.worker.streaming.Watermarks; import org.apache.beam.runners.dataflow.worker.streaming.Work; @@ -81,39 +80,30 @@ public class StreamingWorkScheduler { private final Supplier clock; private final ComputationWorkExecutorFactory computationWorkExecutorFactory; - private final FailureTracker failureTracker; private final WorkFailureProcessor workFailureProcessor; private final StreamingCommitFinalizer commitFinalizer; private final StreamingCounters streamingCounters; private final ConcurrentMap stageInfoMap; private final DataflowExecutionStateSampler sampler; - private final StreamingGlobalConfigHandle globalConfigHandle; private final BoundedQueueExecutor workExecutor; - private final boolean hotKeyLoggingEnabled; public StreamingWorkScheduler( Supplier clock, BoundedQueueExecutor workExecutor, ComputationWorkExecutorFactory computationWorkExecutorFactory, - FailureTracker failureTracker, WorkFailureProcessor workFailureProcessor, StreamingCommitFinalizer commitFinalizer, StreamingCounters streamingCounters, ConcurrentMap stageInfoMap, - DataflowExecutionStateSampler sampler, - StreamingGlobalConfigHandle globalConfigHandle, - boolean hotKeyLoggingEnabled) { + DataflowExecutionStateSampler sampler) { this.clock = clock; this.workExecutor = workExecutor; this.computationWorkExecutorFactory = computationWorkExecutorFactory; - this.failureTracker = failureTracker; this.workFailureProcessor = workFailureProcessor; this.commitFinalizer = commitFinalizer; this.streamingCounters = streamingCounters; this.stageInfoMap = stageInfoMap; this.sampler = sampler; - this.globalConfigHandle = globalConfigHandle; - this.hotKeyLoggingEnabled = hotKeyLoggingEnabled; } public static StreamingWorkScheduler create( @@ -142,30 +132,22 @@ public static StreamingWorkScheduler create( readerCache, stateCacheFactory, sampler, - streamingCounters.pendingDeltaCounters(), + streamingCounters, + failureTracker, idGenerator, globalConfigHandle, hotKeyLogger, sideInputStateFetcherFactory); - List experiments = options.getExperiments(); - boolean hotKeyLoggingEnabled = - options.isHotKeyLoggingEnabled() - || (experiments != null - && experiments.stream().anyMatch("enable_hot_key_logging"::equalsIgnoreCase)); - return new StreamingWorkScheduler( clock, workExecutor, computationWorkExecutorFactory, - failureTracker, workFailureProcessor, StreamingCommitFinalizer.create(workExecutor, commitFinalizerCleanupExecutor), streamingCounters, stageInfoMap, - sampler, - globalConfigHandle, - hotKeyLoggingEnabled); + sampler); } private static long computeShuffleBytesRead(Windmill.WorkItem workItem) { @@ -185,14 +167,6 @@ private static Windmill.WorkItemCommitRequest.Builder initializeOutputBuilder( .setCacheToken(workItem.getCacheToken()); } - private static Windmill.WorkItemCommitRequest buildWorkItemTruncationRequest( - ByteString key, Windmill.WorkItem workItem, int estimatedCommitSize) { - Windmill.WorkItemCommitRequest.Builder outputBuilder = initializeOutputBuilder(key, workItem); - outputBuilder.setExceedsMaxWorkItemCommitBytes(true); - outputBuilder.setEstimatedWorkItemCommitBytes(estimatedCommitSize); - return outputBuilder.build(); - } - /** 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); @@ -315,32 +289,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, @@ -457,17 +405,13 @@ private void commitWorkBatch( 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); - validatedCommitRequest = - validatedCommitRequest + Windmill.WorkItemCommitRequest commitRequestWithAttributions = + commitRequest .toBuilder() .addAllPerWorkItemLatencyAttributions(work.getLatencyAttributions(sampler)) .build(); - work.queueCommit(validatedCommitRequest, computationState); + work.queueCommit(commitRequestWithAttributions, computationState); } private void recordProcessingTime( 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 561596f68d0f..3d11e5097463 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 @@ -63,6 +63,7 @@ import org.apache.beam.runners.dataflow.worker.streaming.config.FakeGlobalConfigHandle; import org.apache.beam.runners.dataflow.worker.streaming.config.StreamingGlobalConfig; import org.apache.beam.runners.dataflow.worker.streaming.config.StreamingGlobalConfigHandle; +import org.apache.beam.runners.dataflow.worker.streaming.harness.StreamingCounters; import org.apache.beam.runners.dataflow.worker.streaming.sideinput.SideInputStateFetcherFactory; import org.apache.beam.runners.dataflow.worker.util.common.worker.WorkExecutor; import org.apache.beam.runners.dataflow.worker.windmill.Windmill; @@ -71,6 +72,7 @@ import org.apache.beam.runners.dataflow.worker.windmill.state.WindmillStateReader; import org.apache.beam.runners.dataflow.worker.windmill.state.WindmillTagEncodingV1; import org.apache.beam.runners.dataflow.worker.windmill.state.WindmillTagEncodingV2; +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; import org.apache.beam.sdk.coders.Coder; @@ -142,6 +144,9 @@ private StreamingModeExecutionContext createExecutionContext( new HotKeyLogger(), /*hotKeyLoggingEnabled=*/ false, /*stepName=*/ "stepName", + /*systemName=*/ "systemName", + StreamingCounters.create(), + mock(FailureTracker.class), "sourceBytesProcessCounterName", SideInputStateFetcherFactory.fromOptions(options)); } 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 31ea1bab07af..2da52bea8eba 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 @@ -94,6 +94,7 @@ import org.apache.beam.runners.dataflow.worker.streaming.config.FixedGlobalConfigHandle; import org.apache.beam.runners.dataflow.worker.streaming.config.StreamingGlobalConfig; import org.apache.beam.runners.dataflow.worker.streaming.config.StreamingGlobalConfigHandle; +import org.apache.beam.runners.dataflow.worker.streaming.harness.StreamingCounters; import org.apache.beam.runners.dataflow.worker.streaming.sideinput.SideInputStateFetcherFactory; import org.apache.beam.runners.dataflow.worker.testing.TestCountingSource; import org.apache.beam.runners.dataflow.worker.util.common.worker.NativeReader; @@ -103,6 +104,7 @@ 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; import org.apache.beam.sdk.coders.BigEndianIntegerCoder; @@ -641,6 +643,9 @@ public void testReadUnboundedReader() throws Exception { new HotKeyLogger(), /*hotKeyLoggingEnabled=*/ false, /*stepName=*/ "stepName", + /*systemName=*/ "systemName", + StreamingCounters.create(), + mock(FailureTracker.class), "sourceBytesProcessCounterName", SideInputStateFetcherFactory.fromOptions(options)); @@ -1014,6 +1019,9 @@ public void testFailedWorkItemsAbort() throws Exception { new HotKeyLogger(), /*hotKeyLoggingEnabled=*/ false, /*stepName=*/ "stepName", + /*systemName=*/ "systemName", + StreamingCounters.create(), + mock(FailureTracker.class), "sourceBytesProcessCounterName", SideInputStateFetcherFactory.fromOptions(options)); From a3b7163a8111a5dd533f4846c4bf1a64f2235f0a Mon Sep 17 00:00:00 2001 From: Ryan Wigglesworth Date: Fri, 17 Jul 2026 20:05:03 +0000 Subject: [PATCH 21/21] 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 {