[
https://issues.apache.org/jira/browse/ARROW-17593?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=17599216#comment-17599216
]
Weston Pace commented on ARROW-17593:
-------------------------------------
The dataset writer already has this decoupled with {{min_rows_per_group}} which
should handle what you are describing. I think the write path is in ok shape.
I'm mostly worried about the "start with datasets but do further downstream
work in-memory in python" path.
> [C++] Try and maintain input shape in Acero
> -------------------------------------------
>
> Key: ARROW-17593
> URL: https://issues.apache.org/jira/browse/ARROW-17593
> Project: Apache Arrow
> Issue Type: Bug
> Components: C++
> Reporter: Weston Pace
> Priority: Major
>
> Data is scanned in large chunks based on the format. For example, CSV scans
> chunks based on a chunk_size while parquet scans entire row groups.
> Then, upon entry into Acero, these chunks are sliced into morsels (~L3 size)
> for parallelism and batches (~L1-L2 size) for cache efficient processing.
> However, the way it is currently done, means that the output of Acero is a
> stream of tiny batches. This is somewhat undesirable in many cases.
> For example, if a pyarrow user calls pq.read_table they might expect to get
> one batch per row group. If they were to turn around and write out that
> table to a new parquet file then either they end up with a non-ideal parquet
> file (tiny row groups) or they are forced to concatenate the batches (which
> is an allocation + copy).
> Even if the user is doing their own streaming processing (e.g. in pyarrow)
> these small batch sizes are undesirable as the overhead of python means that
> streaming processing should be done in larger batches.
> Instead, there should be a configurable max_batch_size, independent of row
> group size and morsel size, which is configurable, and quite large by default
> (1Mi or 64Mi rows). This control exists for users that want to do their own
> streaming processing and need to be able to tune for RAM usage.
> Acero will read in data based on the format, as it does today (e.g. CSV chunk
> size, row group size). If the source data is very large (bigger than
> max_batch_size) it will be sliced. From that point on, any morsels or
> batches should simply be views into this larger output batch. For example,
> when doing a projection to add a new column, we should allocate a
> max_batch_size array and then populate it over many runs of the project node.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)