From e66be4167032041fb01ea014226838ee88002371 Mon Sep 17 00:00:00 2001 From: Shunping Huang Date: Fri, 24 Jul 2026 22:44:36 -0400 Subject: [PATCH 01/11] Fix DataflowOutputCounter calculation for ValueInEmptyWindows When processing shuffle or streaming data in Dataflow Legacy Runner (e.g., from GroupingShuffleReader or WindowingWindmillReader), KeyedWorkItems are wrapped inside a ValueInEmptyWindows (windows.size() == 0). Previously, DataflowOutputCounter.update() counted these as 1 element. This caused inaccurate element counts because: 1. A KeyedWorkItem can contain multiple elements. 2. Elements may belong to multiple windows and need to be fanned out accordingly. 3. KeyedWorkItems containing only timers were incorrectly incrementing element counters. --- .../worker/DataflowOutputCounter.java | 21 +++++++++++++++---- 1 file changed, 17 insertions(+), 4 deletions(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounter.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounter.java index 7c5859a9d324..9c84f880357d 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounter.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounter.java @@ -18,6 +18,7 @@ package org.apache.beam.runners.dataflow.worker; import org.apache.beam.runners.core.ElementByteSizeObservable; +import org.apache.beam.runners.core.KeyedWorkItem; import org.apache.beam.runners.dataflow.worker.counters.Counter; import org.apache.beam.runners.dataflow.worker.counters.CounterFactory; import org.apache.beam.runners.dataflow.worker.counters.CounterName; @@ -63,11 +64,23 @@ public void update(Object elem) throws Exception { objectAndByteCounter.update(elem); long windowsSize = ((WindowedValue) elem).getWindows().size(); if (windowsSize == 0) { - // GroupingShuffleReader produces ValueInEmptyWindows. - // For now, we count the element at least once to keep the current counter - // behavior. - elementCount.addValue(1L); + // ValueInEmptyWindows occurs when processing shuffle/streaming work items + // (e.g. GroupingShuffleReader or WindowingWindmillReader). KeyedWorkItems contain elements + // and timers across multiple windows, so the wrapper ValueInEmptyWindows has 0 windows. + Object value = ((WindowedValue) elem).getValue(); + if (value instanceof KeyedWorkItem) { + KeyedWorkItem keyedWorkItem = (KeyedWorkItem) value; + long totalElementCount = 0; + // Iterate only through elementsIterable and ignore timers in KeyedWorkItem. + for (WindowedValue element : keyedWorkItem.elementsIterable()) { + long elementWindowsSize = element.getWindows().size(); + // Fan out for windows. + totalElementCount += (elementWindowsSize == 0 ? 1L : elementWindowsSize); + } + elementCount.addValue(totalElementCount); + } } else { + // Standard WindowedValue. elementCount.addValue(windowsSize); } } From 420ced52ad7c886e7546d903f5c49e7d797b2cc8 Mon Sep 17 00:00:00 2001 From: Shunping Huang Date: Fri, 24 Jul 2026 23:38:47 -0400 Subject: [PATCH 02/11] Address non keyedworkitems --- .../runners/dataflow/worker/DataflowOutputCounter.java | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounter.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounter.java index 9c84f880357d..7edebc699ecd 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounter.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounter.java @@ -65,10 +65,10 @@ public void update(Object elem) throws Exception { long windowsSize = ((WindowedValue) elem).getWindows().size(); if (windowsSize == 0) { // ValueInEmptyWindows occurs when processing shuffle/streaming work items - // (e.g. GroupingShuffleReader or WindowingWindmillReader). KeyedWorkItems contain elements - // and timers across multiple windows, so the wrapper ValueInEmptyWindows has 0 windows. Object value = ((WindowedValue) elem).getValue(); if (value instanceof KeyedWorkItem) { + // KeyedWorkItem wrapped in ValueInEmptyWindows + // (e.g. WindowingWindmillReader for Streaming GBK) KeyedWorkItem keyedWorkItem = (KeyedWorkItem) value; long totalElementCount = 0; // Iterate only through elementsIterable and ignore timers in KeyedWorkItem. @@ -78,6 +78,10 @@ public void update(Object elem) throws Exception { totalElementCount += (elementWindowsSize == 0 ? 1L : elementWindowsSize); } elementCount.addValue(totalElementCount); + } else { + // Non-KeyedWorkItem wrapped in ValueInEmptyWindows + // (e.g. GroupingShuffleReader KV output for Batch GBK) + elementCount.addValue(1L); } } else { // Standard WindowedValue. From 16177c854f72a7554d304551a9d0d8fbd29580d8 Mon Sep 17 00:00:00 2001 From: Shunping Huang Date: Sat, 25 Jul 2026 00:25:03 -0400 Subject: [PATCH 03/11] Spotless --- .../beam/runners/dataflow/worker/DataflowOutputCounter.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounter.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounter.java index 7edebc699ecd..add1b72a0c90 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounter.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounter.java @@ -67,7 +67,7 @@ public void update(Object elem) throws Exception { // ValueInEmptyWindows occurs when processing shuffle/streaming work items Object value = ((WindowedValue) elem).getValue(); if (value instanceof KeyedWorkItem) { - // KeyedWorkItem wrapped in ValueInEmptyWindows + // KeyedWorkItem wrapped in ValueInEmptyWindows // (e.g. WindowingWindmillReader for Streaming GBK) KeyedWorkItem keyedWorkItem = (KeyedWorkItem) value; long totalElementCount = 0; @@ -79,7 +79,7 @@ public void update(Object elem) throws Exception { } elementCount.addValue(totalElementCount); } else { - // Non-KeyedWorkItem wrapped in ValueInEmptyWindows + // Non-KeyedWorkItem wrapped in ValueInEmptyWindows // (e.g. GroupingShuffleReader KV output for Batch GBK) elementCount.addValue(1L); } From a620b1c6d59fb1581cde1c97665d04190d36e8b7 Mon Sep 17 00:00:00 2001 From: Shunping Huang Date: Wed, 29 Jul 2026 10:02:36 -0400 Subject: [PATCH 04/11] Add elementWindowsIterable to only decode window metadata and use it in DataflowOutputCounter --- .../beam/runners/core/KeyedWorkItem.java | 9 ++++ .../worker/DataflowOutputCounter.java | 2 +- .../worker/WindmillKeyedWorkItem.java | 50 +++++++++++++++++++ 3 files changed, 60 insertions(+), 1 deletion(-) diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/KeyedWorkItem.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/KeyedWorkItem.java index 4901c5cbed5b..cd7b6c10e500 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/KeyedWorkItem.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/KeyedWorkItem.java @@ -35,4 +35,13 @@ public interface KeyedWorkItem { /** Returns an iterable containing the elements. */ Iterable> elementsIterable(); + + /** + * Returns an iterable containing windowed values without guaranteeing element payload decoding. + * Useful for lightweight inspection of windowing metadata without payload deserialization + * overhead. + */ + default Iterable> elementWindowsIterable() { + return elementsIterable(); + } } diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounter.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounter.java index add1b72a0c90..0f8858a5fb31 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounter.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounter.java @@ -72,7 +72,7 @@ public void update(Object elem) throws Exception { KeyedWorkItem keyedWorkItem = (KeyedWorkItem) value; long totalElementCount = 0; // Iterate only through elementsIterable and ignore timers in KeyedWorkItem. - for (WindowedValue element : keyedWorkItem.elementsIterable()) { + for (WindowedValue element : keyedWorkItem.elementWindowsIterable()) { long elementWindowsSize = element.getWindows().size(); // Fan out for windows. totalElementCount += (elementWindowsSize == 0 ? 1L : elementWindowsSize); 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 c4c0b6ed92d3..1df2fc5eb43b 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 @@ -184,6 +184,56 @@ public Iterable timersIterable() { } } + @SuppressWarnings("nullness") + private @Nullable WindowedValue parseElemWindowOnly(Windmill.Message message) { + try { + Instant timestamp = WindmillTimeUtils.windmillToHarnessTimestamp(message.getTimestamp()); + Collection windows = + WindmillSink.decodeMetadataWindows(windowsCoder, message.getMetadata()); + PaneInfo paneInfo = WindmillSink.decodeMetadataPane(message.getMetadata()); + CausedByDrain drainingValueFromUpstream = CausedByDrain.NORMAL; + ValueKind valueKind = ValueKind.INSERT; + if (WindowedValues.WindowedValueCoder.isMetadataSupported()) { + BeamFnApi.Elements.ElementMetadata elementMetadata = + WindmillSink.decodeAdditionalMetadata(windowsCoder, message.getMetadata()); + drainingValueFromUpstream = + elementMetadata.getDrain() == BeamFnApi.Elements.DrainMode.Enum.DRAINING + ? CausedByDrain.CAUSED_BY_DRAIN + : CausedByDrain.NORMAL; + valueKind = WindmillValueKindHelper.fromProto(elementMetadata.getValueKind()); + } + return WindowedValues.of( + (ElemT) null, + timestamp, + windows, + paneInfo, + null, + null, + drainingValueFromUpstream, + null, + valueKind); + } catch (RuntimeException | IOException e) { + if (!skipUndecodableElements) { + throw new RuntimeException(e); + } + LOG.error( + "Skipping input element for work token {} on sharding key {} due to decoding error", + workItem.getWorkToken(), + workItem.getShardingKey(), + e); + return null; + } + } + + @Override + @SuppressWarnings("nullness") + public Iterable> elementWindowsIterable() { + return FluentIterable.from(workItem.getMessageBundlesList()) + .transformAndConcat(Windmill.InputMessageBundle::getMessagesList) + .transform(this::parseElemWindowOnly) + .filter(Objects::nonNull); + } + @Override @SuppressWarnings("nullness") public Iterable> elementsIterable() { From c0aa622a0315a4984db2e67221d4c4774ecbe678 Mon Sep 17 00:00:00 2001 From: Shunping Huang Date: Wed, 29 Jul 2026 10:23:27 -0400 Subject: [PATCH 05/11] Use the elementWindowsIterable in ReduceFnRunner --- .../GroupAlsoByWindowViaWindowSetNewDoFn.java | 2 +- .../apache/beam/runners/core/ReduceFnRunner.java | 15 +++++++++++++-- .../StreamingGroupAlsoByWindowViaWindowSetFn.java | 2 +- 3 files changed, 15 insertions(+), 4 deletions(-) diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/GroupAlsoByWindowViaWindowSetNewDoFn.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/GroupAlsoByWindowViaWindowSetNewDoFn.java index f242c7d10003..349e109930f0 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/GroupAlsoByWindowViaWindowSetNewDoFn.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/GroupAlsoByWindowViaWindowSetNewDoFn.java @@ -109,7 +109,7 @@ public void processElement(ProcessContext c) throws Exception { reduceFn, c.getPipelineOptions()); - reduceFnRunner.processElements(keyedWorkItem.elementsIterable()); + reduceFnRunner.processElements(keyedWorkItem); reduceFnRunner.onTimers(keyedWorkItem.timersIterable()); reduceFnRunner.persist(); } diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java index 7fe3b711aa0a..52faa2f815b9 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java @@ -361,13 +361,24 @@ private Collection windowsThatShouldFire(Set windows) throws Exception { * setting holds, and invoking {@link ReduceFn#onTrigger}. * */ + public void processElements(KeyedWorkItem keyedWorkItem) throws Exception { + processElementsInternal( + keyedWorkItem.elementWindowsIterable(), keyedWorkItem.elementsIterable()); + } + public void processElements(Iterable> values) throws Exception { - if (!values.iterator().hasNext()) { + processElementsInternal(values, values); + } + + private void processElementsInternal( + Iterable> elementWindows, Iterable> values) + throws Exception { + if (!elementWindows.iterator().hasNext()) { return; } // Determine all the windows for elements. - Set windows = collectWindows(values); + Set windows = collectWindows(elementWindows); // If an incoming element introduces a new window, attempt to merge it into an existing // window eagerly. Map windowToMergeResult = mergeWindows(windows); diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingGroupAlsoByWindowViaWindowSetFn.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingGroupAlsoByWindowViaWindowSetFn.java index ec36644d1e68..a183df19b6e7 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingGroupAlsoByWindowViaWindowSetFn.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingGroupAlsoByWindowViaWindowSetFn.java @@ -93,7 +93,7 @@ public void processElement( reduceFn, options); - reduceFnRunner.processElements(keyedWorkItem.elementsIterable()); + reduceFnRunner.processElements(keyedWorkItem); reduceFnRunner.onTimers(keyedWorkItem.timersIterable()); reduceFnRunner.persist(); } From c1fb2d6235fcc5c13548fe4acccfcc12be4484d0 Mon Sep 17 00:00:00 2001 From: Shunping Huang Date: Wed, 29 Jul 2026 10:39:11 -0400 Subject: [PATCH 06/11] Refactor: separate batch and streaming DataflowOutputCounter implementations --- .../worker/BatchDataflowOutputCounter.java | 53 +++++++++ .../worker/DataflowOutputCounter.java | 48 ++++---- .../IntrinsicMapTaskExecutorFactory.java | 14 ++- .../dataflow/worker/SimpleParDoFnHelpers.java | 5 +- .../StreamingDataflowOutputCounter.java | 67 +++++++++++ .../worker/DataflowOutputCounterTest.java | 108 ++++++++++++++++++ 6 files changed, 269 insertions(+), 26 deletions(-) create mode 100644 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/BatchDataflowOutputCounter.java create mode 100644 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowOutputCounter.java create mode 100644 runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounterTest.java diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/BatchDataflowOutputCounter.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/BatchDataflowOutputCounter.java new file mode 100644 index 000000000000..b512794fdfed --- /dev/null +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/BatchDataflowOutputCounter.java @@ -0,0 +1,53 @@ +/* + * 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. + */ +package org.apache.beam.runners.dataflow.worker; + +import org.apache.beam.runners.core.ElementByteSizeObservable; +import org.apache.beam.runners.dataflow.worker.counters.CounterFactory; +import org.apache.beam.runners.dataflow.worker.counters.NameContext; + +/** + * A Dataflow output counter specific to Batch pipelines. In batch pipelines, empty-window elements + * (e.g. GroupingShuffleReader emitting KV) represent a single PCollection + * element output. + */ +@SuppressWarnings({ + "nullness" // TODO(https://github.com/apache/beam/issues/20497) +}) +public class BatchDataflowOutputCounter extends DataflowOutputCounter { + + public BatchDataflowOutputCounter( + String outputName, CounterFactory counterFactory, NameContext nameContext) { + super(outputName, counterFactory, nameContext); + } + + public BatchDataflowOutputCounter( + String outputName, + ElementByteSizeObservable elementByteSizeObservable, + CounterFactory counterFactory, + NameContext nameContext) { + super(outputName, elementByteSizeObservable, counterFactory, nameContext); + } + + @Override + protected void updateEmptyWindows(Object elem) { + // Non-KeyedWorkItem wrapped in ValueInEmptyWindows (e.g. GroupingShuffleReader KV output for + // Batch GBK) + elementCount.addValue(1L); + } +} diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounter.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounter.java index 0f8858a5fb31..04b4c9b4f32d 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounter.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounter.java @@ -18,7 +18,6 @@ package org.apache.beam.runners.dataflow.worker; import org.apache.beam.runners.core.ElementByteSizeObservable; -import org.apache.beam.runners.core.KeyedWorkItem; import org.apache.beam.runners.dataflow.worker.counters.Counter; import org.apache.beam.runners.dataflow.worker.counters.CounterFactory; import org.apache.beam.runners.dataflow.worker.counters.CounterName; @@ -41,7 +40,28 @@ public class DataflowOutputCounter implements ElementCounter { private static final String MEAN_BYTE_COUNTER_NAME = "-MeanByteCount"; private OutputObjectAndByteCounter objectAndByteCounter; - private Counter elementCount; + protected Counter elementCount; + + public static DataflowOutputCounter create( + String outputName, + ElementByteSizeObservable elementByteSizeObservable, + CounterFactory counterFactory, + NameContext nameContext, + boolean isStreaming) { + return isStreaming + ? new StreamingDataflowOutputCounter( + outputName, elementByteSizeObservable, counterFactory, nameContext) + : new BatchDataflowOutputCounter( + outputName, elementByteSizeObservable, counterFactory, nameContext); + } + + public static DataflowOutputCounter create( + String outputName, + CounterFactory counterFactory, + NameContext nameContext, + boolean isStreaming) { + return create(outputName, null, counterFactory, nameContext, isStreaming); + } public DataflowOutputCounter( String outputName, CounterFactory counterFactory, NameContext nameContext) { @@ -64,31 +84,17 @@ public void update(Object elem) throws Exception { objectAndByteCounter.update(elem); long windowsSize = ((WindowedValue) elem).getWindows().size(); if (windowsSize == 0) { - // ValueInEmptyWindows occurs when processing shuffle/streaming work items - Object value = ((WindowedValue) elem).getValue(); - if (value instanceof KeyedWorkItem) { - // KeyedWorkItem wrapped in ValueInEmptyWindows - // (e.g. WindowingWindmillReader for Streaming GBK) - KeyedWorkItem keyedWorkItem = (KeyedWorkItem) value; - long totalElementCount = 0; - // Iterate only through elementsIterable and ignore timers in KeyedWorkItem. - for (WindowedValue element : keyedWorkItem.elementWindowsIterable()) { - long elementWindowsSize = element.getWindows().size(); - // Fan out for windows. - totalElementCount += (elementWindowsSize == 0 ? 1L : elementWindowsSize); - } - elementCount.addValue(totalElementCount); - } else { - // Non-KeyedWorkItem wrapped in ValueInEmptyWindows - // (e.g. GroupingShuffleReader KV output for Batch GBK) - elementCount.addValue(1L); - } + updateEmptyWindows(elem); } else { // Standard WindowedValue. elementCount.addValue(windowsSize); } } + protected void updateEmptyWindows(Object elem) throws Exception { + elementCount.addValue(1L); + } + @Override public void finishLazyUpdate(Object elem) { objectAndByteCounter.finishLazyUpdate(elem); diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/IntrinsicMapTaskExecutorFactory.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/IntrinsicMapTaskExecutorFactory.java index d3f2aacc74d0..56cc785cb28e 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/IntrinsicMapTaskExecutorFactory.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/IntrinsicMapTaskExecutorFactory.java @@ -63,6 +63,7 @@ import org.apache.beam.sdk.coders.KvCoder; import org.apache.beam.sdk.fn.IdGenerator; import org.apache.beam.sdk.options.PipelineOptions; +import org.apache.beam.sdk.options.StreamingOptions; import org.apache.beam.sdk.util.common.ElementByteSizeObserver; import org.apache.beam.sdk.values.TupleTag; import org.apache.beam.sdk.values.WindowedValues.WindowedValueCoder; @@ -102,8 +103,9 @@ public DataflowMapTaskExecutor create( IdGenerator idGenerator) { // Swap out all the InstructionOutput nodes with OutputReceiver nodes + boolean isStreaming = options.as(StreamingOptions.class).isStreaming(); Networks.replaceDirectedNetworkNodes( - network, createOutputReceiversTransform(stageName, counterSet)); + network, createOutputReceiversTransform(stageName, counterSet, isStreaming)); // Swap out all the ParallelInstruction nodes with Operation nodes. While updating the network, // we keep track of @@ -346,6 +348,11 @@ OperationNode createFlattenOperation( */ static Function createOutputReceiversTransform( final String stageName, final CounterFactory counterFactory) { + return createOutputReceiversTransform(stageName, counterFactory, false); + } + + static Function createOutputReceiversTransform( + final String stageName, final CounterFactory counterFactory, final boolean isStreaming) { return new TypeSafeNodeFunction(InstructionOutputNode.class) { @Override public Node typedApply(InstructionOutputNode input) { @@ -355,7 +362,7 @@ public Node typedApply(InstructionOutputNode input) { CloudObjects.coderFromCloudObject(CloudObject.fromSpec(cloudOutput.getCodec())); ElementCounter outputCounter = - new DataflowOutputCounter( + DataflowOutputCounter.create( cloudOutput.getName(), new ElementByteSizeObservableCoder<>(coder), counterFactory, @@ -363,7 +370,8 @@ public Node typedApply(InstructionOutputNode input) { stageName, cloudOutput.getOriginalName(), cloudOutput.getSystemName(), - cloudOutput.getName())); + cloudOutput.getName()), + isStreaming); outputReceiver.addOutputCounter(outputCounter); return OutputReceiverNode.create(outputReceiver, coder, input.getPcollectionId()); diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/SimpleParDoFnHelpers.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/SimpleParDoFnHelpers.java index 964cf2323d51..15bfba9bbc49 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/SimpleParDoFnHelpers.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/SimpleParDoFnHelpers.java @@ -190,9 +190,10 @@ public void output(TupleTag tag, WindowedValue output) { // doesn't today.) OutputReceiver undeclaredReceiver = new OutputReceiver(); + boolean isStreaming = options.as(StreamingOptions.class).isStreaming(); ElementCounter outputCounter = - new DataflowOutputCounter( - outputName, counterFactory, stepContext.getNameContext()); + DataflowOutputCounter.create( + outputName, counterFactory, stepContext.getNameContext(), isStreaming); undeclaredReceiver.addOutputCounter(outputCounter); undeclaredOutputs.put(tag, undeclaredReceiver); receiver = undeclaredReceiver; diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowOutputCounter.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowOutputCounter.java new file mode 100644 index 000000000000..1acd2917c32d --- /dev/null +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowOutputCounter.java @@ -0,0 +1,67 @@ +/* + * 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. + */ +package org.apache.beam.runners.dataflow.worker; + +import org.apache.beam.runners.core.ElementByteSizeObservable; +import org.apache.beam.runners.core.KeyedWorkItem; +import org.apache.beam.runners.dataflow.worker.counters.CounterFactory; +import org.apache.beam.runners.dataflow.worker.counters.NameContext; +import org.apache.beam.sdk.values.WindowedValue; + +/** + * A Dataflow output counter specific to Streaming pipelines. Unpacks {@link KeyedWorkItem}s in + * empty windows (e.g. emitted by WindowingWindmillReader) and counts element windows (ignoring + * timers) using lightweight metadata-only iteration via {@link + * KeyedWorkItem#elementWindowsIterable()}. + */ +@SuppressWarnings({ + "nullness" // TODO(https://github.com/apache/beam/issues/20497) +}) +public class StreamingDataflowOutputCounter extends DataflowOutputCounter { + + public StreamingDataflowOutputCounter( + String outputName, CounterFactory counterFactory, NameContext nameContext) { + super(outputName, counterFactory, nameContext); + } + + public StreamingDataflowOutputCounter( + String outputName, + ElementByteSizeObservable elementByteSizeObservable, + CounterFactory counterFactory, + NameContext nameContext) { + super(outputName, elementByteSizeObservable, counterFactory, nameContext); + } + + @Override + protected void updateEmptyWindows(Object elem) { + Object value = ((WindowedValue) elem).getValue(); + if (value instanceof KeyedWorkItem) { + KeyedWorkItem keyedWorkItem = (KeyedWorkItem) value; + long totalElementCount = 0; + // Iterate only through elementWindowsIterable and ignore timers in KeyedWorkItem. + for (WindowedValue element : keyedWorkItem.elementWindowsIterable()) { + long elementWindowsSize = element.getWindows().size(); + // Fan out for windows. + totalElementCount += (elementWindowsSize == 0 ? 1L : elementWindowsSize); + } + elementCount.addValue(totalElementCount); + } else { + elementCount.addValue(1L); + } + } +} diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounterTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounterTest.java new file mode 100644 index 000000000000..c32964c58587 --- /dev/null +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounterTest.java @@ -0,0 +1,108 @@ +/* + * 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. + */ +package org.apache.beam.runners.dataflow.worker; + +import static org.junit.Assert.assertEquals; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +import java.util.Collections; +import org.apache.beam.runners.core.KeyedWorkItem; +import org.apache.beam.runners.dataflow.worker.counters.CounterName; +import org.apache.beam.runners.dataflow.worker.counters.CounterSet; +import org.apache.beam.runners.dataflow.worker.counters.NameContext; +import org.apache.beam.runners.dataflow.worker.util.ValueInEmptyWindows; +import org.apache.beam.sdk.values.KV; +import org.apache.beam.sdk.values.WindowedValue; +import org.apache.beam.sdk.values.WindowedValues; +import org.junit.Before; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.junit.runners.JUnit4; + +/** Tests for {@link BatchDataflowOutputCounter} and {@link StreamingDataflowOutputCounter}. */ +@RunWith(JUnit4.class) +public class DataflowOutputCounterTest { + private static final String OUTPUT_NAME = "test_output"; + private CounterSet counterSet; + private NameContext nameContext; + + @Before + public void setUp() { + counterSet = new CounterSet(); + nameContext = NameContext.create("stage", "original", "system", OUTPUT_NAME); + } + + @Test + public void testBatchOutputCounterWithEmptyWindows() throws Exception { + DataflowOutputCounter batchCounter = + DataflowOutputCounter.create(OUTPUT_NAME, counterSet, nameContext, false); + + // Non-KeyedWorkItem value in empty windows (e.g. GroupingShuffleReader output) + ValueInEmptyWindows> shuffleValue = + new ValueInEmptyWindows<>(KV.of("key", "value")); + batchCounter.update(shuffleValue); + + long elementCount = + (Long) + counterSet + .getExistingCounter( + CounterName.named(DataflowOutputCounter.getElementCounterName(OUTPUT_NAME))) + .getAggregate(); + assertEquals(1L, elementCount); + } + + @Test + public void testStreamingOutputCounterWithKeyedWorkItem() throws Exception { + DataflowOutputCounter streamingCounter = + DataflowOutputCounter.create(OUTPUT_NAME, counterSet, nameContext, true); + + KeyedWorkItem kwi = mock(KeyedWorkItem.class); + WindowedValue element1 = WindowedValues.valueInGlobalWindow("v1"); + when(kwi.elementWindowsIterable()).thenReturn(Collections.singletonList(element1)); + + ValueInEmptyWindows> streamingValue = + new ValueInEmptyWindows<>(kwi); + streamingCounter.update(streamingValue); + + long elementCount = + (Long) + counterSet + .getExistingCounter( + CounterName.named(DataflowOutputCounter.getElementCounterName(OUTPUT_NAME))) + .getAggregate(); + assertEquals(1L, elementCount); + } + + @Test + public void testStandardWindowedValueCounting() throws Exception { + DataflowOutputCounter counter = + DataflowOutputCounter.create(OUTPUT_NAME, counterSet, nameContext, false); + + WindowedValue standardValue = WindowedValues.valueInGlobalWindow("v1"); + counter.update(standardValue); + + long elementCount = + (Long) + counterSet + .getExistingCounter( + CounterName.named(DataflowOutputCounter.getElementCounterName(OUTPUT_NAME))) + .getAggregate(); + assertEquals(1L, elementCount); + } +} From ddaa606acd2fe5ea88ed8c4d82371d4a95352878 Mon Sep 17 00:00:00 2001 From: Shunping Huang Date: Wed, 29 Jul 2026 13:23:59 -0400 Subject: [PATCH 07/11] Minor change on tests. --- .../runners/dataflow/worker/DataflowOutputCounterTest.java | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounterTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounterTest.java index c32964c58587..f95da5b13c80 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounterTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounterTest.java @@ -21,7 +21,7 @@ import static org.mockito.Mockito.mock; import static org.mockito.Mockito.when; -import java.util.Collections; +import java.util.Arrays; import org.apache.beam.runners.core.KeyedWorkItem; import org.apache.beam.runners.dataflow.worker.counters.CounterName; import org.apache.beam.runners.dataflow.worker.counters.CounterSet; @@ -74,7 +74,8 @@ public void testStreamingOutputCounterWithKeyedWorkItem() throws Exception { KeyedWorkItem kwi = mock(KeyedWorkItem.class); WindowedValue element1 = WindowedValues.valueInGlobalWindow("v1"); - when(kwi.elementWindowsIterable()).thenReturn(Collections.singletonList(element1)); + WindowedValue element2 = WindowedValues.valueInGlobalWindow("v2"); + when(kwi.elementWindowsIterable()).thenReturn(Arrays.asList(element1, element2)); ValueInEmptyWindows> streamingValue = new ValueInEmptyWindows<>(kwi); @@ -86,7 +87,7 @@ public void testStreamingOutputCounterWithKeyedWorkItem() throws Exception { .getExistingCounter( CounterName.named(DataflowOutputCounter.getElementCounterName(OUTPUT_NAME))) .getAggregate(); - assertEquals(1L, elementCount); + assertEquals(2L, elementCount); } @Test From 86cf801ed2775936befad3f201cebc652c8928fe Mon Sep 17 00:00:00 2001 From: Shunping Huang Date: Wed, 29 Jul 2026 15:43:24 -0400 Subject: [PATCH 08/11] Address reviewer comments --- .../beam/runners/core/KeyedWorkItem.java | 4 +- .../beam/runners/core/ReduceFnRunner.java | 7 +- .../worker/BatchDataflowOutputCounter.java | 53 --------------- .../worker/DataflowOutputCounter.java | 65 ++++++++++++++---- .../StreamingDataflowOutputCounter.java | 67 ------------------- .../worker/WindmillKeyedWorkItem.java | 60 +++++------------ .../worker/DataflowOutputCounterTest.java | 6 +- 7 files changed, 78 insertions(+), 184 deletions(-) delete mode 100644 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/BatchDataflowOutputCounter.java delete mode 100644 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowOutputCounter.java diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/KeyedWorkItem.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/KeyedWorkItem.java index cd7b6c10e500..2be8b0790301 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/KeyedWorkItem.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/KeyedWorkItem.java @@ -41,7 +41,7 @@ public interface KeyedWorkItem { * Useful for lightweight inspection of windowing metadata without payload deserialization * overhead. */ - default Iterable> elementWindowsIterable() { - return elementsIterable(); + default Iterable> elementWindowsIterable() { + return (Iterable) elementsIterable(); } } diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java index 52faa2f815b9..e49c858393f4 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java @@ -60,6 +60,7 @@ import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.FluentIterable; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableSet; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Iterables; import org.checkerframework.checker.nullness.qual.Nullable; import org.joda.time.Duration; import org.joda.time.Instant; @@ -371,9 +372,9 @@ public void processElements(Iterable> values) throws Excep } private void processElementsInternal( - Iterable> elementWindows, Iterable> values) + Iterable> elementWindows, Iterable> values) throws Exception { - if (!elementWindows.iterator().hasNext()) { + if (Iterables.isEmpty(elementWindows)) { return; } @@ -437,7 +438,7 @@ public void persist() { } /** Extract the windows associated with the values. */ - private Set collectWindows(Iterable> values) throws Exception { + private Set collectWindows(Iterable> values) throws Exception { Set windows = new HashSet<>(); for (WindowedValue value : values) { for (BoundedWindow untypedWindow : value.getWindows()) { diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/BatchDataflowOutputCounter.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/BatchDataflowOutputCounter.java deleted file mode 100644 index b512794fdfed..000000000000 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/BatchDataflowOutputCounter.java +++ /dev/null @@ -1,53 +0,0 @@ -/* - * 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. - */ -package org.apache.beam.runners.dataflow.worker; - -import org.apache.beam.runners.core.ElementByteSizeObservable; -import org.apache.beam.runners.dataflow.worker.counters.CounterFactory; -import org.apache.beam.runners.dataflow.worker.counters.NameContext; - -/** - * A Dataflow output counter specific to Batch pipelines. In batch pipelines, empty-window elements - * (e.g. GroupingShuffleReader emitting KV) represent a single PCollection - * element output. - */ -@SuppressWarnings({ - "nullness" // TODO(https://github.com/apache/beam/issues/20497) -}) -public class BatchDataflowOutputCounter extends DataflowOutputCounter { - - public BatchDataflowOutputCounter( - String outputName, CounterFactory counterFactory, NameContext nameContext) { - super(outputName, counterFactory, nameContext); - } - - public BatchDataflowOutputCounter( - String outputName, - ElementByteSizeObservable elementByteSizeObservable, - CounterFactory counterFactory, - NameContext nameContext) { - super(outputName, elementByteSizeObservable, counterFactory, nameContext); - } - - @Override - protected void updateEmptyWindows(Object elem) { - // Non-KeyedWorkItem wrapped in ValueInEmptyWindows (e.g. GroupingShuffleReader KV output for - // Batch GBK) - elementCount.addValue(1L); - } -} diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounter.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounter.java index 04b4c9b4f32d..002783f38d89 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounter.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounter.java @@ -18,6 +18,7 @@ package org.apache.beam.runners.dataflow.worker; import org.apache.beam.runners.core.ElementByteSizeObservable; +import org.apache.beam.runners.core.KeyedWorkItem; import org.apache.beam.runners.dataflow.worker.counters.Counter; import org.apache.beam.runners.dataflow.worker.counters.CounterFactory; import org.apache.beam.runners.dataflow.worker.counters.CounterName; @@ -40,7 +41,8 @@ public class DataflowOutputCounter implements ElementCounter { private static final String MEAN_BYTE_COUNTER_NAME = "-MeanByteCount"; private OutputObjectAndByteCounter objectAndByteCounter; - protected Counter elementCount; + private Counter elementCount; + private final boolean isStreaming; public static DataflowOutputCounter create( String outputName, @@ -48,11 +50,8 @@ public static DataflowOutputCounter create( CounterFactory counterFactory, NameContext nameContext, boolean isStreaming) { - return isStreaming - ? new StreamingDataflowOutputCounter( - outputName, elementByteSizeObservable, counterFactory, nameContext) - : new BatchDataflowOutputCounter( - outputName, elementByteSizeObservable, counterFactory, nameContext); + return new DataflowOutputCounter( + outputName, elementByteSizeObservable, counterFactory, nameContext, isStreaming); } public static DataflowOutputCounter create( @@ -65,7 +64,15 @@ public static DataflowOutputCounter create( public DataflowOutputCounter( String outputName, CounterFactory counterFactory, NameContext nameContext) { - this(outputName, null, counterFactory, nameContext); + this(outputName, null, counterFactory, nameContext, false); + } + + public DataflowOutputCounter( + String outputName, + CounterFactory counterFactory, + NameContext nameContext, + boolean isStreaming) { + this(outputName, null, counterFactory, nameContext, isStreaming); } public DataflowOutputCounter( @@ -73,9 +80,19 @@ public DataflowOutputCounter( ElementByteSizeObservable elementByteSizeObservable, CounterFactory counterFactory, NameContext nameContext) { - objectAndByteCounter = + this(outputName, elementByteSizeObservable, counterFactory, nameContext, false); + } + + public DataflowOutputCounter( + String outputName, + ElementByteSizeObservable elementByteSizeObservable, + CounterFactory counterFactory, + NameContext nameContext, + boolean isStreaming) { + this.isStreaming = isStreaming; + this.objectAndByteCounter = new OutputObjectAndByteCounter(elementByteSizeObservable, counterFactory, nameContext); - objectAndByteCounter.countMeanByte(outputName + MEAN_BYTE_COUNTER_NAME); + this.objectAndByteCounter.countMeanByte(outputName + MEAN_BYTE_COUNTER_NAME); createElementCounter(counterFactory, outputName + ELEMENT_COUNTER_NAME); } @@ -84,15 +101,39 @@ public void update(Object elem) throws Exception { objectAndByteCounter.update(elem); long windowsSize = ((WindowedValue) elem).getWindows().size(); if (windowsSize == 0) { - updateEmptyWindows(elem); + updateEmptyWindows((WindowedValue) elem); } else { // Standard WindowedValue. elementCount.addValue(windowsSize); } } - protected void updateEmptyWindows(Object elem) throws Exception { - elementCount.addValue(1L); + protected void updateEmptyWindows(WindowedValue elem) { + if (isStreaming) { + Object value = elem.getValue(); + if (value instanceof KeyedWorkItem) { + // KeyedWorkItem wrapped in ValueInEmptyWindows + // (e.g. WindowingWindmillReader for Streaming GBK) + KeyedWorkItem keyedWorkItem = (KeyedWorkItem) value; + long totalElementCount = 0; + // Iterate through elementWindowsIterable and ignore timers in KeyedWorkItem. + // Uses lightweight metadata-only iteration without payload deserialization overhead. + for (WindowedValue element : keyedWorkItem.elementWindowsIterable()) { + long elementWindowsSize = element.getWindows().size(); + // Fan out for windows. + totalElementCount += (elementWindowsSize == 0 ? 1L : elementWindowsSize); + } + elementCount.addValue(totalElementCount); + } else { + // NOTE: in streaming mode, this should not normally happen. + // Counting as 1 element serves as a fallback to maintain counter behavior without failing execution. + elementCount.addValue(1L); + } + } else { + // Non-KeyedWorkItem wrapped in ValueInEmptyWindows + // (e.g. GroupingShuffleReader KV output for Batch GBK) + elementCount.addValue(1L); + } } @Override diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowOutputCounter.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowOutputCounter.java deleted file mode 100644 index 1acd2917c32d..000000000000 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowOutputCounter.java +++ /dev/null @@ -1,67 +0,0 @@ -/* - * 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. - */ -package org.apache.beam.runners.dataflow.worker; - -import org.apache.beam.runners.core.ElementByteSizeObservable; -import org.apache.beam.runners.core.KeyedWorkItem; -import org.apache.beam.runners.dataflow.worker.counters.CounterFactory; -import org.apache.beam.runners.dataflow.worker.counters.NameContext; -import org.apache.beam.sdk.values.WindowedValue; - -/** - * A Dataflow output counter specific to Streaming pipelines. Unpacks {@link KeyedWorkItem}s in - * empty windows (e.g. emitted by WindowingWindmillReader) and counts element windows (ignoring - * timers) using lightweight metadata-only iteration via {@link - * KeyedWorkItem#elementWindowsIterable()}. - */ -@SuppressWarnings({ - "nullness" // TODO(https://github.com/apache/beam/issues/20497) -}) -public class StreamingDataflowOutputCounter extends DataflowOutputCounter { - - public StreamingDataflowOutputCounter( - String outputName, CounterFactory counterFactory, NameContext nameContext) { - super(outputName, counterFactory, nameContext); - } - - public StreamingDataflowOutputCounter( - String outputName, - ElementByteSizeObservable elementByteSizeObservable, - CounterFactory counterFactory, - NameContext nameContext) { - super(outputName, elementByteSizeObservable, counterFactory, nameContext); - } - - @Override - protected void updateEmptyWindows(Object elem) { - Object value = ((WindowedValue) elem).getValue(); - if (value instanceof KeyedWorkItem) { - KeyedWorkItem keyedWorkItem = (KeyedWorkItem) value; - long totalElementCount = 0; - // Iterate only through elementWindowsIterable and ignore timers in KeyedWorkItem. - for (WindowedValue element : keyedWorkItem.elementWindowsIterable()) { - long elementWindowsSize = element.getWindows().size(); - // Fan out for windows. - totalElementCount += (elementWindowsSize == 0 ? 1L : elementWindowsSize); - } - elementCount.addValue(totalElementCount); - } else { - elementCount.addValue(1L); - } - } -} 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 1df2fc5eb43b..7ab18867fdd8 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 @@ -139,6 +139,16 @@ public Iterable timersIterable() { } private @Nullable WindowedValue parseElem(Windmill.Message message) { + return parseElemInternal(message, true); + } + + private @Nullable WindowedValue parseElemWindowOnly(Windmill.Message message) { + return parseElemInternal(message, false); + } + + @SuppressWarnings("nullness") + private @Nullable WindowedValue parseElemInternal( + Windmill.Message message, boolean parseValue) { try { Instant timestamp = WindmillTimeUtils.windmillToHarnessTimestamp(message.getTimestamp()); Collection windows = @@ -159,51 +169,13 @@ public Iterable timersIterable() { : CausedByDrain.NORMAL; valueKind = WindmillValueKindHelper.fromProto(elementMetadata.getValueKind()); } - InputStream inputStream = message.getData().newInput(); - ElemT value = valueCoder.decode(inputStream, Coder.Context.OUTER); - return WindowedValues.of( - value, - timestamp, - windows, - paneInfo, - null, - null, - drainingValueFromUpstream, - null, - valueKind); - } catch (RuntimeException | IOException e) { - if (!skipUndecodableElements) { - throw new RuntimeException(e); - } - LOG.error( - "Skipping input element for work token {} on sharding key {} due to decoding error", - workItem.getWorkToken(), - workItem.getShardingKey(), - e); - return null; - } - } - - @SuppressWarnings("nullness") - private @Nullable WindowedValue parseElemWindowOnly(Windmill.Message message) { - try { - Instant timestamp = WindmillTimeUtils.windmillToHarnessTimestamp(message.getTimestamp()); - Collection windows = - WindmillSink.decodeMetadataWindows(windowsCoder, message.getMetadata()); - PaneInfo paneInfo = WindmillSink.decodeMetadataPane(message.getMetadata()); - CausedByDrain drainingValueFromUpstream = CausedByDrain.NORMAL; - ValueKind valueKind = ValueKind.INSERT; - if (WindowedValues.WindowedValueCoder.isMetadataSupported()) { - BeamFnApi.Elements.ElementMetadata elementMetadata = - WindmillSink.decodeAdditionalMetadata(windowsCoder, message.getMetadata()); - drainingValueFromUpstream = - elementMetadata.getDrain() == BeamFnApi.Elements.DrainMode.Enum.DRAINING - ? CausedByDrain.CAUSED_BY_DRAIN - : CausedByDrain.NORMAL; - valueKind = WindmillValueKindHelper.fromProto(elementMetadata.getValueKind()); + ElemT value = null; + if (parseValue) { + InputStream inputStream = message.getData().newInput(); + value = valueCoder.decode(inputStream, Coder.Context.OUTER); } return WindowedValues.of( - (ElemT) null, + value, timestamp, windows, paneInfo, @@ -227,7 +199,7 @@ public Iterable timersIterable() { @Override @SuppressWarnings("nullness") - public Iterable> elementWindowsIterable() { + public Iterable> elementWindowsIterable() { return FluentIterable.from(workItem.getMessageBundlesList()) .transformAndConcat(Windmill.InputMessageBundle::getMessagesList) .transform(this::parseElemWindowOnly) diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounterTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounterTest.java index f95da5b13c80..e11412b88655 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounterTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounterTest.java @@ -18,8 +18,8 @@ package org.apache.beam.runners.dataflow.worker; import static org.junit.Assert.assertEquals; +import static org.mockito.Mockito.doReturn; import static org.mockito.Mockito.mock; -import static org.mockito.Mockito.when; import java.util.Arrays; import org.apache.beam.runners.core.KeyedWorkItem; @@ -35,7 +35,7 @@ import org.junit.runner.RunWith; import org.junit.runners.JUnit4; -/** Tests for {@link BatchDataflowOutputCounter} and {@link StreamingDataflowOutputCounter}. */ +/** Tests for {@link DataflowOutputCounter}. */ @RunWith(JUnit4.class) public class DataflowOutputCounterTest { private static final String OUTPUT_NAME = "test_output"; @@ -75,7 +75,7 @@ public void testStreamingOutputCounterWithKeyedWorkItem() throws Exception { KeyedWorkItem kwi = mock(KeyedWorkItem.class); WindowedValue element1 = WindowedValues.valueInGlobalWindow("v1"); WindowedValue element2 = WindowedValues.valueInGlobalWindow("v2"); - when(kwi.elementWindowsIterable()).thenReturn(Arrays.asList(element1, element2)); + doReturn(Arrays.asList(element1, element2)).when(kwi).elementWindowsIterable(); ValueInEmptyWindows> streamingValue = new ValueInEmptyWindows<>(kwi); From 24e6401b2e8b09199a2549288f86f0e1ba71c464 Mon Sep 17 00:00:00 2001 From: Shunping Huang Date: Wed, 29 Jul 2026 15:49:13 -0400 Subject: [PATCH 09/11] Spotless --- .../beam/runners/dataflow/worker/DataflowOutputCounter.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounter.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounter.java index 002783f38d89..8255049115de 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounter.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounter.java @@ -126,7 +126,8 @@ protected void updateEmptyWindows(WindowedValue elem) { elementCount.addValue(totalElementCount); } else { // NOTE: in streaming mode, this should not normally happen. - // Counting as 1 element serves as a fallback to maintain counter behavior without failing execution. + // Counting as 1 element serves as a fallback to maintain counter behavior without failing + // execution. elementCount.addValue(1L); } } else { From c9f888e2a576cf80d2c07462d1ca25f15be61b63 Mon Sep 17 00:00:00 2001 From: Shunping Huang Date: Wed, 29 Jul 2026 15:50:46 -0400 Subject: [PATCH 10/11] Remove unnecessary comments --- .../beam/runners/dataflow/worker/DataflowOutputCounterTest.java | 1 - 1 file changed, 1 deletion(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounterTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounterTest.java index e11412b88655..b5c49ee639b5 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounterTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounterTest.java @@ -53,7 +53,6 @@ public void testBatchOutputCounterWithEmptyWindows() throws Exception { DataflowOutputCounter batchCounter = DataflowOutputCounter.create(OUTPUT_NAME, counterSet, nameContext, false); - // Non-KeyedWorkItem value in empty windows (e.g. GroupingShuffleReader output) ValueInEmptyWindows> shuffleValue = new ValueInEmptyWindows<>(KV.of("key", "value")); batchCounter.update(shuffleValue); From b53bad224961350cc53c4061c02d79fa49649bbc Mon Sep 17 00:00:00 2001 From: Shunping Huang Date: Wed, 29 Jul 2026 20:48:49 -0400 Subject: [PATCH 11/11] Address comments --- .../worker/DataflowOutputCounter.java | 29 ++++--------------- .../IntrinsicMapTaskExecutorFactory.java | 5 ---- .../IntrinsicMapTaskExecutorFactoryTest.java | 14 +++++---- 3 files changed, 14 insertions(+), 34 deletions(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounter.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounter.java index 8255049115de..1a927c03c61c 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounter.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounter.java @@ -25,6 +25,7 @@ import org.apache.beam.runners.dataflow.worker.counters.NameContext; import org.apache.beam.runners.dataflow.worker.util.common.worker.ElementCounter; import org.apache.beam.runners.dataflow.worker.util.common.worker.OutputObjectAndByteCounter; +import org.apache.beam.sdk.annotations.Internal; import org.apache.beam.sdk.values.WindowedValue; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting; @@ -34,6 +35,7 @@ @SuppressWarnings({ "nullness" // TODO(https://github.com/apache/beam/issues/20497) }) +@Internal public class DataflowOutputCounter implements ElementCounter { /** Number of logical element and single window pairs that were processed. */ private static final String ELEMENT_COUNTER_NAME = "-ElementCount"; @@ -59,31 +61,10 @@ public static DataflowOutputCounter create( CounterFactory counterFactory, NameContext nameContext, boolean isStreaming) { - return create(outputName, null, counterFactory, nameContext, isStreaming); + return new DataflowOutputCounter(outputName, null, counterFactory, nameContext, isStreaming); } - public DataflowOutputCounter( - String outputName, CounterFactory counterFactory, NameContext nameContext) { - this(outputName, null, counterFactory, nameContext, false); - } - - public DataflowOutputCounter( - String outputName, - CounterFactory counterFactory, - NameContext nameContext, - boolean isStreaming) { - this(outputName, null, counterFactory, nameContext, isStreaming); - } - - public DataflowOutputCounter( - String outputName, - ElementByteSizeObservable elementByteSizeObservable, - CounterFactory counterFactory, - NameContext nameContext) { - this(outputName, elementByteSizeObservable, counterFactory, nameContext, false); - } - - public DataflowOutputCounter( + private DataflowOutputCounter( String outputName, ElementByteSizeObservable elementByteSizeObservable, CounterFactory counterFactory, @@ -108,7 +89,7 @@ public void update(Object elem) throws Exception { } } - protected void updateEmptyWindows(WindowedValue elem) { + private void updateEmptyWindows(WindowedValue elem) { if (isStreaming) { Object value = elem.getValue(); if (value instanceof KeyedWorkItem) { diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/IntrinsicMapTaskExecutorFactory.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/IntrinsicMapTaskExecutorFactory.java index 56cc785cb28e..3ea29787eb3b 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/IntrinsicMapTaskExecutorFactory.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/IntrinsicMapTaskExecutorFactory.java @@ -346,11 +346,6 @@ OperationNode createFlattenOperation( /** * Returns a function which can convert {@link InstructionOutput}s into {@link OutputReceiver}s. */ - static Function createOutputReceiversTransform( - final String stageName, final CounterFactory counterFactory) { - return createOutputReceiversTransform(stageName, counterFactory, false); - } - static Function createOutputReceiversTransform( final String stageName, final CounterFactory counterFactory, final boolean isStreaming) { return new TypeSafeNodeFunction(InstructionOutputNode.class) { diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/IntrinsicMapTaskExecutorFactoryTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/IntrinsicMapTaskExecutorFactoryTest.java index 3443ae0022bc..d3a424758f66 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/IntrinsicMapTaskExecutorFactoryTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/IntrinsicMapTaskExecutorFactoryTest.java @@ -330,7 +330,8 @@ public void testCreateReadOperation() throws Exception { when(network.successors(instructionNode)) .thenReturn( ImmutableSet.of( - IntrinsicMapTaskExecutorFactory.createOutputReceiversTransform(STAGE, counterSet) + IntrinsicMapTaskExecutorFactory.createOutputReceiversTransform( + STAGE, counterSet, false) .apply( InstructionOutputNode.create( instructionNode.getParallelInstruction().getOutputs().get(0), @@ -535,7 +536,7 @@ public void testCreateParDoOperation() throws Exception { ExecutionLocation.UNKNOWN); Node outputReceiverNode = - IntrinsicMapTaskExecutorFactory.createOutputReceiversTransform(STAGE, counterSet) + IntrinsicMapTaskExecutorFactory.createOutputReceiversTransform(STAGE, counterSet, false) .apply( InstructionOutputNode.create( instructionNode.getParallelInstruction().getOutputs().get(0), PCOLLECTION_ID)); @@ -614,7 +615,8 @@ public void testCreatePartialGroupByKeyOperation() throws Exception { when(network.successors(instructionNode)) .thenReturn( ImmutableSet.of( - IntrinsicMapTaskExecutorFactory.createOutputReceiversTransform(STAGE, counterSet) + IntrinsicMapTaskExecutorFactory.createOutputReceiversTransform( + STAGE, counterSet, false) .apply( InstructionOutputNode.create( instructionNode.getParallelInstruction().getOutputs().get(0), @@ -669,7 +671,8 @@ public void testCreatePartialGroupByKeyOperationWithCombine() throws Exception { when(network.successors(instructionNode)) .thenReturn( ImmutableSet.of( - IntrinsicMapTaskExecutorFactory.createOutputReceiversTransform(STAGE, counterSet) + IntrinsicMapTaskExecutorFactory.createOutputReceiversTransform( + STAGE, counterSet, false) .apply( InstructionOutputNode.create( instructionNode.getParallelInstruction().getOutputs().get(0), @@ -750,7 +753,8 @@ public void testCreateFlattenOperation() throws Exception { when(network.successors(instructionNode)) .thenReturn( ImmutableSet.of( - IntrinsicMapTaskExecutorFactory.createOutputReceiversTransform(STAGE, counterSet) + IntrinsicMapTaskExecutorFactory.createOutputReceiversTransform( + STAGE, counterSet, false) .apply( InstructionOutputNode.create( instructionNode.getParallelInstruction().getOutputs().get(0),