# Is parallel Arrow data querying possible?

**URL:** <https://discourse.julialang.org/t/is-parallel-arrow-data-querying-possible/60462>\
**Category:** Data\
**Tags:** parallel, arrow\
**Created:** [May 3, 2021, 1:39pm UTC](https://discourse.julialang.org/t/is-parallel-arrow-data-querying-possible/60462 "2021-05-03T13:39:08Z")\
**Posts on this page:** 5\
**Page:** 1

<div class="post-metadata">

**Author:** ![Sami](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/sami/32/20726_2.png) [@Sami](https://discourse.julialang.org/u/Sami)\
**Post date:** [May 3, 2021, 1:39pm UTC](https://discourse.julialang.org/t/is-parallel-arrow-data-querying-possible/60462/1 "2021-05-03T13:39:08Z")

</div>

I have 16 cores on my machine and queries seem to utilize only one of them. I would like to be 16 times faster with my queries - so is it possible to run parallel queries on Arrow data (read in with Arrow.jl)??

For example, is there some kind of macro-magic (that is how I see macros currently as a newcomer to Julia world) that could be used? Like done [here](https://discourse.julialang.org/t/parallel-dataframe-processing/29397)

---

<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:** [May 3, 2021, 2:36pm UTC](https://discourse.julialang.org/t/is-parallel-arrow-data-querying-possible/60462/2 "2021-05-03T14:36:49Z")

</div>

Arrow.jl returns an `Arrow.Table`, which is made up of concrete subtypes of `ArrowVector`, which implement the `AbstractArray` interface. So, for example, when doing `at = Arrow.Table(file); col1 = at.col1`, `col1` is an object that is a “view” into the raw arrow data in `file`. Indexing like `col1[1]` computes the exact byte offset of the value in the raw arrow data and returns the data.

All that is to say, there’s nothing automatic in Arrow.jl to utilize multiple cores, but you’re completely free and flexible to do parallel/concurrent processing however you’d like. You could spawn multithreaded tasks to operate over an array; you could assign different processors to process arrays separately, etc.

Currently, DataFrames.jl defines some operations to process data in parallel using multiple threads when the conditions are right (i.e. when it would actually benefit performance on large datasets), and I believe the Transducers.jl framework as some nice parallelism workflows (cc: @tkf). But yeah, it really depends on your workflow and what you’re trying to do.

---

<div class="post-metadata">

**Author:** ![djholiver](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/djholiver/32/50470_2.png) [@djholiver](https://discourse.julialang.org/u/djholiver)\
**Post date:** [May 3, 2021, 2:45pm UTC](https://discourse.julialang.org/t/is-parallel-arrow-data-querying-possible/60462/3 "2021-05-03T14:45:59Z")

</div>

Hi,

I actually use the arrow format for a reasonably heavy analytics layer and process across both rows and columns in a mutithreaded/parallel way.

Recall that you have access to the actual indices of a column (one of the genuinely incredible aspects or arrow imo) and can therefore access a column in a partitioned thread approach like “Threads.@threads” to query over a range.

Regards

---

<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:** [May 3, 2021, 6:06pm UTC](https://discourse.julialang.org/t/is-parallel-arrow-data-querying-possible/60462/4 "2021-05-03T18:06:47Z")

</div>

Could you give an example of that? Sounds interesting

---

<div class="post-metadata">

**Author:** ![Sami](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/sami/32/20726_2.png) [@Sami](https://discourse.julialang.org/u/Sami)\
**Post date:** [May 4, 2021, 7:30am UTC](https://discourse.julialang.org/t/is-parallel-arrow-data-querying-possible/60462/5 "2021-05-04T07:30:32Z")

</div>

Thank you @quinnj for simplifying Arrow inner workings, and thanks @djholiver for telling me that it could be done. However, I am new to Julia and currently just capable of utilizing the ready packages/apis (like DataFrames.jl or DataFramesMera.jl). If there would be a easier way to utilizes the modern hardware to the max, I believe that would help in gaining attraction.

I wonder whether there are pointers to record batches also and could those record batches be processed in map-reduce fashion (and also in zero-copy fashion)? Furthermore, could those pointers be provided to the GPU and let it to the processing?

Interesting post about [data parallelism](https://juliafolds.github.io/data-parallelism/tutorials/quick-introduction/)
