Skip to content

Commit 03db1a0

Browse files
authored
Fix DataflowOutputCounter calculation for ValueInEmptyWindows (#39487)
* 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. * Address non keyedworkitems * Spotless * Add elementWindowsIterable to only decode window metadata and use it in DataflowOutputCounter * Use the elementWindowsIterable in ReduceFnRunner * Refactor: separate batch and streaming DataflowOutputCounter implementations * Minor change on tests. * Address reviewer comments * Spotless * Remove unnecessary comments * Address comments
1 parent 7b9380b commit 03db1a0

10 files changed

Lines changed: 234 additions & 29 deletions

File tree

runners/core-java/src/main/java/org/apache/beam/runners/core/GroupAlsoByWindowViaWindowSetNewDoFn.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -109,7 +109,7 @@ public void processElement(ProcessContext c) throws Exception {
109109
reduceFn,
110110
c.getPipelineOptions());
111111

112-
reduceFnRunner.processElements(keyedWorkItem.elementsIterable());
112+
reduceFnRunner.processElements(keyedWorkItem);
113113
reduceFnRunner.onTimers(keyedWorkItem.timersIterable());
114114
reduceFnRunner.persist();
115115
}

runners/core-java/src/main/java/org/apache/beam/runners/core/KeyedWorkItem.java

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -35,4 +35,13 @@ public interface KeyedWorkItem<K, ElemT> {
3535

3636
/** Returns an iterable containing the elements. */
3737
Iterable<WindowedValue<ElemT>> elementsIterable();
38+
39+
/**
40+
* Returns an iterable containing windowed values without guaranteeing element payload decoding.
41+
* Useful for lightweight inspection of windowing metadata without payload deserialization
42+
* overhead.
43+
*/
44+
default Iterable<WindowedValue<?>> elementWindowsIterable() {
45+
return (Iterable) elementsIterable();
46+
}
3847
}

runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java

Lines changed: 15 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -60,6 +60,7 @@
6060
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting;
6161
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.FluentIterable;
6262
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableSet;
63+
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Iterables;
6364
import org.checkerframework.checker.nullness.qual.Nullable;
6465
import org.joda.time.Duration;
6566
import org.joda.time.Instant;
@@ -361,13 +362,24 @@ private Collection<W> windowsThatShouldFire(Set<W> windows) throws Exception {
361362
* setting holds, and invoking {@link ReduceFn#onTrigger}.
362363
* </ol>
363364
*/
365+
public void processElements(KeyedWorkItem<?, InputT> keyedWorkItem) throws Exception {
366+
processElementsInternal(
367+
keyedWorkItem.elementWindowsIterable(), keyedWorkItem.elementsIterable());
368+
}
369+
364370
public void processElements(Iterable<WindowedValue<InputT>> values) throws Exception {
365-
if (!values.iterator().hasNext()) {
371+
processElementsInternal(values, values);
372+
}
373+
374+
private void processElementsInternal(
375+
Iterable<? extends WindowedValue<?>> elementWindows, Iterable<WindowedValue<InputT>> values)
376+
throws Exception {
377+
if (Iterables.isEmpty(elementWindows)) {
366378
return;
367379
}
368380

369381
// Determine all the windows for elements.
370-
Set<W> windows = collectWindows(values);
382+
Set<W> windows = collectWindows(elementWindows);
371383
// If an incoming element introduces a new window, attempt to merge it into an existing
372384
// window eagerly.
373385
Map<W, W> windowToMergeResult = mergeWindows(windows);
@@ -426,7 +438,7 @@ public void persist() {
426438
}
427439

428440
/** Extract the windows associated with the values. */
429-
private Set<W> collectWindows(Iterable<WindowedValue<InputT>> values) throws Exception {
441+
private Set<W> collectWindows(Iterable<? extends WindowedValue<?>> values) throws Exception {
430442
Set<W> windows = new HashSet<>();
431443
for (WindowedValue<?> value : values) {
432444
for (BoundedWindow untypedWindow : value.getWindows()) {

runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounter.java

Lines changed: 57 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -18,12 +18,14 @@
1818
package org.apache.beam.runners.dataflow.worker;
1919

2020
import org.apache.beam.runners.core.ElementByteSizeObservable;
21+
import org.apache.beam.runners.core.KeyedWorkItem;
2122
import org.apache.beam.runners.dataflow.worker.counters.Counter;
2223
import org.apache.beam.runners.dataflow.worker.counters.CounterFactory;
2324
import org.apache.beam.runners.dataflow.worker.counters.CounterName;
2425
import org.apache.beam.runners.dataflow.worker.counters.NameContext;
2526
import org.apache.beam.runners.dataflow.worker.util.common.worker.ElementCounter;
2627
import org.apache.beam.runners.dataflow.worker.util.common.worker.OutputObjectAndByteCounter;
28+
import org.apache.beam.sdk.annotations.Internal;
2729
import org.apache.beam.sdk.values.WindowedValue;
2830
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting;
2931

@@ -33,6 +35,7 @@
3335
@SuppressWarnings({
3436
"nullness" // TODO(https://github.com/apache/beam/issues/20497)
3537
})
38+
@Internal
3639
public class DataflowOutputCounter implements ElementCounter {
3740
/** Number of logical element and single window pairs that were processed. */
3841
private static final String ELEMENT_COUNTER_NAME = "-ElementCount";
@@ -41,20 +44,36 @@ public class DataflowOutputCounter implements ElementCounter {
4144

4245
private OutputObjectAndByteCounter objectAndByteCounter;
4346
private Counter<Long, ?> elementCount;
47+
private final boolean isStreaming;
4448

45-
public DataflowOutputCounter(
46-
String outputName, CounterFactory counterFactory, NameContext nameContext) {
47-
this(outputName, null, counterFactory, nameContext);
49+
public static DataflowOutputCounter create(
50+
String outputName,
51+
ElementByteSizeObservable<?> elementByteSizeObservable,
52+
CounterFactory counterFactory,
53+
NameContext nameContext,
54+
boolean isStreaming) {
55+
return new DataflowOutputCounter(
56+
outputName, elementByteSizeObservable, counterFactory, nameContext, isStreaming);
57+
}
58+
59+
public static DataflowOutputCounter create(
60+
String outputName,
61+
CounterFactory counterFactory,
62+
NameContext nameContext,
63+
boolean isStreaming) {
64+
return new DataflowOutputCounter(outputName, null, counterFactory, nameContext, isStreaming);
4865
}
4966

50-
public DataflowOutputCounter(
67+
private DataflowOutputCounter(
5168
String outputName,
5269
ElementByteSizeObservable<?> elementByteSizeObservable,
5370
CounterFactory counterFactory,
54-
NameContext nameContext) {
55-
objectAndByteCounter =
71+
NameContext nameContext,
72+
boolean isStreaming) {
73+
this.isStreaming = isStreaming;
74+
this.objectAndByteCounter =
5675
new OutputObjectAndByteCounter(elementByteSizeObservable, counterFactory, nameContext);
57-
objectAndByteCounter.countMeanByte(outputName + MEAN_BYTE_COUNTER_NAME);
76+
this.objectAndByteCounter.countMeanByte(outputName + MEAN_BYTE_COUNTER_NAME);
5877
createElementCounter(counterFactory, outputName + ELEMENT_COUNTER_NAME);
5978
}
6079

@@ -63,15 +82,42 @@ public void update(Object elem) throws Exception {
6382
objectAndByteCounter.update(elem);
6483
long windowsSize = ((WindowedValue<?>) elem).getWindows().size();
6584
if (windowsSize == 0) {
66-
// GroupingShuffleReader produces ValueInEmptyWindows.
67-
// For now, we count the element at least once to keep the current counter
68-
// behavior.
69-
elementCount.addValue(1L);
85+
updateEmptyWindows((WindowedValue<?>) elem);
7086
} else {
87+
// Standard WindowedValue.
7188
elementCount.addValue(windowsSize);
7289
}
7390
}
7491

92+
private void updateEmptyWindows(WindowedValue<?> elem) {
93+
if (isStreaming) {
94+
Object value = elem.getValue();
95+
if (value instanceof KeyedWorkItem<?, ?>) {
96+
// KeyedWorkItem wrapped in ValueInEmptyWindows
97+
// (e.g. WindowingWindmillReader for Streaming GBK)
98+
KeyedWorkItem<?, ?> keyedWorkItem = (KeyedWorkItem<?, ?>) value;
99+
long totalElementCount = 0;
100+
// Iterate through elementWindowsIterable and ignore timers in KeyedWorkItem.
101+
// Uses lightweight metadata-only iteration without payload deserialization overhead.
102+
for (WindowedValue<?> element : keyedWorkItem.elementWindowsIterable()) {
103+
long elementWindowsSize = element.getWindows().size();
104+
// Fan out for windows.
105+
totalElementCount += (elementWindowsSize == 0 ? 1L : elementWindowsSize);
106+
}
107+
elementCount.addValue(totalElementCount);
108+
} else {
109+
// NOTE: in streaming mode, this should not normally happen.
110+
// Counting as 1 element serves as a fallback to maintain counter behavior without failing
111+
// execution.
112+
elementCount.addValue(1L);
113+
}
114+
} else {
115+
// Non-KeyedWorkItem wrapped in ValueInEmptyWindows
116+
// (e.g. GroupingShuffleReader KV output for Batch GBK)
117+
elementCount.addValue(1L);
118+
}
119+
}
120+
75121
@Override
76122
public void finishLazyUpdate(Object elem) {
77123
objectAndByteCounter.finishLazyUpdate(elem);

runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/IntrinsicMapTaskExecutorFactory.java

Lines changed: 7 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -63,6 +63,7 @@
6363
import org.apache.beam.sdk.coders.KvCoder;
6464
import org.apache.beam.sdk.fn.IdGenerator;
6565
import org.apache.beam.sdk.options.PipelineOptions;
66+
import org.apache.beam.sdk.options.StreamingOptions;
6667
import org.apache.beam.sdk.util.common.ElementByteSizeObserver;
6768
import org.apache.beam.sdk.values.TupleTag;
6869
import org.apache.beam.sdk.values.WindowedValues.WindowedValueCoder;
@@ -102,8 +103,9 @@ public DataflowMapTaskExecutor create(
102103
IdGenerator idGenerator) {
103104

104105
// Swap out all the InstructionOutput nodes with OutputReceiver nodes
106+
boolean isStreaming = options.as(StreamingOptions.class).isStreaming();
105107
Networks.replaceDirectedNetworkNodes(
106-
network, createOutputReceiversTransform(stageName, counterSet));
108+
network, createOutputReceiversTransform(stageName, counterSet, isStreaming));
107109

108110
// Swap out all the ParallelInstruction nodes with Operation nodes. While updating the network,
109111
// we keep track of
@@ -345,7 +347,7 @@ OperationNode createFlattenOperation(
345347
* Returns a function which can convert {@link InstructionOutput}s into {@link OutputReceiver}s.
346348
*/
347349
static Function<Node, Node> createOutputReceiversTransform(
348-
final String stageName, final CounterFactory counterFactory) {
350+
final String stageName, final CounterFactory counterFactory, final boolean isStreaming) {
349351
return new TypeSafeNodeFunction<InstructionOutputNode>(InstructionOutputNode.class) {
350352
@Override
351353
public Node typedApply(InstructionOutputNode input) {
@@ -355,15 +357,16 @@ public Node typedApply(InstructionOutputNode input) {
355357
CloudObjects.coderFromCloudObject(CloudObject.fromSpec(cloudOutput.getCodec()));
356358

357359
ElementCounter outputCounter =
358-
new DataflowOutputCounter(
360+
DataflowOutputCounter.create(
359361
cloudOutput.getName(),
360362
new ElementByteSizeObservableCoder<>(coder),
361363
counterFactory,
362364
NameContext.create(
363365
stageName,
364366
cloudOutput.getOriginalName(),
365367
cloudOutput.getSystemName(),
366-
cloudOutput.getName()));
368+
cloudOutput.getName()),
369+
isStreaming);
367370
outputReceiver.addOutputCounter(outputCounter);
368371

369372
return OutputReceiverNode.create(outputReceiver, coder, input.getPcollectionId());

runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/SimpleParDoFnHelpers.java

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -190,9 +190,10 @@ public <TagT> void output(TupleTag<TagT> tag, WindowedValue<TagT> output) {
190190
// doesn't today.)
191191
OutputReceiver undeclaredReceiver = new OutputReceiver();
192192

193+
boolean isStreaming = options.as(StreamingOptions.class).isStreaming();
193194
ElementCounter outputCounter =
194-
new DataflowOutputCounter(
195-
outputName, counterFactory, stepContext.getNameContext());
195+
DataflowOutputCounter.create(
196+
outputName, counterFactory, stepContext.getNameContext(), isStreaming);
196197
undeclaredReceiver.addOutputCounter(outputCounter);
197198
undeclaredOutputs.put(tag, undeclaredReceiver);
198199
receiver = undeclaredReceiver;

runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingGroupAlsoByWindowViaWindowSetFn.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -93,7 +93,7 @@ public void processElement(
9393
reduceFn,
9494
options);
9595

96-
reduceFnRunner.processElements(keyedWorkItem.elementsIterable());
96+
reduceFnRunner.processElements(keyedWorkItem);
9797
reduceFnRunner.onTimers(keyedWorkItem.timersIterable());
9898
reduceFnRunner.persist();
9999
}

runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItem.java

Lines changed: 24 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -140,6 +140,16 @@ public Iterable<TimerData> timersIterable() {
140140
}
141141

142142
private @Nullable WindowedValue<ElemT> parseElem(Windmill.Message message) {
143+
return parseElemInternal(message, true);
144+
}
145+
146+
private @Nullable WindowedValue<?> parseElemWindowOnly(Windmill.Message message) {
147+
return parseElemInternal(message, false);
148+
}
149+
150+
@SuppressWarnings("nullness")
151+
private @Nullable WindowedValue<ElemT> parseElemInternal(
152+
Windmill.Message message, boolean parseValue) {
143153
try {
144154
Instant timestamp = WindmillTimeUtils.windmillToHarnessTimestamp(message.getTimestamp());
145155
Collection<? extends BoundedWindow> windows =
@@ -162,8 +172,11 @@ public Iterable<TimerData> timersIterable() {
162172
valueKind = WindmillValueKindHelper.fromProto(elementMetadata.getValueKind());
163173
openTelemetryContext = WindmillOpenTelemetryContextPropagator.read(elementMetadata);
164174
}
165-
InputStream inputStream = message.getData().newInput();
166-
ElemT value = valueCoder.decode(inputStream, Coder.Context.OUTER);
175+
ElemT value = null;
176+
if (parseValue) {
177+
InputStream inputStream = message.getData().newInput();
178+
value = valueCoder.decode(inputStream, Coder.Context.OUTER);
179+
}
167180
return WindowedValues.of(
168181
value,
169182
timestamp,
@@ -187,6 +200,15 @@ public Iterable<TimerData> timersIterable() {
187200
}
188201
}
189202

203+
@Override
204+
@SuppressWarnings("nullness")
205+
public Iterable<WindowedValue<?>> elementWindowsIterable() {
206+
return FluentIterable.from(workItem.getMessageBundlesList())
207+
.transformAndConcat(Windmill.InputMessageBundle::getMessagesList)
208+
.transform(this::parseElemWindowOnly)
209+
.filter(Objects::nonNull);
210+
}
211+
190212
@Override
191213
@SuppressWarnings("nullness")
192214
public Iterable<WindowedValue<ElemT>> elementsIterable() {

0 commit comments

Comments
 (0)