# Efficiently merge multiple .arrow files (12 files totaling 2.5 TB) into a single DataFrame for effective data analysis while minimizing memory overhead

**URL:** <https://discourse.julialang.org/t/efficiently-merge-multiple-arrow-files-12-files-totaling-2-5-tb-into-a-single-dataframe-for-effective-data-analysis-while-minimizing-memory-overhead/104110>\
**Category:** Performance\
**Tags:** question, dataframes\
**Created:** [September 21, 2023, 2:12pm UTC](https://discourse.julialang.org/t/efficiently-merge-multiple-arrow-files-12-files-totaling-2-5-tb-into-a-single-dataframe-for-effective-data-analysis-while-minimizing-memory-overhead/104110 "2023-09-21T14:12:15Z")\
**Posts on this page:** 15\
**Page:** 1

<div class="post-metadata">

**Author:** ![ujjwal\_singh](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/ujjwal_singh/32/53073_2.png) [@ujjwal\_singh](https://discourse.julialang.org/u/ujjwal_singh)\
**Post date:** [September 21, 2023, 2:12pm UTC](https://discourse.julialang.org/t/efficiently-merge-multiple-arrow-files-12-files-totaling-2-5-tb-into-a-single-dataframe-for-effective-data-analysis-while-minimizing-memory-overhead/104110/1 "2023-09-21T14:12:15Z")

</div>

Could someone please guide me on how to merge multiple large .arrow files (12 files totaling 2.5 TB) into a single DataFrame with minimal memory overhead for big data analysis? I have a total of 1 TB of RAM; can I manage data analysis with the existing RAM? Here’s a two-line snippet of my code:

filenames = [“ct1\_first.arrow”, “ct1\_snd.arrow”]  
append!(filenames, [“ct$(i).arrow” for i in 2:11])

---

<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:** [September 21, 2023, 2:20pm UTC](https://discourse.julialang.org/t/efficiently-merge-multiple-arrow-files-12-files-totaling-2-5-tb-into-a-single-dataframe-for-effective-data-analysis-while-minimizing-memory-overhead/104110/2 "2023-09-21T14:20:21Z")

</div>

you can do (pseudocode):

```julia
df = DataFrame()
for arrow_table in vector_of_arrow_tables
    append!(df, arrow_table)
end

```

This should be efficient enough hopefully.

---

<div class="post-metadata">

**Author:** ![ujjwal\_singh](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/ujjwal_singh/32/53073_2.png) [@ujjwal\_singh](https://discourse.julialang.org/u/ujjwal_singh)\
**Post date:** [September 21, 2023, 2:21pm UTC](https://discourse.julialang.org/t/efficiently-merge-multiple-arrow-files-12-files-totaling-2-5-tb-into-a-single-dataframe-for-effective-data-analysis-while-minimizing-memory-overhead/104110/3 "2023-09-21T14:21:34Z")

</div>

One more thing: all the .arrow files have the same columns, as shown in the attached screenshot

 ![gg](https://global.discourse-cdn.com/julialang/original/3X/2/5/252b09fd1be3fba17ead1a01a6b8eeb62dfe9036.png)

---

<div class="post-metadata">

**Author:** ![Palli](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/palli/32/3380_2.png) [@Palli](https://discourse.julialang.org/u/Palli)\
**Post date:** [September 21, 2023, 3:28pm UTC](https://discourse.julialang.org/t/efficiently-merge-multiple-arrow-files-12-files-totaling-2-5-tb-into-a-single-dataframe-for-effective-data-analysis-while-minimizing-memory-overhead/104110/4 "2023-09-21T15:28:19Z")

</div>

I think Polar.jl might help, about to be registered (unless suggested name-change goes though; you can still install with that name):

> <https://github.com/JuliaRegistries/General/pull/91065>
>
> \- Registering package: Polars
> \- Repository: https://github.com/Pangoraw/Polars.j…l
> \- Created by: @Pangoraw
> \- Version: v0.1.0
> \- Commit: 1860f4dd5a8750ffe2ab78268b72e7008cfe169d
> \- Reviewed by: @Pangoraw
> \- Reference: https://github.com/Pangoraw/Polars.jl/commit/1860f4dd5a8750ffe2ab78268b72e7008cfe169d#commitcomment-129076209
> \- Description: :bear: Julia wrapper around the polars library

It has out-of-core (and Arrow) support. There’s also Arrow.jl.

> [@bkamins](#):
>
> you can do (pseudocode):
> 
> ```julia
> df = DataFrame()
> 
> ```

I wasn’t sure I thought it (and Arrows.jl? nor JuliaDB.jl) doesn’t have out-of-core? Am I right? Polars has been faster in some cases, and doesn’t run out of memory because of that I believe as fast?

> [@bkamins](#):
>
> This should be efficient enough hopefully.

Why do you think that then, or working at all? I guess you mean relying on paging/swapping. Likely ok, and efficient enough for a one-time thing this seems to be.

[FYI: also intriguing wrapper already registered RustRegex.jl, also wrapping fast Rust code; and other packages in the pipline e.g. [Julia] ObjectSystem.jl.]

---

<div class="post-metadata">

**Author:** ![ujjwal\_singh](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/ujjwal_singh/32/53073_2.png) [@ujjwal\_singh](https://discourse.julialang.org/u/ujjwal_singh)\
**Post date:** [September 21, 2023, 4:35pm UTC](https://discourse.julialang.org/t/efficiently-merge-multiple-arrow-files-12-files-totaling-2-5-tb-into-a-single-dataframe-for-effective-data-analysis-while-minimizing-memory-overhead/104110/5 "2023-09-21T16:35:49Z")

</div>

![p1](https://global.discourse-cdn.com/julialang/original/3X/7/3/731964535bd5d3986f1daa8e0cb4d23e514ec8f1.png)  
This code is causing memory overload. Could you please suggest a more efficient solution that uses less memory, preferably within 1 TB of RAM?

---

<div class="post-metadata">

**Author:** ![ufechner7](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/ufechner7/32/51363_2.png) [@ufechner7](https://discourse.julialang.org/u/ufechner7)\
**Post date:** [September 21, 2023, 4:43pm UTC](https://discourse.julialang.org/t/efficiently-merge-multiple-arrow-files-12-files-totaling-2-5-tb-into-a-single-dataframe-for-effective-data-analysis-while-minimizing-memory-overhead/104110/6 "2023-09-21T16:43:24Z")

</div>

On which OS are you? If you are on Linux, did you try zram?

> **[How to Configure ZRAM on Your Ubuntu Computer - Make Tech Easier](https://www.maketecheasier.com/configure-zram-ubuntu/)**
>
> ZRAM is a great solution to trade some CPU horsepower to gain more RAM. Here is how you can configure ZRAM in Ubuntu for better performance.

---

<div class="post-metadata">

**Author:** ![ujjwal\_singh](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/ujjwal_singh/32/53073_2.png) [@ujjwal\_singh](https://discourse.julialang.org/u/ujjwal_singh)\
**Post date:** [September 21, 2023, 4:44pm UTC](https://discourse.julialang.org/t/efficiently-merge-multiple-arrow-files-12-files-totaling-2-5-tb-into-a-single-dataframe-for-effective-data-analysis-while-minimizing-memory-overhead/104110/7 "2023-09-21T16:44:54Z")

</div>

I am using Debian 11 OS.

---

<div class="post-metadata">

**Author:** ![ufechner7](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/ufechner7/32/51363_2.png) [@ufechner7](https://discourse.julialang.org/u/ufechner7)\
**Post date:** [September 21, 2023, 4:48pm UTC](https://discourse.julialang.org/t/efficiently-merge-multiple-arrow-files-12-files-totaling-2-5-tb-into-a-single-dataframe-for-effective-data-analysis-while-minimizing-memory-overhead/104110/8 "2023-09-21T16:48:29Z")

</div>

So try to enable zram and play with the settings, e.g. compression algorithm and size of compressed memory…

---

<div class="post-metadata">

**Author:** ![tbeason](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/tbeason/32/15898_2.png) [@tbeason](https://discourse.julialang.org/u/tbeason)\
**Post date:** [September 21, 2023, 5:57pm UTC](https://discourse.julialang.org/t/efficiently-merge-multiple-arrow-files-12-files-totaling-2-5-tb-into-a-single-dataframe-for-effective-data-analysis-while-minimizing-memory-overhead/104110/9 "2023-09-21T17:57:02Z")

</div>

Two things:

- make sure you are using the `heap-size-hint` flag when you start Julia
- look into using DuckDB

---

<div class="post-metadata">

**Author:** ![ujjwal\_singh](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/ujjwal_singh/32/53073_2.png) [@ujjwal\_singh](https://discourse.julialang.org/u/ujjwal_singh)\
**Post date:** [September 22, 2023, 6:30am UTC](https://discourse.julialang.org/t/efficiently-merge-multiple-arrow-files-12-files-totaling-2-5-tb-into-a-single-dataframe-for-effective-data-analysis-while-minimizing-memory-overhead/104110/10 "2023-09-22T06:30:00Z")

</div>

Your suggestion requires 637GB for processing a single file. I don’t think processing 12 files will fit within 1TB. The screenshot of the code is…

 ![jds2](https://global.discourse-cdn.com/julialang/original/3X/3/4/34c85019a6aed4c293a3defb1c8c71a65069f1cc.png)

---

<div class="post-metadata">

**Author:** ![DrChainsaw](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/drchainsaw/32/8497_2.png) [@DrChainsaw](https://discourse.julialang.org/u/DrChainsaw)\
**Post date:** [September 22, 2023, 6:36am UTC](https://discourse.julialang.org/t/efficiently-merge-multiple-arrow-files-12-files-totaling-2-5-tb-into-a-single-dataframe-for-effective-data-analysis-while-minimizing-memory-overhead/104110/11 "2023-09-22T06:36:02Z")

</div>

Isnt it so that Arrow.jl uses mmap:ed arrays, and those materialize in ram when concatenated into the big DataFrame?

I have tried using an ArrayPartition from [RecursiveArrayTools.jl](https://github.com/SciML/RecursiveArrayTools.jl) to try to prevent this column by column.

In my case the extra compiletime overhead made me abandon the attempt before getting to the bottom of whether it helped with memory to any significant degree in practice. On problem could for example be that almost any operation on the dataframe would cause the array to materialize.

If you have access to multiple machines there is [GitHub - JuliaParallel/DTables.jl: Distributed table structures and data manipulation operations built on top of Dagger.jl](https://github.com/JuliaParallel/DTables.jl) although I have never used it.

---

<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:** [September 22, 2023, 6:46am UTC](https://discourse.julialang.org/t/efficiently-merge-multiple-arrow-files-12-files-totaling-2-5-tb-into-a-single-dataframe-for-effective-data-analysis-while-minimizing-memory-overhead/104110/12 "2023-09-22T06:46:38Z")

</div>

DTables.jl would also work on a single machine and is designed for these cases.

> Your suggestion requires 637GB for processing a single file. I don’t think processing 12 files will fit within 1TB.

But this means that you cannot count to fit all 12 files into a single data frame. You need an off-core solution then (like DTables.jl). Also the question is that maybe it is enough for you to process the data table by table and then merge the results?

---

<div class="post-metadata">

**Author:** ![ujjwal\_singh](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/ujjwal_singh/32/53073_2.png) [@ujjwal\_singh](https://discourse.julialang.org/u/ujjwal_singh)\
**Post date:** [September 22, 2023, 7:18am UTC](https://discourse.julialang.org/t/efficiently-merge-multiple-arrow-files-12-files-totaling-2-5-tb-into-a-single-dataframe-for-effective-data-analysis-while-minimizing-memory-overhead/104110/13 "2023-09-22T07:18:21Z")

</div>

I want to merge all the dataframes into a single dataframe before processing.

---

<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:** [September 22, 2023, 7:47am UTC](https://discourse.julialang.org/t/efficiently-merge-multiple-arrow-files-12-files-totaling-2-5-tb-into-a-single-dataframe-for-effective-data-analysis-while-minimizing-memory-overhead/104110/14 "2023-09-22T07:47:56Z")

</div>

I understand, but following what you have said, you do not have enough RAM to do so as your 12 tables in total have more than 1TB of data (unless I misunderstood something). That is why I ask if processing them can be done using map-reduce (or similar) pattern.

Note that any off-core solution (like DTables.jl or other technologies, e.g. databases) anyway will have to do some kind of map-reduce (or similar approach) on your data to process it as it does not fit into RAM.

---

<div class="post-metadata">

**Author:** ![tbeason](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/tbeason/32/15898_2.png) [@tbeason](https://discourse.julialang.org/u/tbeason)\
**Post date:** [September 22, 2023, 9:17am UTC](https://discourse.julialang.org/t/efficiently-merge-multiple-arrow-files-12-files-totaling-2-5-tb-into-a-single-dataframe-for-effective-data-analysis-while-minimizing-memory-overhead/104110/15 "2023-09-22T09:17:48Z")

</div>

I think you have a bit more to learn regarding analyzing big data. You are trying to build the entire table and you just can’t do that. What you want to be doing is to query the whole dataset but materialize only the result of the query.
