This is an automated email from the ASF dual-hosted git repository.
stankiewicz pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git
The following commit(s) were added to refs/heads/master by this push:
new dedbf9c8522 Prevent CommitTooLargeException in OrderedEventProcessor
on duplicate floods (#40242)
dedbf9c8522 is described below
commit dedbf9c8522ff316cfdb9ab2ec51118c85ee195f
Author: Vishal More <[email protected]>
AuthorDate: Mon Oct 5 20:03:43 2026 +0530
Prevent CommitTooLargeException in OrderedEventProcessor on duplicate
floods (#40242)
* Fix: Prevent CommitTooLargeException in OrderedEventProcessor on
duplicate floods
* Apply spotless formatting fixes
---
.../ordered/GlobalSequencesProcessorDoFn.java | 2 +
.../sdk/extensions/ordered/ProcessingState.java | 4 ++
.../beam/sdk/extensions/ordered/ProcessorDoFn.java | 20 ++++++---
.../ordered/SequencePerKeyProcessorDoFn.java | 2 +
.../OrderedEventProcessorPerKeySequenceTest.java | 48 ++++++++++++++++++++++
5 files changed, 70 insertions(+), 6 deletions(-)
diff --git
a/sdks/java/extensions/ordered/src/main/java/org/apache/beam/sdk/extensions/ordered/GlobalSequencesProcessorDoFn.java
b/sdks/java/extensions/ordered/src/main/java/org/apache/beam/sdk/extensions/ordered/GlobalSequencesProcessorDoFn.java
index 3c6c72eb486..0f93329407b 100644
---
a/sdks/java/extensions/ordered/src/main/java/org/apache/beam/sdk/extensions/ordered/GlobalSequencesProcessorDoFn.java
+++
b/sdks/java/extensions/ordered/src/main/java/org/apache/beam/sdk/extensions/ordered/GlobalSequencesProcessorDoFn.java
@@ -175,6 +175,7 @@ class GlobalSequencesProcessorDoFn<
if (numberOfResultsBeforeBundleStart == null) {
// Per key processing is synchronized by Beam. There is no need to have
it here.
numberOfResultsBeforeBundleStart = processingState.getResultCount();
+ numberOfDuplicatesBeforeBundleStart = processingState.getDuplicates();
}
processingState.eventReceived();
@@ -262,6 +263,7 @@ class GlobalSequencesProcessorDoFn<
}
this.numberOfResultsBeforeBundleStart = processingState.getResultCount();
+ this.numberOfDuplicatesBeforeBundleStart = processingState.getDuplicates();
state =
processBufferedEventRange(
diff --git
a/sdks/java/extensions/ordered/src/main/java/org/apache/beam/sdk/extensions/ordered/ProcessingState.java
b/sdks/java/extensions/ordered/src/main/java/org/apache/beam/sdk/extensions/ordered/ProcessingState.java
index 2980c2d614f..f7c8490de6c 100644
---
a/sdks/java/extensions/ordered/src/main/java/org/apache/beam/sdk/extensions/ordered/ProcessingState.java
+++
b/sdks/java/extensions/ordered/src/main/java/org/apache/beam/sdk/extensions/ordered/ProcessingState.java
@@ -313,6 +313,10 @@ class ProcessingState<KeyT> {
return resultCount - numberOfResultsBeforeBundleStart;
}
+ public long duplicatesProducedInBundle(long
numberOfDuplicatesBeforeBundleStart) {
+ return duplicates - numberOfDuplicatesBeforeBundleStart;
+ }
+
public void updateGlobalSequenceDetails(ContiguousSequenceRange updated) {
if (thereAreGloballySequencedEventsToBeProcessed()) {
// We don't update the timer if we can already process events in the
onTimer batch.
diff --git
a/sdks/java/extensions/ordered/src/main/java/org/apache/beam/sdk/extensions/ordered/ProcessorDoFn.java
b/sdks/java/extensions/ordered/src/main/java/org/apache/beam/sdk/extensions/ordered/ProcessorDoFn.java
index 3e97f85cd59..ead50556697 100644
---
a/sdks/java/extensions/ordered/src/main/java/org/apache/beam/sdk/extensions/ordered/ProcessorDoFn.java
+++
b/sdks/java/extensions/ordered/src/main/java/org/apache/beam/sdk/extensions/ordered/ProcessorDoFn.java
@@ -63,6 +63,7 @@ abstract class ProcessorDoFn<
private final long maxNumberOfResultsToProduce;
protected @Nullable Long numberOfResultsBeforeBundleStart = 0L;
+ protected @Nullable Long numberOfDuplicatesBeforeBundleStart = 0L;
ProcessorDoFn(
EventExaminer<EventT, StateT> eventExaminer,
@@ -85,12 +86,14 @@ abstract class ProcessorDoFn<
@StartBundle
public void onBundleStart() {
numberOfResultsBeforeBundleStart = null;
+ numberOfDuplicatesBeforeBundleStart = null;
}
@FinishBundle
public void onBundleFinish() {
// This might be necessary because this field is also used in a Timer
numberOfResultsBeforeBundleStart = null;
+ numberOfDuplicatesBeforeBundleStart = null;
}
/** Returns true if each event needs to be examined. */
@@ -267,10 +270,14 @@ abstract class ProcessorDoFn<
protected boolean reachedMaxResultCountForBundle(
ProcessingState<EventKeyT> processingState, Timer
largeBatchEmissionTimer) {
- boolean exceeded =
+ long resultsEmitted =
processingState.resultsProducedInBundle(
- numberOfResultsBeforeBundleStart == null ? 0 :
numberOfResultsBeforeBundleStart)
- >= maxNumberOfResultsToProduce;
+ numberOfResultsBeforeBundleStart == null ? 0 :
numberOfResultsBeforeBundleStart);
+ long duplicatesEmitted =
+ processingState.duplicatesProducedInBundle(
+ numberOfDuplicatesBeforeBundleStart == null ? 0 :
numberOfDuplicatesBeforeBundleStart);
+
+ boolean exceeded = (resultsEmitted + duplicatesEmitted) >=
maxNumberOfResultsToProduce;
if (exceeded) {
if (LOG.isTraceEnabled()) {
LOG.trace(
@@ -354,9 +361,10 @@ abstract class ProcessorDoFn<
beforeInitialSequence
? Reason.before_initial_sequence
: Reason.duplicate))));
- // TODO: When there is a large number of duplicates this can cause a
situation where
- // we produce too much output and the runner will start throwing
unrecoverable errors.
- // Need to add counting logic to accumulate both the normal and DLQ
outputs.
+ if (reachedMaxResultCountForBundle(processingState,
largeBatchEmissionTimer)) {
+ endClearRange = fromLong(eventSequence + 1);
+ break;
+ }
continue;
}
diff --git
a/sdks/java/extensions/ordered/src/main/java/org/apache/beam/sdk/extensions/ordered/SequencePerKeyProcessorDoFn.java
b/sdks/java/extensions/ordered/src/main/java/org/apache/beam/sdk/extensions/ordered/SequencePerKeyProcessorDoFn.java
index 486622d2069..d42e84079fd 100644
---
a/sdks/java/extensions/ordered/src/main/java/org/apache/beam/sdk/extensions/ordered/SequencePerKeyProcessorDoFn.java
+++
b/sdks/java/extensions/ordered/src/main/java/org/apache/beam/sdk/extensions/ordered/SequencePerKeyProcessorDoFn.java
@@ -159,6 +159,7 @@ class SequencePerKeyProcessorDoFn<
if (numberOfResultsBeforeBundleStart == null) {
// Per key processing is synchronized by Beam. There is no need to have
it here.
numberOfResultsBeforeBundleStart = processingState.getResultCount();
+ numberOfDuplicatesBeforeBundleStart = processingState.getDuplicates();
}
processingState.eventReceived();
@@ -259,6 +260,7 @@ class SequencePerKeyProcessorDoFn<
LOG.debug("Starting to process batch for key '{}'",
processingState.getKey());
this.numberOfResultsBeforeBundleStart = processingState.getResultCount();
+ this.numberOfDuplicatesBeforeBundleStart = processingState.getDuplicates();
processBufferedEvents(
processingState, state, bufferedEventsState, outputReceiver,
largeBatchEmissionTimer);
diff --git
a/sdks/java/extensions/ordered/src/test/java/org/apache/beam/sdk/extensions/ordered/OrderedEventProcessorPerKeySequenceTest.java
b/sdks/java/extensions/ordered/src/test/java/org/apache/beam/sdk/extensions/ordered/OrderedEventProcessorPerKeySequenceTest.java
index 5dad7ac1852..45a4a891adc 100644
---
a/sdks/java/extensions/ordered/src/test/java/org/apache/beam/sdk/extensions/ordered/OrderedEventProcessorPerKeySequenceTest.java
+++
b/sdks/java/extensions/ordered/src/test/java/org/apache/beam/sdk/extensions/ordered/OrderedEventProcessorPerKeySequenceTest.java
@@ -246,6 +246,54 @@ public class OrderedEventProcessorPerKeySequenceTest
extends OrderedEventProcess
DONT_PRODUCE_STATUS_ON_EVERY_EVENT);
}
+ @Test
+ public void testLargeNumberOfDuplicatesPaginatesCorrectly() throws
CannotProvideCoderException {
+ int maxResultsPerOutput = 10;
+ int duplicateCount = 25;
+ List<Event> events = new ArrayList<>();
+ events.add(Event.create(0, "id-1", "a"));
+
+ // Add 25 duplicates of the next event. The maxResultsPerOutput is 10, so
it should
+ // paginate the duplicates across multiple bundles.
+ for (int i = 0; i < duplicateCount; i++) {
+ events.add(Event.create(1, "id-1", "b"));
+ }
+
+ Collection<KV<String, OrderedProcessingStatus>> expectedStatuses = new
ArrayList<>();
+ expectedStatuses.add(
+ KV.of(
+ "id-1",
+ OrderedProcessingStatus.create(
+ 1L,
+ 0,
+ null,
+ null,
+ events.size(),
+ 2L,
+ duplicateCount - 1, // one is processed, the rest are
duplicates
+ false,
+ NOT_USED_FOR_TESTING)));
+
+ Collection<KV<String, String>> expectedOutput = new ArrayList<>();
+ expectedOutput.add(KV.of("id-1", "a"));
+ expectedOutput.add(KV.of("id-1", "ab"));
+
+ Collection<KV<String, KV<Long, UnprocessedEvent<String>>>> duplicates =
new ArrayList<>();
+ for (int i = 0; i < duplicateCount - 1; i++) {
+ duplicates.add(KV.of("id-1", KV.of(1L, UnprocessedEvent.create("b",
Reason.duplicate))));
+ }
+
+ testPerKeySequenceProcessing(
+ events.toArray(new Event[0]),
+ expectedStatuses,
+ expectedOutput,
+ duplicates,
+ EMISSION_FREQUENCY_ON_EVERY_ELEMENT,
+ INITIAL_SEQUENCE_OF_0,
+ maxResultsPerOutput,
+ DONT_PRODUCE_STATUS_ON_EVERY_EVENT);
+ }
+
@Test
public void testHandlingOfCheckedExceptions() throws
CannotProvideCoderException {
Event[] events = {