Skip to content

[Bug]: Unrecoverable Runner Errors (CommitTooLargeException) When Emitting High Volume of Duplicate/Unprocessed Events in OrderedEventProcessor #40240

Description

@vishalmore90

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:

  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 Components

  • Component: Python SDK
  • 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
  • Component: Google Cloud Dataflow Runner

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions