Thomas Groh created BEAM-4228:
---------------------------------

             Summary: The FlinkRunner shouldn't require all of the values for a 
key to fit in memory
                 Key: BEAM-4228
                 URL: https://issues.apache.org/jira/browse/BEAM-4228
             Project: Beam
          Issue Type: New Feature
          Components: runner-flink
            Reporter: Thomas Groh


The use of a reducer that adds all of the elements that it consumes to a list 
is the primary way in which this occurs - if instead, we produce a filtered 
iterable, or a collection of filtered iterables, we can lazily iterate over all 
of the contained elements without having to buffer all of the elements.

 

For an example of where this occurs, see {{Concatenate}} in  
{{FlinkBatchPortablePipelineTranslator}}.



--
This message was sent by Atlassian JIRA
(v7.6.3#76005)

Reply via email to