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..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 @@ -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,27 @@ 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 + 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.elementsIterable()) { + 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); + } } else { + // Standard WindowedValue. elementCount.addValue(windowsSize); } }