vishalmore90 opened a new issue, #40240:
URL: https://github.com/apache/beam/issues/40240

   ### What happened?
   
   
   
   **Description**
   In the `extensions/ordered` module, `ProcessorDoFn` is responsible for 
processing ordered sequences of events. When handling buffered events in 
`processBufferedEventRange()`, the function detects duplicate events (or events 
before the initial sequence) and emits them to a Dead-Letter Queue (DLQ) via 
`unprocessedEventsTupleTag`. 
   
   Currently, all duplicate elements encountered within the 
`bufferedEventsState.readRange(...)` iterator are emitted sequentially in the 
same execution bundle. If a high volume of duplicate events with the same 
sequence number is encountered, this loop will emit an unbounded number of 
elements into the bundle. This leads to unrecoverable runner exceptions (e.g., 
`CommitTooLargeException` or OOM issues on Google Cloud Dataflow), stalling the 
pipeline as the runner continuously retries and fails to commit the bundle. 
   
   There is an existing, untracked `TODO` in the codebase acknowledging this 
critical flaw.
   
   **Code Pointer:** 
   
[`sdks/java/extensions/ordered/src/main/java/org/apache/beam/sdk/extensions/ordered/ProcessorDoFn.java#L344-L361`](https://github.com/apache/beam/blob/master/sdks/java/extensions/ordered/src/main/java/org/apache/beam/sdk/extensions/ordered/ProcessorDoFn.java#L344-L361)
   
   ```java
         if (skipProcessing) {
           outputReceiver
               .get(unprocessedEventsTupleTag)
               .output(
                   KV.of(
                       processingState.getKey(),
                       KV.of(
                           eventSequence,
                           UnprocessedEvent.create(
                               bufferedEvent,
                               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.
           continue;
         }
   ```
   
   **Steps to Reproduce:**
   1. Create a pipeline utilizing the `OrderedEventProcessor` transform.
   2. Inject a high-volume stream of events (e.g., > 100,000 events) that all 
share the same `sequence` number and `key`, simulating a heavily duplicated or 
stuck upstream sender.
   3. The buffering state will collect these elements.
   4. Once processed, `processBufferedEventRange()` iterates through the entire 
chunk and blindly emits all duplicate events to the DLQ output receiver within 
a single timer/bundle execution, exceeding the runner's maximum bundle commit 
size.
   
   **Proposed Solution**
   Implement counting/batching logic combined with pagination via a stateful 
timer to manage DLQ output emissions. Define a configurable 
`MAX_EMISSIONS_PER_BUNDLE`. During `processBufferedEventRange`, keep track of 
the number of elements emitted. If the count reaches the limit, gracefully 
pause the iteration, schedule a continuation timer, and return early—deferring 
the remaining processing to the subsequent timer executions.
   
   ### Issue Priority
   
   Priority: 2 (default / most bugs should be filed as P2)
   
   
   ### Issue Priority
   
   Priority: 2 (default / most bugs should be filed as P2)
   
   ### Issue Components
   
   - [ ] Component: Python SDK
   - [x] Component: Java SDK
   - [ ] Component: Go SDK
   - [ ] Component: Typescript SDK
   - [ ] Component: IO connector
   - [ ] Component: Beam YAML
   - [ ] Component: Beam examples
   - [ ] Component: Beam playground
   - [ ] Component: Beam katas
   - [ ] Component: Website
   - [ ] Component: Infrastructure
   - [ ] Component: Spark Runner
   - [ ] Component: Flink Runner
   - [ ] Component: Prism Runner
   - [ ] Component: Twister2 Runner
   - [ ] Component: Hazelcast Jet Runner
   - [x] Component: Google Cloud Dataflow Runner


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to