From e66be4167032041fb01ea014226838ee88002371 Mon Sep 17 00:00:00 2001 From: Shunping Huang Date: Fri, 24 Jul 2026 22:44:36 -0400 Subject: [PATCH 1/3] 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 2/3] 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 3/3] 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); }