PyArrow Table: Filter rows

Viewed 2972

I have a RecordBatch from a Plasma DataStore which I can read into either a pyarrow.RecordBatch or a pyarrow.Table. I am now trying to filter out rows before converting it to pandas (to_pandas).

Is there a way of using the filter methods from the new Dataset API (that you can use on ParquetDataset) on a pyarrow.Table? This would allow me to us a filter like this:

[[('date', '=', '2020-01-01')]]

Looking at the source code both pyarrow.Table and pyarrow.RecordBatch appears to have a filter function but at least RecordBatch requires a boolean mask.

Is this possible? The reason is that the dataset contains a lot of strings (and/or categories) which are not zero-copy, so running to_pandas actually introduces significant latency and I'm only every looking for about 20% of the dataset.

Regards,
Niklas

2 Answers

This is now possible:

import pyarrow as pa

my_table = pa.Table.from_arrays(
    [pa.array(['foo', 'bar', 'foo'], pa.string())],
    names=['col1']
)

filtered_table = my_table.filter(pa.compute.equal(my_table['col1'], 'foo'))

The question above was for the equivalent of

WHERE date = '2020-01-01'

It is worth mentioning the range of PyArrow functions available for giving a wider range of selection conditions on https://arrow.apache.org/docs/python/api/compute.html#containment-tests

# import pyarrow.compute  as pc 
# WHERE test_table.name LIKE 'T%'
filtered_table = test_table.filter(pc.match_like(test_table["name"],"T%"))

As a filter returns a table structure you can chain filters one after the other.

# WHERE test_table.name LIKE 'T%' AND test_table.date_updated = '2020-01-01'
filtered_table = test_table.filter(
    pc.match_like(test_table["name"],"T%")
).(
    pc.equal(test_table["date_updated"], "2020-01-01")
)
Related