apache / apache/arrow

pd.read_parquet using filters consumes too much memory

Open
#31,771 4 comments 0 reactions 0 assignees View on GitHub
Component: Python Type: bug
Dominant language
C++
Stars
17.1k
Forks
4.3k
Avg merge
3d 13h
Merged PRs (30d)
88

Description

Hello!

 

I have found that pyarrow versions **>= 4.0.1** use more than **2x** memory (RSS) when trying to read a parquet using file-level filters. Using the following dataset:
```java

import pandas as pd
import numpy as np

a = np.random.randint(1,50,(4_000_000,4))

df = pd.DataFrame(a, columns=['A','B','C','D']).to_parquet("test.pq", index=False)
```
and the reader script ({**}read_with_filters.py{**})
```java

import pyarrow as pa
import pandas as pd

print(f"pyarrow version: {pa.__version__}")
print(f"pandas version: {pd.__version__}")

tmp = pd.read_parquet("test.pq", engine='pyarrow', use_legacy_dataset=False, filters=[("B","=",10)])

print(tmp.shape)
```
I get:

 

**Python 3.8.13 (conda), pyarrow 1.0.1 (pip), pandas 1.4.2 (pip)**
```java

gtime -f "RSS (Kb): %M | user (sec): %U | system (sec): %S | real (sec) : %e" python read_with_filters.py
pyarrow version: 1.0.1
pandas version: 1.4.2
(81833, 4)
RSS (Kb): 84876 | user (sec): 0.87 | system (sec): 0.32 | real (sec) : 1.32
```
**Python 3.8.13 (conda), pyarrow 4.0.1 (pip), pandas 1.4.2 (pip)**
```java

gtime -f "RSS (Kb): %M | user (sec): %U | system (sec): %S | real (sec) : %e" python read_with_filters.py
pyarrow version: 4.0.1
pandas version: 1.4.2
(81833, 4)
RSS (Kb): 172816 | user (sec): 0.77 | system (sec): 0.24 | real (sec) : 0.72
```
**Python 3.8.13 (conda), pyarrow 7.0.0 (pip), pandas 1.4.2 (pip)**
```java

gtime -f "RSS (Kb): %M | user (sec): %U | system (sec): %S | real (sec) : %e" python read_with_filters.py
pyarrow version: 7.0.0
pandas version: 1.4.2
(81833, 4)
RSS (Kb): 240112 | user (sec): 0.71 | system (sec): 0.22 | real (sec) : 0.82
```
 

It is more evident when using a larger dataset. However, my personal computer hangs when trying to read larger datasets using **4.0.1** and {**}7.0.0{**}. (That should be a separate issue)

 

It is worth mentioning that you see a relative the same memory usage when I removed the **filters** keyword.

**Python 3.8.13 (conda), pyarrow 1.0.1 (pip), pandas 1.4.2 (pip)**
```java

gtime -f "RSS (Kb): %M | user (sec): %U | system (sec): %S | real (sec) : %e" python read_with_filters.py
pyarrow version: 1.0.1
pandas version: 1.4.2
(4000000, 4)
RSS (Kb): 331424 | user (sec): 0.89 | system (sec): 0.39 | real (sec) : 1.07
```
{**}Python 3.8.13 (conda), pyarrow 4.0.1 (pip), pandas 1.4.2 (pip){**}{**}{{**}}
```java

gtime -f "RSS (Kb): %M | user (sec): %U | system (sec): %S | real (sec) : %e" python read_with_filters.py
pyarrow version: 4.0.1
pandas version: 1.4.2
(4000000, 4)
RSS (Kb): 405916 | user (sec): 0.81 | system (sec): 0.42 | real (sec) : 0.81
```
{**}Python 3.8.13 (conda), pyarrow 7.0.0 (pip), pandas 1.4.2 (pip){**}{**}{{**}}
```java

gtime -f "RSS (Kb): %M | user (sec): %U | system (sec): %S | real (sec) : %e" python read_with_filters.py
pyarrow version: 7.0.0
pandas version: 1.4.2
(4000000, 4)
RSS (Kb): 364152 | user (sec): 0.78 | system (sec): 0.45 | real (sec) : 1.27
```
Thank you,

Luis

**Environment**: Hardware Overview:
Model Name: MacBook Pro
Model Identifier: MacBookPro12,1
Processor Name: Dual-Core Intel Core i5
Processor Speed: 2.7 GHz
Number of Processors: 1
Total Number of Cores: 2
L2 Cache (per Core): 256 KB
L3 Cache: 3 MB
Hyper-Threading Technology: Enabled
Memory: 8 GB

System Software Overview:
System Version: macOS 10.15.7 (19H1217)
Kernel Version: Darwin 19.6.0
Boot Volume: Macintosh HD
Boot Mode: Normal

**Reporter**: [Luis E Pastrana](https://issues.apache.org/jira/browse/ARROW-16391)

**Note**: *This issue was originally created as [ARROW-16391](https://issues.apache.org/jira/browse/ARROW-16391). Please see the [migration documentation](https://github.com/apache/arrow/issues/14542) for further details.*

Contributor guide

Open the contributing guide

Research direction

Start by running the provided read_with_filters.py reproduction with the listed Python, pandas, and pyarrow versions, then compare filtered and unfiltered reads using the generated test.pq dataset. Trace the file-level filter path in the Arrow repository and measure RSS while preserving the reported result shape; done means filtered reads no longer consume substantially more memory than expected.

Written by the indexing model from the issue text.

Assessment

Tech stack
pandas, python
Domain
data-engineering
Issue type
Bug
Difficulty
4/5
Estimated time
3-5 days
Activity status
Stale
Clarity
Mostly clear
Newbie friendliness
35/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.