From f2700d66157b9d3212b2311bf67bfba8040248b6 Mon Sep 17 00:00:00 2001 From: Sania Parveen Date: Wed, 5 Aug 2026 21:25:45 +0000 Subject: [PATCH 1/4] Make primary channel timeout configurable --- .../worker/StreamingDataflowWorker.java | 14 ++- .../client/grpc/stubs/FailoverChannel.java | 51 +++++++---- .../grpc/stubs/FailoverChannelTest.java | 88 ++++++++++++++++++- .../windmill/src/main/proto/windmill.proto | 4 + 4 files changed, 137 insertions(+), 20 deletions(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java index 9e82343474c6..b9217b520749 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 @@ -34,6 +34,7 @@ import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicLong; import java.util.concurrent.atomic.AtomicReference; import java.util.function.Consumer; import java.util.function.Function; @@ -825,6 +826,7 @@ private static ChannelCache createChannelCache( DataflowWorkerHarnessOptions workerOptions, ComputationConfig.Fetcher configFetcher, GrpcDispatcherClient dispatcherClient) { + AtomicLong primaryNotReadyWaitNanos = new AtomicLong(TimeUnit.SECONDS.toNanos(15)); ChannelCache channelCache = ChannelCache.create( (currentFlowControlSettings, serviceAddress) -> { @@ -844,16 +846,20 @@ private static ChannelCache createChannelCache( workerOptions.getWindmillServiceRpcChannelAliveTimeoutSec(), currentFlowControlSettings), MoreCallCredentials.from( - new VendoredCredentialsAdapter(workerOptions.getGcpCredential()))), + new VendoredCredentialsAdapter(workerOptions.getGcpCredential())), + primaryNotReadyWaitNanos::get), currentFlowControlSettings.getOnReadyThresholdBytes()); }); configFetcher .getGlobalConfigHandle() .registerConfigObserver( - config -> - channelCache.consumeFlowControlSettings( - config.userWorkerJobSettings().getFlowControlSettings())); + config -> { + primaryNotReadyWaitNanos.set( + config.userWorkerJobSettings().getPrimaryNotReadyWaitNanos()); + channelCache.consumeFlowControlSettings( + config.userWorkerJobSettings().getFlowControlSettings()); + }); return channelCache; } diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/stubs/FailoverChannel.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/stubs/FailoverChannel.java index faa08c497c8f..8bb83e12eaf5 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/stubs/FailoverChannel.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/stubs/FailoverChannel.java @@ -17,6 +17,7 @@ */ package org.apache.beam.runners.dataflow.worker.windmill.client.grpc.stubs; +import java.net.http.WebSocket.Listener; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; @@ -45,9 +46,9 @@ *

Routes requests to either primary or fallback channel based on two independent failover modes: * *