PyArrow filter table based on calculcated column

Viewed 92

so I am trying to calculate the days between the date column and today. And filter table where the diff is more than 5. In spark, you could do something like

datediff(lit(today),df.date) > 5

In pyarrow what I am doing is following

dates = pa.compute.days_between(df['date'], today)
df = df.append_column('days_diff' , dates)
filtered = df.filter(pc.field('days_diff') > 5)
df = df.remove_column('days_diff')

But this creates a new column which is memory overhead. Is it possible to have calculated column for filter only?

3 Answers

The dataset api can filter (and project) a table without intermediate arrays.

import pyarrow.dataset as ds

expr = ds.field('date') < (today - timedelta(days=5))
ds.dataset(table).scanner(filter=expr).to_table()

Given all pyarrow compute functions work with arrays as input/output, there isn't much you can do about the memory overhead of creating a new array.

At the API level, you can avoid appending a new column to your table, but it's not going to save any memory:

dates_diff = pa.compute.days_between(table['date'], today)
dates_filter = pa.compute.greater(dates_diff, 5)
filtered_table = pa.compute.filter(table, dates_filter)

If memory is really an issue you can do the filtering in small batches:

sub_tables = []
for batch in table.to_batches():
    dates_diff = pa.compute.days_between(batch['date'], today)
    indices = pa.compute.greater(dates_diff, 5)
    sub_table = pa.compute.filter(batch, indices)
    sub_tables.append(sub_table)
pa.Table.from_batches(sub_tables, schema=table.schema)

But this is not very practical, and only works for simple operations. Also if you don't have enough memory to allocate a couple of new columns, it means you have very little wiggle room and won't be able to achieve much.

Lastly the overhead of allocating an array is offset by the fact that operation on arrays can be parallelised and greatly optimised in arrow.

The Table.filter() method can accept a boolean expression since pyarrow 9.0.0 (just released, August 2022), in addition to an actual materialized boolean array.
(and so you can now do this directly with a Table without the need to wrap the table in a dataset, as shown in the answer of A. Coady https://stackoverflow.com/a/73281960/653364)

And based on your example code in the question, you are actually already using that. But you can construct a more complex expression that avoids the need to create this intermediate "days_diff" column:

import datetime
import pyarrow.compute as pc

expr = pc.days_between(pc.field("date"), pc.scalar(today)) > 5
# <pyarrow.compute.Expression (days_between(date, 2022-08-09) > 5)>

df.filter(expr)

In used the expression as you had in your question. But you could check which of the two expressions gives the best performance (the above expression or pc.field('date') < (today - timedelta(days=5)), although if "date" is a timestamp column, it might give different results)

Related