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 = {

Reply via email to