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]
