[
https://issues.apache.org/jira/browse/ARROW-18400?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=17650932#comment-17650932
]
Joris Van den Bossche commented on ARROW-18400:
-----------------------------------------------
Using Alenka's script, I explored it a bit further, and noticed some things:
when reading the parquet file using the dataset reader, it is using more memory
to convert to pandas, but based on `memray` profiler, it didn't seem to take
any other code path, all the code paths just allocate twice the memory (twice
in my case, but for large dataset it might be x4 or x8 etc). So there needed
to be _something_ different about the table created from reading the parquet
file with the legacy API vs dataset API. And it seems that with the dataset
API, it is returning multiple chunks, but each of those chunks is actually a
slice of a single buffer. And in the conversion layer, there is something not
taking into account this offset into the buffer/.
Illustrating this, reading the parquet file in two ways (I have been using
{{nrows = 1_024_000 // 4}}, so the file is a bit smaller, less chunks):
{code:python}
import pyarrow.parquet as pq
import pyarrow.dataset as ds
table1 = pq.read_table("memory_testing.parquet", use_legacy_dataset=True)
dataset = ds.dataset("memory_testing.parquet", format="parquet")
table2 = dataset.to_table()
{code}
Table 1 has a single chunk, while table 2 (from reading with dataset API) has
two chunks:
{code}
>>> table1["c1"].num_chunks
1
>>> table2["c1"].num_chunks
2
{code}
Taking the first chunk of each of those, and then looking at those arrays:
{code:python}
arr1 = table1["c1"].chunk(0)
arr2 = table2["c1"].chunk(0)
>>> len(arr1)
256000
>>> len(arr2) # around half the number of rows (since there are two chunks in
>>> this table)
131072
>>> arr1.get_total_buffer_size()
110624012
>>> arr2.get_total_buffer_size() # but still using the same total memory!
110624012
{code}
So the smaller chunk of table2 is not using less memory. That is because the
two chunks of table2 are actually each a slice into the same underlying buffers:
{code:python}
>>> table2["c1"].chunk(0).buffers()[1]
<pyarrow.Buffer address=0x7fc5cc907340 size=1024004 is_cpu=True is_mutable=True>
>>> table2["c1"].chunk(1).buffers()[1] # second chunk points to same memory
>>> address and has same size as first chunk
<pyarrow.Buffer address=0x7fc5cc907340 size=1024004 is_cpu=True is_mutable=True>
>>> table2["c1"].chunk(1).offset # and the second chunk has an offset to
>>> account for that
131072
{code}
And somehow the conversion code for ListArray to numpy (which creates a numpy
array of numpy arrays, by first creating one numpy array of the flat values,
and then creating slices into that flat array) doesn't seem to take into
account this offset, and ends up converting the full parent buffer twice (in my
case twice, because of having 2 chunks, but this can grow quadratically).
---
The reason this happens for parquet and not for feather in this case, is
because the Parquet file actually consists of a single row group (and I assume
the dataset API will therefore still read that in one go, and then slice output
batches from that to return the expected batch size in the dataset API), while
the feather file already consists of multiple batches on disk (and thus doesn't
result in sliced batches in memory).
> [Python] Quadratic memory usage of Table.to_pandas with nested data
> -------------------------------------------------------------------
>
> Key: ARROW-18400
> URL: https://issues.apache.org/jira/browse/ARROW-18400
> Project: Apache Arrow
> Issue Type: Bug
> Components: Python
> Affects Versions: 10.0.1
> Environment: Python 3.10.8 on Fedora Linux 36. AMD Ryzen 9 5900 X
> with 64 GB RAM
> Reporter: Adam Reeve
> Assignee: Alenka Frim
> Priority: Critical
> Fix For: 11.0.0
>
> Attachments: test_memory.py
>
>
> Reading nested Parquet data and then converting it to a Pandas DataFrame
> shows quadratic memory usage and will eventually run out of memory for
> reasonably small files. I had initially thought this was a regression since
> 7.0.0, but it looks like 7.0.0 has similar quadratic memory usage that kicks
> in at higher row counts.
> Example code to generate nested Parquet data:
> {code:python}
> import numpy as np
> import random
> import string
> import pandas as pd
> _characters = string.ascii_uppercase + string.digits + string.punctuation
> def make_random_string(N=10):
> return ''.join(random.choice(_characters) for _ in range(N))
> nrows = 1_024_000
> filename = 'nested.parquet'
> arr_len = 10
> nested_col = []
> for i in range(nrows):
> nested_col.append(np.array(
> [{
> 'a': None if i % 1000 == 0 else np.random.choice(10000,
> size=3).astype(np.int64),
> 'b': None if i % 100 == 0 else random.choice(range(100)),
> 'c': None if i % 10 == 0 else make_random_string(5)
> } for i in range(arr_len)]
> ))
> df = pd.DataFrame({'c1': nested_col})
> df.to_parquet(filename)
> {code}
> And then read into a DataFrame with:
> {code:python}
> import pyarrow.parquet as pq
> table = pq.read_table(filename)
> df = table.to_pandas()
> {code}
> Only reading to an Arrow table isn't a problem, it's the to_pandas method
> that exhibits the large memory usage. I haven't tested generating nested
> Arrow data in memory without writing Parquet from Pandas but I assume the
> problem probably isn't Parquet specific.
> Memory usage I see when reading different sized files on a machine with 64 GB
> RAM:
> ||Num rows||Memory used with 10.0.1 (MB)||Memory used with 7.0.0 (MB)||
> |32,000|362|361|
> |64,000|531|531|
> |128,000|1,152|1,101|
> |256,000|2,888|1,402|
> |512,000|10,301|3,508|
> |1,024,000|38,697|5,313|
> |2,048,000|OOM|20,061|
> |4,096,000| |OOM|
> With Arrow 10.0.1, memory usage approximately quadruples when row count
> doubles above 256k rows. With Arrow 7.0.0 memory usage is more linear but
> then quadruples from 1024k to 2048k rows.
> PyArrow 8.0.0 shows similar memory usage to 10.0.1 so it looks like something
> changed between 7.0.0 and 8.0.0.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)