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
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:
- Create a pipeline utilizing the
OrderedEventProcessor transform.
- 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.
- The buffering state will collect these elements.
- 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 Components
What happened?
Description
In the
extensions/orderedmodule,ProcessorDoFnis responsible for processing ordered sequences of events. When handling buffered events inprocessBufferedEventRange(), the function detects duplicate events (or events before the initial sequence) and emits them to a Dead-Letter Queue (DLQ) viaunprocessedEventsTupleTag.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.,CommitTooLargeExceptionor 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
TODOin the codebase acknowledging this critical flaw.Code Pointer:
sdks/java/extensions/ordered/src/main/java/org/apache/beam/sdk/extensions/ordered/ProcessorDoFn.java#L344-L361Steps to Reproduce:
OrderedEventProcessortransform.sequencenumber andkey, simulating a heavily duplicated or stuck upstream sender.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. DuringprocessBufferedEventRange, 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 Components