apache / apache/beam

Use WorkItemCommitRequest protobuf fields to signal that a WorkItem needs to be broken up

Open
#19,956 0 comments 0 reactions 0 assignees View on GitHub
dataflow improvement P3 runners
Dominant language
Java
Stars
8.7k
Forks
4.7k
Avg merge
1d 20h
Merged PRs (30d)
196

Description

****Background:****

When a WorkItemCommitRequest is generated that's bigger than the permitted size (\> ~180 MB), a KeyCommitTooLargeException is logged (_not thrown_) and the request is still sent to the service.  The service rejects the commit, but breaks up input messages that were bundled together and adds them to new, smaller work items that will later be pulled and re-tried - likely without generating another commit that is too large.

When a WorkItemCommitRequest is generated that's too large to be sent back to the service (\> 2 GB), a KeyCommitTooLargeException is thrown and nothing is sent back to the service.

 

****Proposed Improvement****

In both cases, prevent the doomed, large commit item from being sent back to the service.  Instead send flags in the commit request signaling that the current work item led to a commit that is too large and the work item should be broken up.  

Imported from Jira [BEAM-8554](https://issues.apache.org/jira/browse/BEAM-8554). Original Jira may contain additional context.
Reported by: stevekoonce.

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.