From 204175ba5a8f9deed26f4f9b1e0a6be146fed331 Mon Sep 17 00:00:00 2001 From: Radek Stankiewicz Date: Thu, 9 Oct 2025 14:16:11 +0200 Subject: [PATCH 1/3] propagate drainmode up to timerData creation in windmillTimerToTimerData --- .../worker/StreamingDataflowWorker.java | 2 ++ .../worker/StreamingModeExecutionContext.java | 14 +++++++++-- .../worker/UngroupedWindmillReader.java | 21 +++++++++++++--- .../worker/WindmillKeyedWorkItem.java | 24 +++++++++++++++---- .../worker/WindmillTimerInternals.java | 6 ++++- .../worker/WindowingWindmillReader.java | 3 ++- .../dataflow/worker/streaming/Work.java | 11 ++++++++- .../harness/SingleSourceWorkerHarness.java | 3 +++ .../grpc/GetWorkResponseChunkAssembler.java | 5 +++- .../client/grpc/GrpcDirectGetWorkStream.java | 1 + .../client/grpc/GrpcGetWorkStream.java | 1 + .../windmill/work/WorkItemReceiver.java | 1 + .../windmill/work/WorkItemScheduler.java | 2 ++ .../processing/StreamingWorkScheduler.java | 4 +++- .../dataflow/worker/FakeWindmillServer.java | 1 + .../worker/StreamingDataflowWorkerTest.java | 2 ++ .../StreamingGroupAlsoByWindowFnsTest.java | 2 +- ...ngGroupAlsoByWindowsReshuffleDoFnTest.java | 2 +- .../StreamingModeExecutionContextTest.java | 1 + .../worker/WindmillKeyedWorkItemTest.java | 4 ++-- .../worker/WindmillTimerInternalsTest.java | 6 +++-- .../worker/WorkerCustomSourcesTest.java | 2 ++ .../worker/streaming/ActiveWorkStateTest.java | 2 ++ .../streaming/ComputationStateCacheTest.java | 1 + ...anOutStreamingEngineWorkerHarnessTest.java | 1 + .../harness/WindmillStreamSenderTest.java | 1 + .../worker/util/BoundedQueueExecutorTest.java | 1 + .../StreamingApplianceWorkCommitterTest.java | 1 + .../StreamingEngineWorkCommitterTest.java | 1 + .../grpc/GrpcDirectGetWorkStreamTest.java | 12 ++++++++-- .../client/grpc/GrpcWindmillServerTest.java | 2 ++ .../failures/WorkFailureProcessorTest.java | 1 + .../work/refresh/ActiveWorkRefresherTest.java | 1 + 33 files changed, 120 insertions(+), 22 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 83e924514b59..aad27b869863 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 @@ -389,6 +389,7 @@ private StreamingWorkerHarnessFactoryOutput createFanOutStreamingEngineWorkerHar serializedWorkItemSize, watermarks, processingContext, + drainMode, getWorkStreamLatencies) -> computationStateCache .get(processingContext.computationId()) @@ -401,6 +402,7 @@ private StreamingWorkerHarnessFactoryOutput createFanOutStreamingEngineWorkerHar serializedWorkItemSize, watermarks, processingContext, + drainMode, getWorkStreamLatencies); }), ChannelCachingRemoteStubFactory.create(options.getGcpCredential(), channelCache), 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 f5157bb46956..e3424e3d6670 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 @@ -196,6 +196,10 @@ public boolean workIsFailed() { return work != null && work.isFailed(); } + public boolean getDrainMode() { + return work != null ? work.getDrainMode() : false; + } + public boolean offsetBasedDeduplicationSupported() { return activeReader != null && activeReader.getCurrentSource().offsetBasedDeduplicationSupported(); @@ -820,7 +824,10 @@ public TimerData getNextFiredTimer(Coder windowCode .transform( timer -> WindmillTimerInternals.windmillTimerToTimerData( - WindmillNamespacePrefix.SYSTEM_NAMESPACE_PREFIX, timer, windowCoder)) + WindmillNamespacePrefix.SYSTEM_NAMESPACE_PREFIX, + timer, + windowCoder, + getDrainMode())) .iterator(); } @@ -880,7 +887,10 @@ public TimerData getNextFiredUserTimer(Coder window .transform( timer -> WindmillTimerInternals.windmillTimerToTimerData( - WindmillNamespacePrefix.USER_NAMESPACE_PREFIX, timer, windowCoder)) + WindmillNamespacePrefix.USER_NAMESPACE_PREFIX, + timer, + windowCoder, + getDrainMode())) .iterator()); } diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/UngroupedWindmillReader.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/UngroupedWindmillReader.java index a9a033c89ad7..3295347249b0 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/UngroupedWindmillReader.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/UngroupedWindmillReader.java @@ -24,6 +24,7 @@ import java.io.InputStream; import java.util.Collection; import java.util.Map; +import org.apache.beam.model.fnexecution.v1.BeamFnApi; import org.apache.beam.runners.dataflow.util.CloudObject; import org.apache.beam.runners.dataflow.worker.util.common.worker.NativeReader; import org.apache.beam.runners.dataflow.worker.windmill.Windmill; @@ -117,8 +118,14 @@ protected WindowedValue decodeMessage(Windmill.Message message) throws IOExce Collection windows = WindmillSink.decodeMetadataWindows(windowsCoder, message.getMetadata()); PaneInfo paneInfo = WindmillSink.decodeMetadataPane(message.getMetadata()); + boolean drainingValueFromUpstream = false; if (WindowedValues.WindowedValueCoder.isMetadataSupported()) { - WindmillSink.decodeAdditionalMetadata(windowsCoder, message.getMetadata()); + BeamFnApi.Elements.ElementMetadata elementMetadata = + WindmillSink.decodeAdditionalMetadata(windowsCoder, message.getMetadata()); + drainingValueFromUpstream = + elementMetadata.hasDrain() + ? (elementMetadata.getDrain() == BeamFnApi.Elements.DrainMode.Enum.DRAINING) + : false; } if (valueCoder instanceof KvCoder) { KvCoder kvCoder = (KvCoder) valueCoder; @@ -129,11 +136,19 @@ protected WindowedValue decodeMessage(Windmill.Message message) throws IOExce T result = (T) KV.of(decode(kvCoder.getKeyCoder(), key), decode(kvCoder.getValueCoder(), data)); // todo #33176 propagate metadata to windowed value - return WindowedValues.of(result, timestampMillis, windows, paneInfo); + return WindowedValues.of( + result, timestampMillis, windows, paneInfo, null, null, drainingValueFromUpstream); } else { notifyElementRead(data.available() + metadata.available()); // todo #33176 propagate metadata to windowed value - return WindowedValues.of(decode(valueCoder, data), timestampMillis, windows, paneInfo); + return WindowedValues.of( + decode(valueCoder, data), + timestampMillis, + windows, + paneInfo, + null, + null, + drainingValueFromUpstream); } } diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItem.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItem.java index 6690377d3de6..5cc3b39df12f 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItem.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItem.java @@ -24,6 +24,7 @@ import java.util.Collection; import java.util.List; import java.util.Objects; +import org.apache.beam.model.fnexecution.v1.BeamFnApi; import org.apache.beam.runners.core.KeyedWorkItem; import org.apache.beam.runners.core.KeyedWorkItemCoder; import org.apache.beam.runners.core.TimerInternals.TimerData; @@ -60,6 +61,7 @@ public class WindmillKeyedWorkItem implements KeyedWorkItem private final Windmill.WorkItem workItem; private final K key; + private final boolean drainMode; private final transient Coder windowCoder; private final transient Coder> windowsCoder; @@ -70,12 +72,14 @@ public WindmillKeyedWorkItem( Windmill.WorkItem workItem, Coder windowCoder, Coder> windowsCoder, - Coder valueCoder) { + Coder valueCoder, + boolean drainMode) { this.key = key; this.workItem = workItem; this.windowCoder = windowCoder; this.windowsCoder = windowsCoder; this.valueCoder = valueCoder; + this.drainMode = drainMode; } @Override @@ -93,7 +97,10 @@ public Iterable timersIterable() { .transform( timer -> WindmillTimerInternals.windmillTimerToTimerData( - WindmillNamespacePrefix.SYSTEM_NAMESPACE_PREFIX, timer, windowCoder)); + WindmillNamespacePrefix.SYSTEM_NAMESPACE_PREFIX, + timer, + windowCoder, + drainMode)); } @Override @@ -108,13 +115,22 @@ public Iterable> elementsIterable() { Collection windows = WindmillSink.decodeMetadataWindows(windowsCoder, message.getMetadata()); PaneInfo paneInfo = WindmillSink.decodeMetadataPane(message.getMetadata()); + // Draining value is based on upstream data + boolean drainingValueFromUpstream = false; if (WindowedValues.WindowedValueCoder.isMetadataSupported()) { - WindmillSink.decodeAdditionalMetadata(windowsCoder, message.getMetadata()); + BeamFnApi.Elements.ElementMetadata elementMetadata = + WindmillSink.decodeAdditionalMetadata(windowsCoder, message.getMetadata()); + drainingValueFromUpstream = + elementMetadata.hasDrain() + ? (elementMetadata.getDrain() + == BeamFnApi.Elements.DrainMode.Enum.DRAINING) + : false; } InputStream inputStream = message.getData().newInput(); ElemT value = valueCoder.decode(inputStream, Coder.Context.OUTER); // todo #33176 specify additional metadata in the future - return WindowedValues.of(value, timestamp, windows, paneInfo); + return WindowedValues.of( + value, timestamp, windows, paneInfo, null, null, drainingValueFromUpstream); } catch (IOException e) { throw new RuntimeException(e); } diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillTimerInternals.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillTimerInternals.java index ee73ac138f0e..f1e09db08e93 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillTimerInternals.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillTimerInternals.java @@ -301,7 +301,10 @@ static Timer timerDataToWindmillTimer( } public static TimerData windmillTimerToTimerData( - WindmillNamespacePrefix prefix, Timer timer, Coder windowCoder) { + WindmillNamespacePrefix prefix, + Timer timer, + Coder windowCoder, + boolean draining) { // The tag is a path-structure string but cheaper to parse than a proper URI. It follows // this pattern, where no component but the ID can contain a slash @@ -395,6 +398,7 @@ public static TimerData windmillTimerToTimerData( timestamp, outputTimestamp, timerTypeToTimeDomain(timer.getType())); + // todo add draining } private static boolean useNewTimerTagEncoding(TimerData timerData) { diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindowingWindmillReader.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindowingWindmillReader.java index d91a5412b917..f4a6eec61cbf 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindowingWindmillReader.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindowingWindmillReader.java @@ -119,7 +119,8 @@ public NativeReaderIterator>> iterator() throw final K key = keyCoder.decode(context.getSerializedKey().newInput(), Coder.Context.OUTER); final WorkItem workItem = context.getWorkItem(); KeyedWorkItem keyedWorkItem = - new WindmillKeyedWorkItem<>(key, workItem, windowCoder, windowsCoder, valueCoder); + new WindmillKeyedWorkItem<>( + key, workItem, windowCoder, windowsCoder, valueCoder, context.getDrainMode()); final boolean isEmptyWorkItem = (Iterables.isEmpty(keyedWorkItem.timersIterable()) && Iterables.isEmpty(keyedWorkItem.elementsIterable())); diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/Work.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/Work.java index 8b41a2d13219..43f355dd7ef8 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/Work.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/Work.java @@ -78,12 +78,14 @@ public final class Work implements RefreshableWork { private volatile TimedState currentState; private volatile boolean isFailed; private volatile String processingThreadName = ""; + private final boolean drainMode; private Work( WorkItem workItem, long serializedWorkItemSize, Watermarks watermarks, ProcessingContext processingContext, + boolean drainMode, Supplier clock) { this.shardedKey = ShardedKey.create(workItem.getKey(), workItem.getShardingKey()); this.workItem = workItem; @@ -91,6 +93,7 @@ private Work( this.processingContext = processingContext; this.watermarks = watermarks; this.clock = clock; + this.drainMode = drainMode; this.startTime = clock.get(); Preconditions.checkState(EMPTY_ENUM_MAP.isEmpty()); // Create by passing EMPTY_ENUM_MAP to avoid recreating @@ -110,8 +113,10 @@ public static Work create( long serializedWorkItemSize, Watermarks watermarks, ProcessingContext processingContext, + boolean drainMode, Supplier clock) { - return new Work(workItem, serializedWorkItemSize, watermarks, processingContext, clock); + return new Work( + workItem, serializedWorkItemSize, watermarks, processingContext, drainMode, clock); } public static ProcessingContext createProcessingContext( @@ -207,6 +212,10 @@ public State getState() { return currentState.state(); } + public boolean getDrainMode() { + return drainMode; + } + public void setState(State state) { Instant now = clock.get(); totalDurationPerState.compute( diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/harness/SingleSourceWorkerHarness.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/harness/SingleSourceWorkerHarness.java index 0de9d130b650..af7746d69028 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/harness/SingleSourceWorkerHarness.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/harness/SingleSourceWorkerHarness.java @@ -155,6 +155,7 @@ private void streamingEngineDispatchLoop( (computationId, inputDataWatermark, synchronizedProcessingTime, + drainMode, workItem, serializedWorkItemSize, getWorkStreamLatencies) -> @@ -178,6 +179,7 @@ private void streamingEngineDispatchLoop( getDataClient, workCommitter::commit, heartbeatSender), + drainMode, getWorkStreamLatencies); })); try { @@ -239,6 +241,7 @@ private void applianceDispatchLoop(Supplier getWorkFn) watermarks.setOutputDataWatermark(workItem.getOutputDataWatermark()).build(), Work.createProcessingContext( computationId, getDataClient, workCommitter::commit, heartbeatSender), + computationWork.getDrainMode(), /* getWorkStreamLatencies= */ ImmutableList.of()); } } diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GetWorkResponseChunkAssembler.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GetWorkResponseChunkAssembler.java index 0ebb4726d3a1..3608bd1ccacd 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GetWorkResponseChunkAssembler.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GetWorkResponseChunkAssembler.java @@ -124,7 +124,8 @@ private static ComputationMetadata fromProto( metadataProto.getComputationId(), WindmillTimeUtils.windmillToHarnessWatermark(metadataProto.getInputDataWatermark()), WindmillTimeUtils.windmillToHarnessWatermark( - metadataProto.getDependentRealtimeInputWatermark())); + metadataProto.getDependentRealtimeInputWatermark()), + metadataProto.getDrainMode()); } abstract String computationId(); @@ -132,6 +133,8 @@ private static ComputationMetadata fromProto( abstract @Nullable Instant inputDataWatermark(); abstract @Nullable Instant synchronizedProcessingTime(); + + abstract boolean drainMode(); } @AutoValue diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcDirectGetWorkStream.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcDirectGetWorkStream.java index 2712bf1bd33d..8eb4c51a2b49 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcDirectGetWorkStream.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcDirectGetWorkStream.java @@ -281,6 +281,7 @@ private void consumeAssembledWorkItem(AssembledWorkItem assembledWorkItem) { assembledWorkItem.bufferedSize(), createWatermarks(workItem, metadata), createProcessingContext(metadata.computationId()), + metadata.drainMode(), assembledWorkItem.latencyAttributions()); budgetTracker.recordBudgetReceived(assembledWorkItem.bufferedSize()); GetWorkBudget extension = budgetTracker.computeBudgetExtension(); diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcGetWorkStream.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcGetWorkStream.java index ae7ce85e13a8..58407ad8147f 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcGetWorkStream.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcGetWorkStream.java @@ -203,6 +203,7 @@ private void consumeAssembledWorkItem(AssembledWorkItem assembledWorkItem) { assembledWorkItem.computationMetadata().computationId(), assembledWorkItem.computationMetadata().inputDataWatermark(), assembledWorkItem.computationMetadata().synchronizedProcessingTime(), + assembledWorkItem.computationMetadata().drainMode(), assembledWorkItem.workItem(), assembledWorkItem.bufferedSize(), assembledWorkItem.latencyAttributions()); diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/WorkItemReceiver.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/WorkItemReceiver.java index e2f69585e48f..71e524a308af 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/WorkItemReceiver.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/WorkItemReceiver.java @@ -30,6 +30,7 @@ void receiveWork( String computation, @Nullable Instant inputDataWatermark, @Nullable Instant synchronizedProcessingTime, + boolean drainMode, Windmill.WorkItem workItem, long serializedWorkItemSize, ImmutableList getWorkStreamLatencies); diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/WorkItemScheduler.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/WorkItemScheduler.java index b9d31fbe501d..4121aa758ba7 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/WorkItemScheduler.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/WorkItemScheduler.java @@ -35,6 +35,7 @@ public interface WorkItemScheduler { * @param workItem {@link WorkItem} to be processed. * @param watermarks processing watermarks for the workItem. * @param processingContext for processing the workItem. + * @param drainMode is job is draining. * @param getWorkStreamLatencies Latencies per processing stage for the WorkItem for reporting * back to Streaming Engine backend. */ @@ -43,5 +44,6 @@ void scheduleWork( long serializedWorkItemSize, Watermarks watermarks, Work.ProcessingContext processingContext, + boolean drainMode, ImmutableList getWorkStreamLatencies); } 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 a4cd5d6d8a6b..0ca820cad8d8 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 @@ -210,10 +210,12 @@ public void scheduleWork( long serializedWorkItemSize, Watermarks watermarks, Work.ProcessingContext processingContext, + boolean drainMode, ImmutableList getWorkStreamLatencies) { computationState.activateWork( ExecutableWork.create( - Work.create(workItem, serializedWorkItemSize, watermarks, processingContext, clock), + Work.create( + workItem, serializedWorkItemSize, watermarks, processingContext, drainMode, clock), work -> processWork(computationState, work, getWorkStreamLatencies))); } diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/FakeWindmillServer.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/FakeWindmillServer.java index a5c8909b8d07..1605e5323b13 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/FakeWindmillServer.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/FakeWindmillServer.java @@ -275,6 +275,7 @@ public boolean awaitTermination(int time, TimeUnit unit) throws InterruptedExcep computationWork.getComputationId(), inputDataWatermark, Instant.now(), + false, workItem, workItem.getSerializedSize(), ImmutableList.of( 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 b21b8e830ae8..22253d661073 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 @@ -372,6 +372,7 @@ private static ExecutableWork createMockWork( Watermarks.builder().setInputDataWatermark(Instant.EPOCH).build(), Work.createProcessingContext( computationId, new FakeGetDataClient(), ignored -> {}, mock(HeartbeatSender.class)), + false, Instant::now), processWorkFn); } @@ -3552,6 +3553,7 @@ public void testLatencyAttributionProtobufsPopulated() { new FakeGetDataClient(), ignored -> {}, mock(HeartbeatSender.class)), + false, clock); clock.sleep(Duration.millis(10)); diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingGroupAlsoByWindowFnsTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingGroupAlsoByWindowFnsTest.java index 1182a2c0b9e9..f7852ec1767d 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingGroupAlsoByWindowFnsTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingGroupAlsoByWindowFnsTest.java @@ -194,7 +194,7 @@ private WindowedValue> createValue( return new ValueInEmptyWindows<>( (KeyedWorkItem) new WindmillKeyedWorkItem<>( - KEY, workItem.build(), windowCoder, wildcardWindowsCoder, valueCoder)); + KEY, workItem.build(), windowCoder, wildcardWindowsCoder, valueCoder, false)); } @Test diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingGroupAlsoByWindowsReshuffleDoFnTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingGroupAlsoByWindowsReshuffleDoFnTest.java index a348c0f00214..52c9844add86 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingGroupAlsoByWindowsReshuffleDoFnTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingGroupAlsoByWindowsReshuffleDoFnTest.java @@ -132,7 +132,7 @@ private WindowedValue> createValue( return new ValueInEmptyWindows<>( (KeyedWorkItem) new WindmillKeyedWorkItem<>( - KEY, workItem.build(), windowCoder, wildcardWindowsCoder, valueCoder)); + KEY, workItem.build(), windowCoder, wildcardWindowsCoder, valueCoder, false)); } @Test 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 e216f912d77f..93b279f0aec5 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 @@ -143,6 +143,7 @@ private static Work createMockWork(Windmill.WorkItem workItem, Watermarks waterm watermarks, Work.createProcessingContext( COMPUTATION_ID, new FakeGetDataClient(), ignored -> {}, mock(HeartbeatSender.class)), + false, Instant::now); } diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItemTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItemTest.java index 53a36722e41c..1e03f196d816 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItemTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItemTest.java @@ -93,7 +93,7 @@ public void testElementIteration() throws Exception { KeyedWorkItem keyedWorkItem = new WindmillKeyedWorkItem<>( - KEY, workItem.build(), WINDOW_CODER, WINDOWS_CODER, VALUE_CODER); + KEY, workItem.build(), WINDOW_CODER, WINDOWS_CODER, VALUE_CODER, false); assertThat( keyedWorkItem.elementsIterable(), @@ -148,7 +148,7 @@ public void testTimerOrdering() throws Exception { .build(); KeyedWorkItem keyedWorkItem = - new WindmillKeyedWorkItem<>(KEY, workItem, WINDOW_CODER, WINDOWS_CODER, VALUE_CODER); + new WindmillKeyedWorkItem<>(KEY, workItem, WINDOW_CODER, WINDOWS_CODER, VALUE_CODER, false); assertThat( keyedWorkItem.timersIterable(), diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WindmillTimerInternalsTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WindmillTimerInternalsTest.java index ec8672b6a75f..4780cd768efb 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WindmillTimerInternalsTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WindmillTimerInternalsTest.java @@ -96,7 +96,8 @@ public void testTimerDataToFromTimer() { WindmillTimerInternals.windmillTimerToTimerData( prefix, WindmillTimerInternals.timerDataToWindmillTimer(stateFamily, prefix, timer), - coder); + coder, + false); // The function itself bounds output, so we dont expect the original input as the // output, we expect it to be bounded TimerData expected = @@ -145,7 +146,8 @@ public void testTimerDataToFromTimer() { prefix, WindmillTimerInternals.timerDataToWindmillTimer( stateFamily, prefix, timer), - coder), + coder, + false), equalTo(expected)); } } 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 df3b959c82c5..334b9414b26b 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 @@ -204,6 +204,7 @@ private static Work createMockWork(Windmill.WorkItem workItem, Watermarks waterm watermarks, Work.createProcessingContext( COMPUTATION_ID, new FakeGetDataClient(), ignored -> {}, mock(HeartbeatSender.class)), + false, Instant::now); } @@ -1014,6 +1015,7 @@ public void testFailedWorkItemsAbort() throws Exception { new FakeGetDataClient(), ignored -> {}, mock(HeartbeatSender.class)), + false, Instant::now); context.start( "key", diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/ActiveWorkStateTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/ActiveWorkStateTest.java index c0cb8241d73e..865ae2612803 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/ActiveWorkStateTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/ActiveWorkStateTest.java @@ -71,6 +71,7 @@ private static ExecutableWork createWork(Windmill.WorkItem workItem) { workItem.getSerializedSize(), Watermarks.builder().setInputDataWatermark(Instant.EPOCH).build(), createWorkProcessingContext(), + false, Instant::now), ignored -> {}); } @@ -82,6 +83,7 @@ private static ExecutableWork expiredWork(Windmill.WorkItem workItem) { workItem.getSerializedSize(), Watermarks.builder().setInputDataWatermark(Instant.EPOCH).build(), createWorkProcessingContext(), + false, () -> Instant.EPOCH), ignored -> {}); } diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCacheTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCacheTest.java index 935b25acb6f2..1c8b8fca131d 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCacheTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCacheTest.java @@ -75,6 +75,7 @@ private static ExecutableWork createWork(ShardedKey shardedKey, long workToken, new FakeGetDataClient(), ignored -> {}, mock(HeartbeatSender.class)), + false, Instant::now), ignored -> {}); } diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/harness/FanOutStreamingEngineWorkerHarnessTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/harness/FanOutStreamingEngineWorkerHarnessTest.java index 65e40f171b0c..94c8f4b75957 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/harness/FanOutStreamingEngineWorkerHarnessTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/harness/FanOutStreamingEngineWorkerHarnessTest.java @@ -128,6 +128,7 @@ private static WorkItemScheduler noOpProcessWorkItemFn() { serializedWorkItemSize, watermarks, processingContext, + drainMode, getWorkStreamLatencies) -> {}; } diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/harness/WindmillStreamSenderTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/harness/WindmillStreamSenderTest.java index b94270ad7bb7..3217c736adb1 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/harness/WindmillStreamSenderTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/harness/WindmillStreamSenderTest.java @@ -68,6 +68,7 @@ public class WindmillStreamSenderTest { serializedWorkItemSize, watermarks, processingContext, + drainMode, getWorkStreamLatencies) -> {}; @Rule public transient Timeout globalTimeout = Timeout.seconds(600); private ManagedChannel inProcessChannel; diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/util/BoundedQueueExecutorTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/util/BoundedQueueExecutorTest.java index a86e6060955c..d7ea039bb809 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/util/BoundedQueueExecutorTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/util/BoundedQueueExecutorTest.java @@ -83,6 +83,7 @@ private static ExecutableWork createWork(Consumer executeWorkFn) { new FakeGetDataClient(), ignored -> {}, mock(HeartbeatSender.class)), + false, Instant::now), executeWorkFn); } diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingApplianceWorkCommitterTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingApplianceWorkCommitterTest.java index 477c764a70ef..5c3132ae471d 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingApplianceWorkCommitterTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingApplianceWorkCommitterTest.java @@ -74,6 +74,7 @@ private static Work createMockWork(long workToken) { throw new UnsupportedOperationException(); }, mock(HeartbeatSender.class)), + false, Instant::now); } diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingEngineWorkCommitterTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingEngineWorkCommitterTest.java index 5748b128f971..b4f63fa71618 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingEngineWorkCommitterTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingEngineWorkCommitterTest.java @@ -104,6 +104,7 @@ private static Work createMockWork(long workToken) { throw new UnsupportedOperationException(); }, mock(HeartbeatSender.class)), + false, Instant::now); } diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcDirectGetWorkStreamTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcDirectGetWorkStreamTest.java index 419000178381..76883bebdac0 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcDirectGetWorkStreamTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcDirectGetWorkStreamTest.java @@ -70,6 +70,7 @@ public class GrpcDirectGetWorkStreamTest { serializedWorkItemSize, watermarks, processingContext, + drainMode, getWorkStreamLatencies) -> {}; private static final Windmill.JobHeader TEST_JOB_HEADER = Windmill.JobHeader.newBuilder() @@ -283,6 +284,7 @@ public void testConsumedWorkItem_computesAndSendsCorrectExtension() throws Inter serializedWorkItemSize, watermarks, processingContext, + drainMode, getWorkStreamLatencies) -> { scheduledWorkItems.add(work); }); @@ -327,8 +329,12 @@ public void testConsumedWorkItem_doesNotSendExtensionIfOutstandingBudgetHigh() createGetWorkStream( testStub, initialBudget, - (work, serializedWorkItemSize, watermarks, processingContext, getWorkStreamLatencies) -> - scheduledWorkItems.add(work)); + (work, + serializedWorkItemSize, + watermarks, + processingContext, + drainMode, + getWorkStreamLatencies) -> scheduledWorkItems.add(work)); Windmill.WorkItem workItem = Windmill.WorkItem.newBuilder() .setKey(ByteString.copyFromUtf8("somewhat_long_key")) @@ -365,6 +371,7 @@ public void testConsumedWorkItems() throws InterruptedException { serializedWorkItemSize, watermarks, processingContext, + drainMode, getWorkStreamLatencies) -> { scheduledWorkItems.add(work); }); @@ -408,6 +415,7 @@ public void testConsumedWorkItems_itemsSplitAcrossResponses() throws Interrupted serializedWorkItemSize, watermarks, processingContext, + drainMode, getWorkStreamLatencies) -> { scheduledWorkItems.add(work); }); diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcWindmillServerTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcWindmillServerTest.java index e52b6e8de4bf..a0fe03a08e6c 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcWindmillServerTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcWindmillServerTest.java @@ -336,6 +336,7 @@ public void onCompleted() { (String computation, @Nullable Instant inputDataWatermark, Instant synchronizedProcessingTime, + boolean drainMode, WorkItem workItem, long serializedWorkItemSize, ImmutableList getWorkStreamLatencies) -> { @@ -469,6 +470,7 @@ public void onCompleted() { (String computation, @Nullable Instant inputDataWatermark, Instant synchronizedProcessingTime, + boolean drainMode, WorkItem workItem, long serializedWorkItemSize, ImmutableList getWorkStreamLatencies) -> { diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/failures/WorkFailureProcessorTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/failures/WorkFailureProcessorTest.java index f55549f7e2d9..41f2230f4a8f 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/failures/WorkFailureProcessorTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/failures/WorkFailureProcessorTest.java @@ -95,6 +95,7 @@ private static ExecutableWork createWork(Supplier clock, Consumer new FakeGetDataClient(), ignored -> {}, mock(HeartbeatSender.class)), + false, clock), processWorkFn); } diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/work/refresh/ActiveWorkRefresherTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/work/refresh/ActiveWorkRefresherTest.java index 115deccf6df4..e88209710022 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/work/refresh/ActiveWorkRefresherTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/work/refresh/ActiveWorkRefresherTest.java @@ -133,6 +133,7 @@ private ExecutableWork createOldWork( Watermarks.builder().setInputDataWatermark(Instant.EPOCH).build(), Work.createProcessingContext( "computationId", new FakeGetDataClient(), ignored -> {}, heartbeatSender), + false, A_LONG_TIME_AGO), processWork); } From 0f10903e73250f0feff942949bbb4d1d5905997c Mon Sep 17 00:00:00 2001 From: Radek Stankiewicz Date: Fri, 7 Nov 2025 19:55:24 +0100 Subject: [PATCH 2/3] review --- .../runners/dataflow/worker/UngroupedWindmillReader.java | 4 +--- .../beam/runners/dataflow/worker/WindmillKeyedWorkItem.java | 5 +---- .../beam/runners/dataflow/worker/FakeWindmillServer.java | 2 +- .../worker/windmill/client/grpc/GrpcWindmillServerTest.java | 4 +++- 4 files changed, 6 insertions(+), 9 deletions(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/UngroupedWindmillReader.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/UngroupedWindmillReader.java index 3295347249b0..7b1d4ae9629c 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/UngroupedWindmillReader.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/UngroupedWindmillReader.java @@ -123,9 +123,7 @@ protected WindowedValue decodeMessage(Windmill.Message message) throws IOExce BeamFnApi.Elements.ElementMetadata elementMetadata = WindmillSink.decodeAdditionalMetadata(windowsCoder, message.getMetadata()); drainingValueFromUpstream = - elementMetadata.hasDrain() - ? (elementMetadata.getDrain() == BeamFnApi.Elements.DrainMode.Enum.DRAINING) - : false; + elementMetadata.getDrain() == BeamFnApi.Elements.DrainMode.Enum.DRAINING; } if (valueCoder instanceof KvCoder) { KvCoder kvCoder = (KvCoder) valueCoder; diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItem.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItem.java index 5cc3b39df12f..aa79be1dc93e 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItem.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItem.java @@ -121,10 +121,7 @@ public Iterable> elementsIterable() { BeamFnApi.Elements.ElementMetadata elementMetadata = WindmillSink.decodeAdditionalMetadata(windowsCoder, message.getMetadata()); drainingValueFromUpstream = - elementMetadata.hasDrain() - ? (elementMetadata.getDrain() - == BeamFnApi.Elements.DrainMode.Enum.DRAINING) - : false; + elementMetadata.getDrain() == BeamFnApi.Elements.DrainMode.Enum.DRAINING; } InputStream inputStream = message.getData().newInput(); ElemT value = valueCoder.decode(inputStream, Coder.Context.OUTER); diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/FakeWindmillServer.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/FakeWindmillServer.java index 1605e5323b13..1c5f7504bf32 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/FakeWindmillServer.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/FakeWindmillServer.java @@ -275,7 +275,7 @@ public boolean awaitTermination(int time, TimeUnit unit) throws InterruptedExcep computationWork.getComputationId(), inputDataWatermark, Instant.now(), - false, + computationWork.getDrainMode(), workItem, workItem.getSerializedSize(), ImmutableList.of( diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcWindmillServerTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcWindmillServerTest.java index a0fe03a08e6c..d417d7d3417c 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcWindmillServerTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcWindmillServerTest.java @@ -413,7 +413,8 @@ public void onNext(StreamingGetWorkRequest request) { ComputationWorkItemMetadata.newBuilder() .setComputationId("comp") .setDependentRealtimeInputWatermark(17000) - .setInputDataWatermark(18000)); + .setInputDataWatermark(18000) + .setDrainMode(true)); int loopVariant = loop % 3; if (loopVariant < 1) { responseChunk.addSerializedWorkItem(serializedResponses.pop()); @@ -477,6 +478,7 @@ public void onCompleted() { assertEquals(inputDataWatermark, new Instant(18)); assertEquals(synchronizedProcessingTime, new Instant(17)); assertEquals(workItem.getKey(), ByteString.copyFromUtf8("somewhat_long_key")); + assertTrue(drainMode); assertTrue(sentResponseIds.containsKey(workItem.getWorkToken())); sentResponseIds.remove(workItem.getWorkToken()); latch.countDown(); From 64c3dd8b4fbd10eaea1cfa10076aa5855f2b8b68 Mon Sep 17 00:00:00 2001 From: Radek Stankiewicz Date: Mon, 24 Nov 2025 12:03:45 +0100 Subject: [PATCH 3/3] add comments and relevant tests --- .../worker/UngroupedWindmillReader.java | 6 +- .../worker/WindmillKeyedWorkItem.java | 7 +- .../worker/WindmillTimerInternals.java | 3 +- .../worker/WindmillKeyedWorkItemTest.java | 68 +++++++++++++++++++ 4 files changed, 79 insertions(+), 5 deletions(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/UngroupedWindmillReader.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/UngroupedWindmillReader.java index 7b1d4ae9629c..c248259a12de 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/UngroupedWindmillReader.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/UngroupedWindmillReader.java @@ -118,6 +118,10 @@ protected WindowedValue decodeMessage(Windmill.Message message) throws IOExce Collection windows = WindmillSink.decodeMetadataWindows(windowsCoder, message.getMetadata()); PaneInfo paneInfo = WindmillSink.decodeMetadataPane(message.getMetadata()); + /** + * https://s.apache.org/beam-drain-mode - propagate drain bit if aggregation/expiry induced by + * drain happened upstream + */ boolean drainingValueFromUpstream = false; if (WindowedValues.WindowedValueCoder.isMetadataSupported()) { BeamFnApi.Elements.ElementMetadata elementMetadata = @@ -133,12 +137,10 @@ protected WindowedValue decodeMessage(Windmill.Message message) throws IOExce @SuppressWarnings("unchecked") T result = (T) KV.of(decode(kvCoder.getKeyCoder(), key), decode(kvCoder.getValueCoder(), data)); - // todo #33176 propagate metadata to windowed value return WindowedValues.of( result, timestampMillis, windows, paneInfo, null, null, drainingValueFromUpstream); } else { notifyElementRead(data.available() + metadata.available()); - // todo #33176 propagate metadata to windowed value return WindowedValues.of( decode(valueCoder, data), timestampMillis, diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItem.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItem.java index aa79be1dc93e..415dab526bb5 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItem.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItem.java @@ -61,6 +61,7 @@ public class WindmillKeyedWorkItem implements KeyedWorkItem private final Windmill.WorkItem workItem; private final K key; + // used to inform that timer was caused by drain private final boolean drainMode; private final transient Coder windowCoder; @@ -115,7 +116,10 @@ public Iterable> elementsIterable() { Collection windows = WindmillSink.decodeMetadataWindows(windowsCoder, message.getMetadata()); PaneInfo paneInfo = WindmillSink.decodeMetadataPane(message.getMetadata()); - // Draining value is based on upstream data + /** + * https://s.apache.org/beam-drain-mode - propagate drain bit if aggregation/expiry + * induced by drain happened upstream + */ boolean drainingValueFromUpstream = false; if (WindowedValues.WindowedValueCoder.isMetadataSupported()) { BeamFnApi.Elements.ElementMetadata elementMetadata = @@ -125,7 +129,6 @@ public Iterable> elementsIterable() { } InputStream inputStream = message.getData().newInput(); ElemT value = valueCoder.decode(inputStream, Coder.Context.OUTER); - // todo #33176 specify additional metadata in the future return WindowedValues.of( value, timestamp, windows, paneInfo, null, null, drainingValueFromUpstream); } catch (IOException e) { diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillTimerInternals.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillTimerInternals.java index f1e09db08e93..8dac9d11715e 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillTimerInternals.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillTimerInternals.java @@ -398,7 +398,8 @@ public static TimerData windmillTimerToTimerData( timestamp, outputTimestamp, timerTypeToTimeDomain(timer.getType())); - // todo add draining + // todo add draining (https://github.com/apache/beam/issues/36884) + } private static boolean useNewTimerTagEncoding(TimerData timerData) { diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItemTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItemTest.java index 1e03f196d816..bbdde4498605 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItemTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItemTest.java @@ -22,6 +22,7 @@ import java.io.IOException; import java.util.Collection; import java.util.Collections; +import java.util.Iterator; import org.apache.beam.model.fnexecution.v1.BeamFnApi; import org.apache.beam.runners.core.KeyedWorkItem; import org.apache.beam.runners.core.StateNamespace; @@ -40,10 +41,12 @@ import org.apache.beam.sdk.transforms.windowing.IntervalWindow; import org.apache.beam.sdk.transforms.windowing.PaneInfo; import org.apache.beam.sdk.transforms.windowing.PaneInfo.Timing; +import org.apache.beam.sdk.values.WindowedValue; import org.apache.beam.sdk.values.WindowedValues; import org.apache.beam.vendor.grpc.v1p69p0.com.google.protobuf.ByteString; import org.hamcrest.Matchers; import org.joda.time.Instant; +import org.junit.Assert; import org.junit.Before; import org.junit.Test; import org.junit.runner.RunWith; @@ -123,6 +126,24 @@ private void addElement( .setMetadata(encodedMetadata); } + private void addElementWithMetadata( + Windmill.InputMessageBundle.Builder chunk, + long timestamp, + String value, + IntervalWindow window, + PaneInfo pane, + BeamFnApi.Elements.ElementMetadata metadata) + throws IOException { + ByteString encodedMetadata = + WindmillSink.encodeMetadata( + WINDOWS_CODER, Collections.singletonList(window), pane, metadata); + chunk + .addMessagesBuilder() + .setTimestamp(WindmillTimeUtils.harnessToWindmillTimestamp(new Instant(timestamp))) + .setData(ByteString.copyFromUtf8(value)) + .setMetadata(encodedMetadata); + } + private PaneInfo paneInfo(int index) { return PaneInfo.createPane(false, false, Timing.EARLY, index, -1); } @@ -186,4 +207,51 @@ public void testCoderIsSerializableWithWellKnownCoderType() { FakeKeyedWorkItemCoder.of( KvCoder.of(GlobalWindow.Coder.INSTANCE, GlobalWindow.Coder.INSTANCE))); } + + @Test + public void testDrainPropagated() throws Exception { + WindowedValues.FullWindowedValueCoder.setMetadataSupported(); + Windmill.WorkItem.Builder workItem = + Windmill.WorkItem.newBuilder() + .setKey(SERIALIZED_KEY) + .setTimers( + Windmill.TimerBundle.newBuilder() + .addTimers( + makeSerializedTimer(STATE_NAMESPACE_2, 3, Windmill.Timer.Type.WATERMARK)) + .build()) + .setWorkToken(17); + Windmill.InputMessageBundle.Builder chunk1 = workItem.addMessageBundlesBuilder(); + chunk1.setSourceComputationId("computation"); + addElementWithMetadata( + chunk1, + 5, + "hello", + WINDOW_1, + paneInfo(0), + BeamFnApi.Elements.ElementMetadata.newBuilder() + .setDrain(BeamFnApi.Elements.DrainMode.Enum.DRAINING) + .build()); + addElementWithMetadata( + chunk1, + 7, + "world", + WINDOW_2, + paneInfo(2), + BeamFnApi.Elements.ElementMetadata.newBuilder() + .setDrain(BeamFnApi.Elements.DrainMode.Enum.NOT_DRAINING) + .build()); + KeyedWorkItem keyedWorkItem = + new WindmillKeyedWorkItem<>( + KEY, workItem.build(), WINDOW_CODER, WINDOWS_CODER, VALUE_CODER, true); + + Iterator> iterator = keyedWorkItem.elementsIterable().iterator(); + Assert.assertTrue(iterator.next().causedByDrain()); + Assert.assertFalse(iterator.next().causedByDrain()); + + // todo add assert for draining once timerdata is filled + // (https://github.com/apache/beam/issues/36884) + assertThat( + keyedWorkItem.timersIterable(), + Matchers.contains(makeTimer(STATE_NAMESPACE_2, 3, TimeDomain.EVENT_TIME))); + } }