# Conditionally read a subset of Arrow-data

**URL:** <https://discourse.julialang.org/t/conditionally-read-a-subset-of-arrow-data/72122>\
**Category:** Data\
**Tags:** dataframes, arrow\
**Created:** [November 26, 2021, 3:53pm UTC](https://discourse.julialang.org/t/conditionally-read-a-subset-of-arrow-data/72122 "2021-11-26T15:53:30Z")\
**Posts on this page:** 10\
**Page:** 1

<div class="post-metadata">

**Author:** ![fbanning](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/fbanning/32/14972_2.png) [@fbanning](https://discourse.julialang.org/u/fbanning)\
**Post date:** [November 26, 2021, 3:53pm UTC](https://discourse.julialang.org/t/conditionally-read-a-subset-of-arrow-data/72122/1 "2021-11-26T15:53:30Z")

</div>

Maybe my head is stuck in a loop right now but I have a question about dealing with arrow data.

I have some data stored in a few large arrow files. I currently read the data via `Arrow.Table("file.arrow") |> DataFrame` and then `filter` the resulting DataFrame to reduce it to the subset of data that I want to work with. However, this is very inefficient: first read all the data, then reduce it again.

How can I search through the Arrow file while accessing it and directly filter the data stream? Thanks for any pointers, I definitely have the feeling that I’m missing something crucial here.

---

<div class="post-metadata">

**Author:** ![trahflow](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/trahflow/32/30585_2.png) [@trahflow](https://discourse.julialang.org/u/trahflow)\
**Post date:** [November 26, 2021, 4:03pm UTC](https://discourse.julialang.org/t/conditionally-read-a-subset-of-arrow-data/72122/2 "2021-11-26T16:03:51Z")

</div>

Hi @fbanning ,  
the [documentation](https://arrow.juliadata.org/dev/manual/#Arrow.Table) on `Arrow.Table` actually has quite some information on this.  
In particular:

> `Arrow.Table` first [“mmapped”](https://en.wikipedia.org/wiki/Mmap) the `data.arrow` file, which is an important technique for dealing with data larger than available RAM on a system. By “mmapping” a file, the OS doesn’t actually load the entire file contents into RAM at the same time, but file contents are “swapped” into RAM as different regions of a file are requested.

and:

> `Tables.datavaluerows(Arrow.Table(file)) |> @map(...) |> @filter(...) |> DataFrame` : use [`Query.jl` 's](https://www.queryverse.org/Query.jl/stable/standalonequerycommands/) row-processing utilities to map, group, filter, mutate, etc. directly over arrow data.

(I haven’t tried it myself, but the way I interpret that example is that data is loaded only row-wise into memory)

---

<div class="post-metadata">

**Author:** ![nalimilan](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/nalimilan/32/147_2.png) [@nalimilan](https://discourse.julialang.org/u/nalimilan)\
**Post date:** [November 26, 2021, 4:20pm UTC](https://discourse.julialang.org/t/conditionally-read-a-subset-of-arrow-data/72122/3 "2021-11-26T16:20:12Z")

</div>

Note that `DataFrame(Arrow.Table(...))` creates a `DataFrame` whose columns are special Arrow types which already benefit from mmapping, without making any copy. So AFAIK calling `filter` on the result is already the fastest and most memory-efficient solution to extract a subset of data. Query .jl is only needed if you want to apply an intermediate processing step via `@map` before taking a subset, and without making an additional copy.

The key point is that Arrow is designed to be read efficiently from disk without any copy, contrary to e.g. CSV.

---

<div class="post-metadata">

**Author:** ![fbanning](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/fbanning/32/14972_2.png) [@fbanning](https://discourse.julialang.org/u/fbanning)\
**Post date:** [November 26, 2021, 4:42pm UTC](https://discourse.julialang.org/t/conditionally-read-a-subset-of-arrow-data/72122/4 "2021-11-26T16:42:48Z")

</div>

> [@trahflow](#):
>
> Tables.datavaluerows(Arrow.Table(file)) |\> @map(…) |\> @filter(…) |\> DataFrame

I think this is the part that I have missed. Will look into how to work with this, thanks.  
For comparison, this is what I’m currently doing `Arrow.Table(file) |> DataFrame |> filter(...)`. I just didn’t know that I can simply switch the pieces of the pipeline - very nice if that’s how it works. 🙂

> [@nalimilan](#):
>
> So AFAIK calling `filter` on the result is already the fastest and most memory-efficient solution to extract a subset of data.

So there’s no need to do as described above? 🤔

Edit: Ah, so the intermediate steps are good to avoid making an additional copy when filtering, is that correct?

---

<div class="post-metadata">

**Author:** ![nalimilan](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/nalimilan/32/147_2.png) [@nalimilan](https://discourse.julialang.org/u/nalimilan)\
**Post date:** [November 26, 2021, 4:59pm UTC](https://discourse.julialang.org/t/conditionally-read-a-subset-of-arrow-data/72122/5 "2021-11-26T16:59:59Z")

</div>

> [@fbanning](#):
>
> So there’s no need to do as described above? 🤔
> 
> Edit: Ah, so the intermediate steps are good to avoid making an additional copy when filtering, is that correct?

If by “intermediate steps” you mean `@map`, then no, it’s not really needed to avoid copies. Generally you can just call `filter` first and then apply any transformations you need to the (smaller) resulting data set. Or if you want to avoid doing any copy at all, you can do `filter(..., view=true)`. Query.jl is nice but I don’t think it’s really more efficient here.

---

<div class="post-metadata">

**Author:** ![bkamins](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/bkamins/32/208538_2.png) [@bkamins](https://discourse.julialang.org/u/bkamins)\
**Post date:** [November 26, 2021, 5:13pm UTC](https://discourse.julialang.org/t/conditionally-read-a-subset-of-arrow-data/72122/6 "2021-11-26T17:13:10Z")

</div>

In general what you should keep in mind that all DataFrames.jl operations are eager. They get materialized immediately (unless you ask for a view when no copying happens).

So for example if you do `filter` in DataFrames.jl without `view=true` it will copy data to a new data frame.

---

<div class="post-metadata">

**Author:** ![nilshg](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/nilshg/32/2283_2.png) [@nilshg](https://discourse.julialang.org/u/nilshg)\
**Post date:** [November 26, 2021, 5:37pm UTC](https://discourse.julialang.org/t/conditionally-read-a-subset-of-arrow-data/72122/7 "2021-11-26T17:37:34Z")

</div>

I think when I last did something like this I used:

```julia
df = DataFrame(Arrow.Table("mytable.arrow"))
df = @view df[(df.col1 .< 5) .& (df.col2 .== "a"), :]
df = copy!(df)

```

and it seemed to me just from the system monitor that this kept resource utilisation at an acceptable level.

---

<div class="post-metadata">

**Author:** ![bkamins](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/bkamins/32/208538_2.png) [@bkamins](https://discourse.julialang.org/u/bkamins)\
**Post date:** [November 26, 2021, 5:39pm UTC](https://discourse.julialang.org/t/conditionally-read-a-subset-of-arrow-data/72122/8 "2021-11-26T17:39:52Z")

</div>

`copy!` is not defined for `AbstractDataFrame`, but indeed `@view` would be efficient.

---

<div class="post-metadata">

**Author:** ![nilshg](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/nilshg/32/2283_2.png) [@nilshg](https://discourse.julialang.org/u/nilshg)\
**Post date:** [November 26, 2021, 6:20pm UTC](https://discourse.julialang.org/t/conditionally-read-a-subset-of-arrow-data/72122/9 "2021-11-26T18:20:13Z")

</div>

Sorry typo, meant without the bang!

---

<div class="post-metadata">

**Author:** ![quinnj](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/quinnj/32/11_2.png) [@quinnj](https://discourse.julialang.org/u/quinnj)\
**Post date:** [December 2, 2021, 4:42am UTC](https://discourse.julialang.org/t/conditionally-read-a-subset-of-arrow-data/72122/10 "2021-12-02T04:42:32Z")

</div>

I appreciate all the replies above and they are mostly spot-on.

I imagine part of the confusion here is probably having worked with arrow data in other language implementations. Most other language implementations I’m aware of include the ability to pass in various “table operations” when reading the arrow data directly; i.e. filter, aggregate, etc. And those operations are “pushed down” and applied while reading the data so that the “table” object that is returned respects those operations.

In the Julia implementation, we could _potentially_ support something like that, but as mentioned above, the current “read” operation is essentially just an mmap of the raw data along with metadata processing, which is extremely cheap. You can even wrap these “arrow arrays/columns” in a DataFrame and get access to the whole family of operations supported there very easily. And this isn’t limited to DataFrames; as also mentioned, you can use the Query.jl row-processing framework, or the operations in DataFramesMeta; TypedTables or IndexedTables are two other “table” frameworks. There’s also the TableOperations.jl package, which supports some common operations like column selection, mapping, and filtering.

All these packages support operating on any generic “table” input, including `Arrow.Table`. And due to the immutable nature of arrow data, unnecessary copies of the data can be avoided by default.

So essentially we don’t support these kind of operations while reading in Arrow.jl because there are quite a few options for processing the data post-read efficiently. In this way, all the code in the Arrow.jl package is purely focused on reading/writing the arrow format and ends up being a significantly smaller/simpler implementation compared to other language implementations with similar format coverage.

Hope that helps!
