From 24e9ed4e0e9541cb92aad22db4ca3a339b83fd04 Mon Sep 17 00:00:00 2001 From: "dependabot[bot]" <49699333+dependabot[bot]@users.noreply.github.com> Date: Wed, 24 Jun 2026 07:58:07 -0400 Subject: [PATCH 1/6] Bump actions/cache from 4 to 6 (#39081) Bumps [actions/cache](https://github.com/actions/cache) from 4 to 6. - [Release notes](https://github.com/actions/cache/releases) - [Changelog](https://github.com/actions/cache/blob/main/RELEASES.md) - [Commits](https://github.com/actions/cache/compare/v4...v6) --- updated-dependencies: - dependency-name: actions/cache dependency-version: '6' 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/playground_frontend_test.yml | 2 +- .github/workflows/tour_of_beam_frontend_test.yml | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/.github/workflows/playground_frontend_test.yml b/.github/workflows/playground_frontend_test.yml index c4c762704869..b54fdc0ecfd6 100644 --- a/.github/workflows/playground_frontend_test.yml +++ b/.github/workflows/playground_frontend_test.yml @@ -48,7 +48,7 @@ jobs: - uses: actions/checkout@v7 - name: 'Cache Flutter Dependencies' - uses: actions/cache@v4 + uses: actions/cache@v6 with: path: /opt/hostedtoolcache/flutter key: ${{ runner.OS }}-flutter-install-cache-${{ env.FLUTTER_VERSION }} diff --git a/.github/workflows/tour_of_beam_frontend_test.yml b/.github/workflows/tour_of_beam_frontend_test.yml index 0812740fabdc..c6df5bf6ea17 100644 --- a/.github/workflows/tour_of_beam_frontend_test.yml +++ b/.github/workflows/tour_of_beam_frontend_test.yml @@ -50,7 +50,7 @@ jobs: - uses: actions/checkout@v7 - name: 'Cache Flutter Dependencies' - uses: actions/cache@v4 + uses: actions/cache@v6 with: path: /opt/hostedtoolcache/flutter key: ${{ runner.OS }}-flutter-install-cache-${{ env.FLUTTER_VERSION }} From f7f0e1655d54e0f23e06614aed10d751488fd409 Mon Sep 17 00:00:00 2001 From: "dependabot[bot]" <49699333+dependabot[bot]@users.noreply.github.com> Date: Wed, 24 Jun 2026 07:58:35 -0400 Subject: [PATCH 2/6] Bump actions/checkout from 6 to 7 (#39080) 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> --- .../workflows/beam_PostCommit_Python_Xlang_Messaging_Direct.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.github/workflows/beam_PostCommit_Python_Xlang_Messaging_Direct.yml b/.github/workflows/beam_PostCommit_Python_Xlang_Messaging_Direct.yml index 943dfbaffd94..ef4c35c189ee 100644 --- a/.github/workflows/beam_PostCommit_Python_Xlang_Messaging_Direct.yml +++ b/.github/workflows/beam_PostCommit_Python_Xlang_Messaging_Direct.yml @@ -63,7 +63,7 @@ jobs: job_name: ["beam_PostCommit_Python_Xlang_Messaging_Direct"] job_phrase: ["Run Python_Xlang_Messaging_Direct PostCommit"] steps: - - uses: actions/checkout@v6 + - uses: actions/checkout@v7 - name: Setup repository uses: ./.github/actions/setup-action with: From 0165344ab138084d2e2998bbc7129decf0350f0f Mon Sep 17 00:00:00 2001 From: Arun Pandian Date: Wed, 24 Jun 2026 05:21:44 -0700 Subject: [PATCH 3/6] Add UnboundedCountingSource::to for bounded reads from an UnboundedCountingSource (#39084) * Add UnboundedCountingSource::to for bounded reads from an UnboundedCountingSource UnboundedCountingSource currently does not support configuring an end limit, this PR adds that capability. GenerateSequence with a from + to + rate boils down to a BoundedReadFromUnboundedSource. BoundedReadFromUnboundedSource splits emit all outputs from a single bundle that does not work well on Streaming jobs. UnboundedCountingSource::to is currently not exposed outside CountingSource. Eventually GenerateSequence can use UnboundedCountingSource::to instead of BoundedReadFromUnboundedSource --- .../apache/beam/sdk/io/CountingSource.java | 40 +++++++++++++++---- .../beam/sdk/io/CountingSourceTest.java | 36 +++++++++++++++++ 2 files changed, 68 insertions(+), 8 deletions(-) diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/io/CountingSource.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/io/CountingSource.java index 9d30efb2f113..4ef928480d75 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/io/CountingSource.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/io/CountingSource.java @@ -39,6 +39,7 @@ import org.apache.beam.sdk.metrics.SourceMetrics; import org.apache.beam.sdk.options.PipelineOptions; import org.apache.beam.sdk.transforms.SerializableFunction; +import org.apache.beam.sdk.transforms.windowing.BoundedWindow; import org.apache.beam.sdk.values.PCollection; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList; import org.checkerframework.checker.nullness.qual.Nullable; @@ -95,7 +96,8 @@ static BoundedSource createSourceForSubrange(long startIndex, long endInde /** Create a new {@link UnboundedCountingSource}. */ // package-private to return a typed UnboundedCountingSource rather than the UnboundedSource type. static UnboundedCountingSource createUnboundedFrom(long start) { - return new UnboundedCountingSource(start, 1, 1L, Duration.ZERO, new NowTimestampFn()); + return new UnboundedCountingSource( + start, 1, Long.MAX_VALUE, 1L, Duration.ZERO, new NowTimestampFn()); } /** @@ -130,7 +132,7 @@ public static UnboundedSource unbounded() { @Deprecated public static UnboundedSource unboundedWithTimestampFn( SerializableFunction timestampFn) { - return new UnboundedCountingSource(0, 1, 1L, Duration.ZERO, timestampFn); + return new UnboundedCountingSource(0, 1, Long.MAX_VALUE, 1L, Duration.ZERO, timestampFn); } ///////////////////////////////////////////////////////////////////////////////////////////// @@ -267,11 +269,13 @@ public void close() throws IOException {} } /** An implementation of {@link CountingSource} that produces an unbounded {@link PCollection}. */ - static class UnboundedCountingSource extends UnboundedSource { + public static class UnboundedCountingSource extends UnboundedSource { /** The first number (>= 0) generated by this {@link UnboundedCountingSource}. */ private final long start; /** The interval between numbers generated by this {@link UnboundedCountingSource}. */ private final long stride; + /** The exclusive limit for the sequence. */ + private final long end; /** The number of elements to produce each period. */ private final long elementsPerPeriod; /** The time between producing numbers from this {@link UnboundedCountingSource}. */ @@ -291,11 +295,13 @@ static class UnboundedCountingSource extends UnboundedSource private UnboundedCountingSource( long start, long stride, + long end, long elementsPerPeriod, Duration period, SerializableFunction timestampFn) { this.start = start; this.stride = stride; + this.end = end; checkArgument( elementsPerPeriod > 0L, "Must produce at least one element per period, got %s", @@ -312,7 +318,8 @@ private UnboundedCountingSource( * will be produced with an interval between them equal to the period. */ public UnboundedCountingSource withRate(long elementsPerPeriod, Duration period) { - return new UnboundedCountingSource(start, stride, elementsPerPeriod, period, timestampFn); + return new UnboundedCountingSource( + start, stride, end, elementsPerPeriod, period, timestampFn); } /** @@ -324,7 +331,18 @@ public UnboundedCountingSource withRate(long elementsPerPeriod, Duration period) public UnboundedCountingSource withTimestampFn( SerializableFunction timestampFn) { checkNotNull(timestampFn); - return new UnboundedCountingSource(start, stride, elementsPerPeriod, period, timestampFn); + return new UnboundedCountingSource( + start, stride, end, elementsPerPeriod, period, timestampFn); + } + + /** + * Returns an {@link UnboundedCountingSource} like this one but with the specified exclusive + * limit. + */ + public UnboundedCountingSource to(long end) { + checkArgument(end >= start, "end (%s) must be >= start (%s)", end, start); + return new UnboundedCountingSource( + start, stride, end, elementsPerPeriod, period, timestampFn); } /** @@ -348,7 +366,7 @@ public List> split( // 0, 2, and 4. splits.add( new UnboundedCountingSource( - start + i * stride, newStride, elementsPerPeriod, period, timestampFn)); + start + i * stride, newStride, end, elementsPerPeriod, period, timestampFn)); } return splits.build(); } @@ -376,6 +394,7 @@ public boolean equals(@Nullable Object other) { UnboundedCountingSource that = (UnboundedCountingSource) other; return this.start == that.start && this.stride == that.stride + && this.end == that.end && this.elementsPerPeriod == that.elementsPerPeriod && Objects.equals(this.period, that.period) && Objects.equals(this.timestampFn, that.timestampFn); @@ -383,7 +402,7 @@ public boolean equals(@Nullable Object other) { @Override public int hashCode() { - return Objects.hash(start, stride, elementsPerPeriod, period, timestampFn); + return Objects.hash(start, stride, end, elementsPerPeriod, period, timestampFn); } } @@ -431,6 +450,9 @@ public boolean advance() throws IOException { return false; } long nextValue = current + source.stride; + if (nextValue >= source.end) { + return false; + } if (expectedValue() < nextValue) { return false; } @@ -453,7 +475,9 @@ private long expectedValue() { @Override public Instant getWatermark() { - return source.timestampFn.apply(current); + return (current >= source.end - source.stride) + ? BoundedWindow.TIMESTAMP_MAX_VALUE + : source.timestampFn.apply(current); } @Override diff --git a/sdks/java/core/src/test/java/org/apache/beam/sdk/io/CountingSourceTest.java b/sdks/java/core/src/test/java/org/apache/beam/sdk/io/CountingSourceTest.java index 70a09083619d..337462340df1 100644 --- a/sdks/java/core/src/test/java/org/apache/beam/sdk/io/CountingSourceTest.java +++ b/sdks/java/core/src/test/java/org/apache/beam/sdk/io/CountingSourceTest.java @@ -21,6 +21,7 @@ import static org.hamcrest.Matchers.is; import static org.hamcrest.Matchers.lessThan; import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertTrue; import java.io.IOException; @@ -315,4 +316,39 @@ public void testUnboundedSourceCheckpointMark() throws Exception { assertEquals(numToSkip + 1, (long) reader.getCurrent()); assertEquals(numToSkip + 1, reader.getCurrentTimestamp().getMillis()); } + + @Test + @Category(NeedsRunner.class) + public void testUnboundedSourceWithFromAndTo() { + long start = 5; + long end = 10; + PCollection input = p.apply(Read.from(CountingSource.createUnboundedFrom(start).to(end))); + + PAssert.that(input).containsInAnyOrder(5L, 6L, 7L, 8L, 9L); + p.run(); + } + + @Test + @Category(NeedsRunner.class) + public void testUnboundedSourceWithFromAndToAndRate() { + long start = 5; + long end = 15; + long elementsPerPeriod = 2; + Duration period = Duration.millis(10); + PCollection input = + p.apply( + Read.from( + CountingSource.createUnboundedFrom(start) + .to(end) + .withRate(elementsPerPeriod, period))); + + PAssert.that(input).containsInAnyOrder(5L, 6L, 7L, 8L, 9L, 10L, 11L, 12L, 13L, 14L); + + Instant startTime = Instant.now(); + p.run(); + Instant endTime = Instant.now(); + + long expectedMinimumMillis = ((end - start) * period.getMillis()) / elementsPerPeriod; + assertFalse(endTime.isBefore(startTime.plus(Duration.millis(expectedMinimumMillis)))); + } } From f33feafb2b2ab46dd73b198a29ea182a1c3c2a00 Mon Sep 17 00:00:00 2001 From: scwhittle Date: Wed, 24 Jun 2026 15:30:10 +0200 Subject: [PATCH 4/6] [Dataflow Java Streaming] Add a timeout to how long commits will retry to the service. (#39085) * [Dataflow Java Streaming] Add a timeout to how long commits will retry to the service. In cases where the service is unavailable this prevents build-up of commits on workers that are no longer relevant. This defaults to 30 minutes but can be disabled via an experiment or by setting the option. --- .../DataflowStreamingPipelineOptions.java | 17 ++++ .../worker/StreamingDataflowWorker.java | 2 + .../client/grpc/GrpcCommitWorkStream.java | 94 ++++++++++++++----- .../grpc/GrpcWindmillStreamFactory.java | 10 ++ .../client/grpc/GrpcCommitWorkStreamTest.java | 82 ++++++++++++++++ 5 files changed, 183 insertions(+), 22 deletions(-) diff --git a/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/options/DataflowStreamingPipelineOptions.java b/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/options/DataflowStreamingPipelineOptions.java index 0897e5263821..cca261fff627 100644 --- a/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/options/DataflowStreamingPipelineOptions.java +++ b/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/options/DataflowStreamingPipelineOptions.java @@ -196,6 +196,13 @@ public interface DataflowStreamingPipelineOptions extends PipelineOptions { void setStuckCommitDurationMillis(int value); + @Description( + "Retry commits on stream errors until this much time has elapsed since the commit was scheduled. If zero, retry forever.") + @Default.InstanceFactory(CommitWorkStreamRetryTimeoutMillisFactory.class) + long getCommitWorkStreamRetryTimeoutMillis(); + + void setCommitWorkStreamRetryTimeoutMillis(long value); + @Description( "Period for sending 'global get config' requests to the service. The duration is " + "specified as seconds in 'PTx.yS' format, e.g. 'PT5.125S'." @@ -343,4 +350,14 @@ public Boolean create(PipelineOptions options) { return ExperimentalOptions.hasExperiment(options, "enable_windmill_service_direct_path"); } } + + class CommitWorkStreamRetryTimeoutMillisFactory implements DefaultValueFactory { + @Override + public Long create(PipelineOptions options) { + if (ExperimentalOptions.hasExperiment(options, "disable_commit_retry_timeout")) { + return 0L; + } + return Duration.standardMinutes(30).getMillis(); + } + } } diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java index 9e82343474c6..9f063d393703 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java @@ -756,6 +756,8 @@ public static StreamingDataflowWorker fromOptions(DataflowWorkerHarnessOptions o new WorkHeartbeatResponseProcessor(computationStateCache::get)) .setHealthCheckIntervalMillis( options.getWindmillServiceStreamingRpcHealthCheckPeriodMs()) + .setCommitWorkStreamRetryTimeout( + java.time.Duration.ofMillis(options.getCommitWorkStreamRetryTimeoutMillis())) .build(); return ConfigFetcherComputationStateCacheAndWindmillClient.builder() .setWindmillDispatcherClient(dispatcherClient) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcCommitWorkStream.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcCommitWorkStream.java index d24676652fd8..160b0cce0133 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcCommitWorkStream.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcCommitWorkStream.java @@ -20,7 +20,6 @@ import static org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkNotNull; import static org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkState; -import com.google.auto.value.AutoValue; import java.io.PrintWriter; import java.time.Duration; import java.util.HashMap; @@ -77,6 +76,7 @@ private static class StreamAndRequest { private final JobHeader jobHeader; private final int streamingRpcBatchLimit; private volatile boolean logMissingResponse = true; + private final Duration maxRetryDuration; private GrpcCommitWorkStream( String backendWorkerToken, @@ -90,6 +90,7 @@ private GrpcCommitWorkStream( AtomicLong idGenerator, int streamingRpcBatchLimit, Duration halfClosePhysicalStreamAfter, + Duration maxRetryDuration, ScheduledExecutorService executor) { super( LOG, @@ -104,6 +105,7 @@ private GrpcCommitWorkStream( this.idGenerator = idGenerator; this.jobHeader = jobHeader; this.streamingRpcBatchLimit = streamingRpcBatchLimit; + this.maxRetryDuration = maxRetryDuration; } static GrpcCommitWorkStream create( @@ -118,6 +120,7 @@ static GrpcCommitWorkStream create( AtomicLong idGenerator, int streamingRpcBatchLimit, Duration halfClosePhysicalStreamAfter, + Duration maxRetryDuration, ScheduledExecutorService executor) { return new GrpcCommitWorkStream( backendWorkerToken, @@ -130,6 +133,7 @@ static GrpcCommitWorkStream create( idGenerator, streamingRpcBatchLimit, halfClosePhysicalStreamAfter, + maxRetryDuration, executor); } @@ -224,14 +228,48 @@ public void onResponse(StreamingCommitResponse response) { failureHandler.throwIfNonEmpty(); } - @Override @SuppressWarnings("ReferenceEquality") + private boolean belongsToThisHandler(StreamAndRequest streamAndRequest) { + return streamAndRequest.handler == this; + } + + @Override public boolean hasPendingRequests() { - return pending.entrySet().stream().anyMatch(e -> e.getValue().handler == this); + return pending.entrySet().stream().anyMatch(e -> belongsToThisHandler(e.getValue())); } @Override + @SuppressWarnings("ReferenceEquality") public void onDone(Status status) { + if (maxRetryDuration.compareTo(Duration.ZERO) > 0) { + // Remove the requests that have exceeded the retry time so they are not retried. + long nowNanos = System.nanoTime(); + long maxRetryDurationNanos = maxRetryDuration.toNanos(); + Iterator> iterator = pending.entrySet().iterator(); + int keptRequests = 0, removedRequests = 0; + while (iterator.hasNext()) { + StreamAndRequest streamAndRequest = checkNotNull(iterator.next().getValue()); + PendingRequest pendingRequest = streamAndRequest.request; + if (!belongsToThisHandler(streamAndRequest) + || nowNanos - pendingRequest.getStartTimeNanos() < maxRetryDurationNanos) { + ++keptRequests; + continue; + } + ++removedRequests; + iterator.remove(); + try { + pendingRequest.completeWithStatus(CommitStatus.ABORTED); + } catch (RuntimeException e) { + LOG.warn("Exception while aborting commit due to retry timeout.", e); + } + } + if (removedRequests > 0) { + LOG.info( + "Aborting {} commits which have exceeded retry deadline, kept {}. Work will be retried as needed by service.", + removedRequests, + keptRequests); + } + } if (status.isOk() && hasPendingRequests()) { LOG.warn("Unexpected requests without responses on drained physical stream, retrying."); } @@ -270,7 +308,7 @@ private void flushInternal(Map requests) if (requests.size() == 1) { Map.Entry elem = requests.entrySet().iterator().next(); - if (elem.getValue().request().getSerializedSize() + if (elem.getValue().getRequest().getSerializedSize() > AbstractWindmillStream.RPC_STREAM_CHUNK_SIZE) { issueMultiChunkRequest(elem.getKey(), elem.getValue()); } else { @@ -286,7 +324,7 @@ private void issueSingleRequest(long id, PendingRequest pendingRequest) StreamingCommitWorkRequest.Builder requestBuilder = StreamingCommitWorkRequest.newBuilder(); requestBuilder .addCommitChunkBuilder() - .setComputationId(pendingRequest.computationId()) + .setComputationId(pendingRequest.getComputationId()) .setRequestId(id) .setShardingKey(pendingRequest.shardingKey()) .setSerializedWorkItemCommit(pendingRequest.serializedCommit()); @@ -311,9 +349,9 @@ private void issueBatchedRequest(Map requests) for (Map.Entry entry : requests.entrySet()) { PendingRequest request = entry.getValue(); StreamingCommitRequestChunk.Builder chunkBuilder = requestBuilder.addCommitChunkBuilder(); - if (lastComputation == null || !lastComputation.equals(request.computationId())) { - chunkBuilder.setComputationId(request.computationId()); - lastComputation = request.computationId(); + if (lastComputation == null || !lastComputation.equals(request.getComputationId())) { + chunkBuilder.setComputationId(request.getComputationId()); + lastComputation = request.getComputationId(); } chunkBuilder .setRequestId(entry.getKey()) @@ -338,7 +376,7 @@ private void issueBatchedRequest(Map requests) private void issueMultiChunkRequest(long id, PendingRequest pendingRequest) throws WindmillStreamShutdownException { - checkNotNull(pendingRequest.computationId(), "Cannot commit WorkItem w/o a computationId."); + checkNotNull(pendingRequest.getComputationId(), "Cannot commit WorkItem w/o a computationId."); ByteString serializedCommit = pendingRequest.serializedCommit(); synchronized (this) { if (isShutdown) { @@ -359,7 +397,7 @@ private void issueMultiChunkRequest(long id, PendingRequest pendingRequest) StreamingCommitRequestChunk.newBuilder() .setRequestId(id) .setSerializedWorkItemCommit(chunk) - .setComputationId(pendingRequest.computationId()) + .setComputationId(pendingRequest.getComputationId()) .setShardingKey(pendingRequest.shardingKey()); int remaining = serializedCommit.size() - end; if (remaining > 0) { @@ -376,34 +414,46 @@ private void issueMultiChunkRequest(long id, PendingRequest pendingRequest) } } - @AutoValue - abstract static class PendingRequest { + private static class PendingRequest { + private final String computationId; + private final WorkItemCommitRequest request; + private final Consumer onDone; + private final long startTimeNanos; // System.nanoTime() of when request began. - private static PendingRequest create( + private PendingRequest( String computationId, WorkItemCommitRequest request, Consumer onDone) { - return new AutoValue_GrpcCommitWorkStream_PendingRequest(computationId, request, onDone); + this.computationId = computationId; + this.request = request; + this.onDone = onDone; + this.startTimeNanos = System.nanoTime(); } - abstract String computationId(); + String getComputationId() { + return computationId; + } - abstract WorkItemCommitRequest request(); + WorkItemCommitRequest getRequest() { + return request; + } - abstract Consumer onDone(); + long getStartTimeNanos() { + return startTimeNanos; + } private long getBytes() { - return (long) request().getSerializedSize() + computationId().length(); + return (long) request.getSerializedSize() + computationId.length(); } private ByteString serializedCommit() { - return request().toByteString(); + return request.toByteString(); } private void completeWithStatus(CommitStatus commitStatus) { - onDone().accept(commitStatus); + onDone.accept(commitStatus); } private long shardingKey() { - return request().getShardingKey(); + return request.getShardingKey(); } private void abort() { @@ -462,7 +512,7 @@ public boolean commitWorkItem( return false; } - PendingRequest request = PendingRequest.create(computation, commitRequest, onDone); + PendingRequest request = new PendingRequest(computation, commitRequest, onDone); add(idGenerator.incrementAndGet(), request); return true; } diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcWindmillStreamFactory.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcWindmillStreamFactory.java index 0184b88d53cd..97ca3c4e83d7 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcWindmillStreamFactory.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcWindmillStreamFactory.java @@ -100,6 +100,7 @@ public class GrpcWindmillStreamFactory implements StatusDataProvider { private final Consumer> processHeartbeatResponses; private final java.time.Duration directStreamingRpcPhysicalStreamHalfCloseAfter; private final Supplier executorServiceSupplier; + private final java.time.Duration commitWorkStreamRetryTimeout; private GrpcWindmillStreamFactory( JobHeader jobHeader, @@ -111,6 +112,7 @@ private GrpcWindmillStreamFactory( Consumer> processHeartbeatResponses, Supplier maxBackOffSupplier, java.time.Duration directStreamingRpcPhysicalStreamHalfCloseAfter, + java.time.Duration commitWorkStreamRetryTimeout, Supplier executorServiceSupplier) { this.jobHeader = jobHeader; this.logEveryNStreamFailures = logEveryNStreamFailures; @@ -132,6 +134,7 @@ private GrpcWindmillStreamFactory( this.directStreamingRpcPhysicalStreamHalfCloseAfter = directStreamingRpcPhysicalStreamHalfCloseAfter; this.executorServiceSupplier = executorServiceSupplier; + this.commitWorkStreamRetryTimeout = commitWorkStreamRetryTimeout; } /** @implNote Used for {@link AutoBuilder} {@link Builder} class, do not call directly. */ @@ -146,6 +149,7 @@ static GrpcWindmillStreamFactory create( Supplier maxBackOffSupplier, int healthCheckIntervalMillis, java.time.Duration directStreamingRpcPhysicalStreamHalfCloseAfter, + java.time.Duration commitWorkStreamRetryTimeout, Supplier scheduledExecutorServiceSupplier) { GrpcWindmillStreamFactory streamFactory = new GrpcWindmillStreamFactory( @@ -158,6 +162,7 @@ static GrpcWindmillStreamFactory create( processHeartbeatResponses, maxBackOffSupplier, directStreamingRpcPhysicalStreamHalfCloseAfter, + commitWorkStreamRetryTimeout, scheduledExecutorServiceSupplier); if (healthCheckIntervalMillis >= 0) { @@ -200,6 +205,7 @@ public static GrpcWindmillStreamFactory.Builder of(JobHeader jobHeader) { .setProcessHeartbeatResponses(ignored -> {}) .setDirectStreamingRpcPhysicalStreamHalfCloseAfter( DEFAULT_DIRECT_STREAMING_RPC_PHYSICAL_STREAM_HALF_CLOSE_AFTER) + .setCommitWorkStreamRetryTimeout(java.time.Duration.ZERO) .setScheduledExecutorServiceSupplier(() -> null); } @@ -347,6 +353,7 @@ public CommitWorkStream createCommitWorkStream(CloudWindmillServiceV1Alpha1Stub streamIdGenerator, streamingRpcBatchLimit, java.time.Duration.ZERO, + commitWorkStreamRetryTimeout, executorForDispatchedStreams("CommitWork")); } @@ -363,6 +370,7 @@ public CommitWorkStream createDirectCommitWorkStream(WindmillConnection connecti streamIdGenerator, streamingRpcBatchLimit, directStreamingRpcPhysicalStreamHalfCloseAfter, + java.time.Duration.ZERO, executorForDirectStreams(connection.backendWorkerToken(), "CommitWork")); } @@ -426,6 +434,8 @@ Builder setProcessHeartbeatResponses( Builder setDirectStreamingRpcPhysicalStreamHalfCloseAfter(java.time.Duration timeout); + Builder setCommitWorkStreamRetryTimeout(java.time.Duration timeout); + Builder setScheduledExecutorServiceSupplier( Supplier scheduledExecutorServiceSupplier); diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcCommitWorkStreamTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcCommitWorkStreamTest.java index e9fd55fa5668..9c3d5c9c3ef3 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcCommitWorkStreamTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcCommitWorkStreamTest.java @@ -38,6 +38,7 @@ import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicReference; import java.util.function.Supplier; import org.apache.beam.runners.dataflow.worker.windmill.CloudWindmillServiceV1Alpha1Grpc; import org.apache.beam.runners.dataflow.worker.windmill.Windmill; @@ -1133,6 +1134,87 @@ public void testCommitWorkItem_multiplePhysicalStreams_multipleHandovers_halfClo assertTrue(commitWorkStream.awaitTermination(10, TimeUnit.SECONDS)); } + @Test + public void testCommitWorkItem_stopsRetriesAfterDuration() throws Exception { + int numCommits = 1; + CountDownLatch commitProcessed = new CountDownLatch(numCommits); + AtomicReference commitStatus = new AtomicReference<>(); + GrpcCommitWorkStream commitWorkStream = + (GrpcCommitWorkStream) + GrpcWindmillStreamFactory.of(TEST_JOB_HEADER) + .setDirectStreamingRpcPhysicalStreamHalfCloseAfter(Duration.ofMinutes(1)) + .setCommitWorkStreamRetryTimeout(Duration.ofNanos(1)) + .build() + .createCommitWorkStream(CloudWindmillServiceV1Alpha1Grpc.newStub(inProcessChannel)); + commitWorkStream.start(); + + FakeWindmillGrpcService.CommitStreamInfo streamInfo = waitForConnectionAndConsumeHeader(); + try (WindmillStream.CommitWorkStream.RequestBatcher batcher = commitWorkStream.batcher()) { + assertTrue( + batcher.commitWorkItem( + COMPUTATION_ID, + workItemCommitRequest(0), + status -> { + commitStatus.set(status); + commitProcessed.countDown(); + })); + } + + // The next request should have some chunks. + assertThat(streamInfo.requests.take().getCommitChunkList()).isNotEmpty(); + + // We won't get responses so we will have some pending requests. + assertThat(commitProcessed.getCount()).isGreaterThan(0); + streamInfo.responseObserver.onError(new IOException("test error")); + commitProcessed.await(); + assertThat(commitStatus.get()).isEqualTo(Windmill.CommitStatus.ABORTED); + } + + @Test + public void testCommitWorkItem_retriesWithLongerCommitRetryTimeout() throws Exception { + int numCommits = 1; + CountDownLatch commitProcessed = new CountDownLatch(numCommits); + AtomicReference commitStatus = new AtomicReference<>(); + GrpcCommitWorkStream commitWorkStream = + (GrpcCommitWorkStream) + GrpcWindmillStreamFactory.of(TEST_JOB_HEADER) + .setDirectStreamingRpcPhysicalStreamHalfCloseAfter(Duration.ofMinutes(1)) + // Verifies that if this is set but is not exceeded that a retry occurs. + .setCommitWorkStreamRetryTimeout(Duration.ofMinutes(10)) + .build() + .createCommitWorkStream(CloudWindmillServiceV1Alpha1Grpc.newStub(inProcessChannel)); + commitWorkStream.start(); + + FakeWindmillGrpcService.CommitStreamInfo streamInfo = waitForConnectionAndConsumeHeader(); + try (WindmillStream.CommitWorkStream.RequestBatcher batcher = commitWorkStream.batcher()) { + assertTrue( + batcher.commitWorkItem( + COMPUTATION_ID, + workItemCommitRequest(0), + status -> { + commitStatus.set(status); + commitProcessed.countDown(); + })); + } + + // The next request should have some chunks. + assertThat(streamInfo.requests.take().getCommitChunkList()).isNotEmpty(); + + // We won't get responses so we will have some pending requests. + assertThat(commitProcessed.getCount()).isGreaterThan(0); + streamInfo.responseObserver.onError(new IOException("test error")); + + // The stream should reconnect and retry the requests. + FakeWindmillGrpcService.CommitStreamInfo reconnectStreamInfo = + waitForConnectionAndConsumeHeader(); + Windmill.StreamingCommitWorkRequest reconnectRequest = reconnectStreamInfo.requests.take(); + assertEquals(1, reconnectRequest.getCommitChunkCount()); + reconnectStreamInfo.responseObserver.onNext( + Windmill.StreamingCommitResponse.newBuilder().addRequestId(1).build()); + commitProcessed.await(); + assertThat(commitStatus.get()).isEqualTo(Windmill.CommitStatus.OK); + } + private FakeWindmillGrpcService.CommitStreamInfo waitForConnectionAndConsumeHeader() { try { FakeWindmillGrpcService.CommitStreamInfo info = fakeService.waitForConnectedCommitStream(); From 857dfc9a87522f435949f118fff5271b3838befa Mon Sep 17 00:00:00 2001 From: Tobias Kaymak Date: Wed, 24 Jun 2026 15:44:06 +0200 Subject: [PATCH 5/6] [mqtt] Fix streaming xlang IT failing the Messaging PostCommit (#39088) * [mqtt] Fix streaming xlang IT: amend test pipeline options, run non-blocking --- ...tCommit_Python_Xlang_Messaging_Direct.json | 2 +- .../io/external/xlang_mqttio_it_test.py | 23 +++++++++++++------ 2 files changed, 17 insertions(+), 8 deletions(-) diff --git a/.github/trigger_files/beam_PostCommit_Python_Xlang_Messaging_Direct.json b/.github/trigger_files/beam_PostCommit_Python_Xlang_Messaging_Direct.json index e3d6056a5de9..c537844dc84a 100644 --- a/.github/trigger_files/beam_PostCommit_Python_Xlang_Messaging_Direct.json +++ b/.github/trigger_files/beam_PostCommit_Python_Xlang_Messaging_Direct.json @@ -1,4 +1,4 @@ { "comment": "Modify this file in a trivial way to cause this test suite to run", - "modification": 1 + "modification": 3 } diff --git a/sdks/python/apache_beam/io/external/xlang_mqttio_it_test.py b/sdks/python/apache_beam/io/external/xlang_mqttio_it_test.py index fa6ebed06efe..889096d1f22b 100644 --- a/sdks/python/apache_beam/io/external/xlang_mqttio_it_test.py +++ b/sdks/python/apache_beam/io/external/xlang_mqttio_it_test.py @@ -33,7 +33,7 @@ import pytest import apache_beam as beam -from apache_beam.options.pipeline_options import PipelineOptions +from apache_beam.options.pipeline_options import PortableOptions from apache_beam.options.pipeline_options import StandardOptions from apache_beam.testing.test_pipeline import TestPipeline from apache_beam.typehints.row_type import RowTypeConstraint @@ -214,12 +214,21 @@ def subscribe(): publisher.start() subscriber.start() - options = PipelineOptions([ - '--runner=PrismRunner', - '--environment_type=LOOPBACK', - '--streaming', - ]) - p = TestPipeline(options=options) + # MqttIO read is unbounded, so this pipeline runs in streaming mode and + # never terminates on its own. Amend the harness-provided pipeline options + # rather than discarding them: enable streaming, run non-blocking so the + # observe-then-cancel logic below can execute, and target the Prism portable + # runner. The latter is required because SwitchingDirectRunner disables its + # Prism delegation for pipelines containing external (cross-language) + # transforms (see runners/direct/direct_runner.py) and falls back to the + # BundleBasedDirectRunner, which cannot execute an unbounded read. + # The runner is instantiated during TestPipeline construction, so it must be + # passed to the constructor; the remaining harness-provided options are + # preserved and only amended (streaming + LOOPBACK environment) afterwards. + p = TestPipeline(runner='PrismRunner', blocking=False) + p.get_pipeline_options().view_as(StandardOptions).streaming = True + p.get_pipeline_options().view_as( + PortableOptions).environment_type = 'LOOPBACK' p.not_use_test_runner_api = True _ = ( p From 0802263e48e842bfbe437ed9b8ec72c2311c0d76 Mon Sep 17 00:00:00 2001 From: Shunping Huang Date: Wed, 24 Jun 2026 10:42:05 -0400 Subject: [PATCH 6/6] Add a github workflow to publish vllm image. (#39089) * Add a github workflow to publish vllm image. * Add the new github workflow the workflows/README.md --- .github/workflows/README.md | 1 + .../beam_Publish_Python_VLLM_Image.yml | 75 +++++++++++++++++++ 2 files changed, 76 insertions(+) create mode 100644 .github/workflows/beam_Publish_Python_VLLM_Image.yml diff --git a/.github/workflows/README.md b/.github/workflows/README.md index 1715365f4ec6..4a32aa271df1 100644 --- a/.github/workflows/README.md +++ b/.github/workflows/README.md @@ -534,6 +534,7 @@ PostCommit Jobs run in a schedule against master branch and generally do not get | [ Publish Beam SDK Snapshots ](https://github.com/apache/beam/actions/workflows/beam_Publish_Beam_SDK_Snapshots.yml) | N/A | [![.github/workflows/beam_Publish_Beam_SDK_Snapshots.yml](https://github.com/apache/beam/actions/workflows/beam_Publish_Beam_SDK_Snapshots.yml/badge.svg?event=schedule)](https://github.com/apache/beam/actions/workflows/beam_Publish_Beam_SDK_Snapshots.yml?query=event%3Aschedule) | | [ Publish BeamMetrics ](https://github.com/apache/beam/actions/workflows/beam_Publish_BeamMetrics.yml) | N/A | [![.github/workflows/beam_Publish_BeamMetrics.yml](https://github.com/apache/beam/actions/workflows/beam_Publish_BeamMetrics.yml/badge.svg?event=schedule)](https://github.com/apache/beam/actions/workflows/beam_Publish_BeamMetrics.yml?query=event%3Aschedule) | [ Publish Docker Snapshots ](https://github.com/apache/beam/actions/workflows/beam_Publish_Docker_Snapshots.yml) | N/A | [![.github/workflows/beam_Publish_Docker_Snapshots.yml](https://github.com/apache/beam/actions/workflows/beam_Publish_Docker_Snapshots.yml/badge.svg?event=schedule)](https://github.com/apache/beam/actions/workflows/beam_Publish_Docker_Snapshots.yml?query=event%3Aschedule) | +| [ Publish Python VLLM Image ](https://github.com/apache/beam/actions/workflows/beam_Publish_Python_VLLM_Image.yml) | N/A | [![.github/workflows/beam_Publish_Python_VLLM_Image.yml](https://github.com/apache/beam/actions/workflows/beam_Publish_Python_VLLM_Image.yml/badge.svg?event=schedule)](https://github.com/apache/beam/actions/workflows/beam_Publish_Python_VLLM_Image.yml?query=event%3Aschedule) | | [ Publish Website ](https://github.com/apache/beam/actions/workflows/beam_Publish_Website.yml) | N/A | [![.github/workflows/beam_Publish_Website.yml](https://github.com/apache/beam/actions/workflows/beam_Publish_Website.yml/badge.svg?event=schedule)](https://github.com/apache/beam/actions/workflows/beam_Publish_Website.yml?query=event%3Aschedule) | | [ Release Nightly Snapshot ](https://github.com/apache/beam/actions/workflows/beam_Release_NightlySnapshot.yml) | N/A | [![.github/workflows/beam_Release_NightlySnapshot.yml](https://github.com/apache/beam/actions/workflows/beam_Release_NightlySnapshot.yml/badge.svg?event=schedule)](https://github.com/apache/beam/actions/workflows/beam_Release_NightlySnapshot.yml?query=event%3Aschedule) | | [ Release Nightly Snapshot Python ](https://github.com/apache/beam/actions/workflows/beam_Release_Python_NightlySnapshot.yml) | N/A | [![.github/workflows/beam_Release_Python_NightlySnapshot.yml](https://github.com/apache/beam/actions/workflows/beam_Release_Python_NightlySnapshot.yml/badge.svg?event=schedule)](https://github.com/apache/beam/actions/workflows/beam_Release_Python_NightlySnapshot.yml?query=event%3Aschedule) | diff --git a/.github/workflows/beam_Publish_Python_VLLM_Image.yml b/.github/workflows/beam_Publish_Python_VLLM_Image.yml new file mode 100644 index 000000000000..21b6dc8d53c8 --- /dev/null +++ b/.github/workflows/beam_Publish_Python_VLLM_Image.yml @@ -0,0 +1,75 @@ +# 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. + +name: Publish Python VLLM Image + +on: + schedule: + - cron: '0 0 * * 0' # Run weekly (Sunday at 00:00) + workflow_dispatch: + +# Setting explicit permissions for the action to avoid the default permissions +permissions: + actions: write + pull-requests: read + checks: read + contents: read + deployments: read + id-token: none + issues: read + discussions: read + packages: read + pages: read + repository-projects: read + security-events: read + statuses: read + +# This allows a subsequently queued workflow run to interrupt previous runs +concurrency: + group: '${{ github.workflow }} @ ${{ github.event.issue.number || github.sha || github.head_ref || github.ref }}-${{ github.event.schedule || github.event.comment.id || github.event.sender.login }}' + cancel-in-progress: true + +jobs: + build_and_push_image: + if: | + github.event_name == 'workflow_dispatch' || + (github.event_name == 'schedule' && github.repository == 'apache/beam') + runs-on: [self-hosted, ubuntu-24.04, main] + timeout-minutes: 60 + steps: + - name: Checkout code + uses: actions/checkout@v7 + - name: Authenticate on GCP + uses: google-github-actions/auth@7c6bc770dae815cd3e89ee6cdf493a5fab2cc093 + with: + service_account: ${{ secrets.GCP_SA_EMAIL }} + credentials_json: ${{ secrets.GCP_SA_KEY }} + - name: Set up Cloud SDK + uses: google-github-actions/setup-gcloud@aa5489c8933f4cc7a4f7d45035b3b1440c9c10db + - name: Set up Docker Buildx + uses: docker/setup-buildx-action@d7f5e7f509e45cec5c76c4d5afdd7de93d0b3df5 + - name: GCloud Docker credential helper + run: | + gcloud auth configure-docker us.gcr.io + - name: Build and Push Multi-Arch Image + run: | + docker buildx build \ + --platform linux/amd64,linux/arm64 \ + -t us.gcr.io/apache-beam-testing/python-postcommit-it/vllm:latest \ + -f sdks/python/apache_beam/ml/inference/test_resources/vllm.dockerfile.old \ + --push \ + .