# Repartitioning 2TB of csv into parquets

**URL:** <https://discourse.julialang.org/t/repartitioning-2tb-of-csv-into-parquets/27716>\
**Category:** Data\
**Tags:** big-data\
**Created:** [August 19, 2019, 2:29pm UTC](https://discourse.julialang.org/t/repartitioning-2tb-of-csv-into-parquets/27716 "2019-08-19T14:29:36Z")\
**Posts on this page:** 20\
**Page:** 1

<div class="post-metadata">

**Author:** ![gabomgp](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/gabomgp/32/1918_2.png) [@gabomgp](https://discourse.julialang.org/u/gabomgp)\
**Post date:** [August 19, 2019, 2:29pm UTC](https://discourse.julialang.org/t/repartitioning-2tb-of-csv-into-parquets/27716/1 "2019-08-19T14:29:37Z")

</div>

Good morning from Colombia.

I’m beginner in big data/data science, and i’m trying to do the next task:

```
We have 2 TB of CSV from one table. We want to try to use a SQL Layer to query that data. 
Currently, the data is stored in Timescale, but Timescale doesn't compress data and the used SSD 
space is growing in a fast pace. So, I and a partner are trying to use Azure Data Lake, with a SQL 
Layer over that, to test if: the performance of queryng is acceptable? the price is better or worse?. 
The two SQL layer that we want to try are Dremio and Azure Data Flow Analytics. But the problem is 
csv are sometime very very very large (100 GB) and sometimes very very tiny (10KB). We want to 
repartition the data first and to write the data to parquet second.

```

To do the task, we tryied:

```
1) To use Pandas, in a very large machine (64 cores, more than 450 GB of RAM). The problem was 
   that Pandas doesn't scale to large machines.
2) To use Azure Data Factory Data Flow, the data flow cluster (Spark cluster really) crash with a 
   System Error (?). So, aborted.

```

Finally, the csv are internally sorted by the key we want to use as partition key. So, maybe, we can use Julia to read the csv in streaming, and to write the partitioned data to parquets. Is that a good idea? you can see problems in that aproximation?

P.D.: Sorry my English. I hope you can understand the problem.

---

<div class="post-metadata">

**Author:** ![Zach\_Christensen](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/zach_christensen/32/7220_2.png) [@Zach\_Christensen](https://discourse.julialang.org/u/Zach_Christensen)\
**Post date:** [August 19, 2019, 9:41pm UTC](https://discourse.julialang.org/t/repartitioning-2tb-of-csv-into-parquets/27716/2 "2019-08-19T21:41:15Z")

</div>

Depends on how you want to partition it. Into multiple files by columns or rows. Take a looked at CSV.jl. There are ways of reading parts of a file in at a time

---

<div class="post-metadata">

**Author:** ![gabomgp](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/gabomgp/32/1918_2.png) [@gabomgp](https://discourse.julialang.org/u/gabomgp)\
**Post date:** [August 20, 2019, 4:14pm UTC](https://discourse.julialang.org/t/repartitioning-2tb-of-csv-into-parquets/27716/3 "2019-08-20T16:14:19Z")

</div>

I want to partition the csv by rows. So, is possible to write parquets in Julia?

---

<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:** [August 21, 2019, 3:58am UTC](https://discourse.julialang.org/t/repartitioning-2tb-of-csv-into-parquets/27716/4 "2019-08-21T03:58:56Z")

</div>

Using the [CSV.jl](https://github.com/JuliaData/CSV.jl) package, you can use the `CSV.Rows(file)` structure to efficiently iterate rows; it shouldn’t matter how big the file is, `CSV.Rows` can efficiently handle iterating rows.

Unfortunately, the [Parquet.jl](https://github.com/JuliaIO/Parquet.jl) package doesn’t support writing parquet files yet. The [Feather.jl](https://github.com/JuliaData/Feather.jl) package, however, is able to write parquet-like binary files that can be very efficient to re-read, so you might check that out if it might work for you.

Hopefully that helps a little.

---

<div class="post-metadata">

**Author:** ![tkluck](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/tkluck/32/15769_2.png) [@tkluck](https://discourse.julialang.org/u/tkluck)\
**Post date:** [August 21, 2019, 7:04am UTC](https://discourse.julialang.org/t/repartitioning-2tb-of-csv-into-parquets/27716/5 "2019-08-21T07:04:33Z")

</div>

> [@gabomgp](#):
>
> To use Pandas, in a very large machine (64 cores, more than 450 GB of RAM). The problem was that Pandas doesn’t scale to large machines.

This seems like a reasonable[1] approach, though. Can you elaborate in which way Pandas “doesn’t scale”?

[1] Yes, there’s probably a smarter/cheaper way, but if this is a one-off, your time is probably more expensive than the hardware.

---

<div class="post-metadata">

**Author:** ![gabomgp](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/gabomgp/32/1918_2.png) [@gabomgp](https://discourse.julialang.org/u/gabomgp)\
**Post date:** [August 21, 2019, 3:15pm UTC](https://discourse.julialang.org/t/repartitioning-2tb-of-csv-into-parquets/27716/6 "2019-08-21T15:15:34Z")

</div>

In my use case:

1. When i’m reading CSV, the reading is not using multiple cores; same when i’m grouping (for partitioning), or sorting (i’m sorting for better compression).
2. I’m parallelize using a Pool, but:  
2.1) The CSV are very large sometimes, and reading many CSV in the same time, consume too much memory.  
2.2) I can read the CSV’s in chunks, but as the process is very slow because we are merging the final partitions with existant parquets. I’m not expert, maybe i’m doing something wrong here.

---

<div class="post-metadata">

**Author:** ![gabomgp](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/gabomgp/32/1918_2.png) [@gabomgp](https://discourse.julialang.org/u/gabomgp)\
**Post date:** [August 21, 2019, 3:22pm UTC](https://discourse.julialang.org/t/repartitioning-2tb-of-csv-into-parquets/27716/7 "2019-08-21T15:22:09Z")

</div>

Is possible to write in streamming the CSV, but without a iterator?

I’m thinking in something as:

1. Read in streamming the CSV.
2. With each row:  
2.1) Check if the partition key has a csv file created.  
2.2) If the csv file is created: append the row to the file.  
2.3) If not created: create the file, in a dictionary save the file as the assigned file of the partition key, and append the row.

Reading and Writting in streamming and without grouping the rows, only using a dictionary of partition keys =\> csv files. Is that possible?

---

<div class="post-metadata">

**Author:** ![xiaodai](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/xiaodai/32/15937_2.png) [@xiaodai](https://discourse.julialang.org/u/xiaodai)\
**Post date:** [April 17, 2020, 7:33am UTC](https://discourse.julialang.org/t/repartitioning-2tb-of-csv-into-parquets/27716/8 "2020-04-17T07:33:02Z")

</div>

Once CSV.jl updates so that `CSV.Rows` accepts `types` then I will register this branch  
[https://github.com/xiaodaigh/DataConvenience.jl/tree/CSV-chunk-reader#csv-chunk-reader](https://github.com/xiaodaigh/DataConvenience.jl/tree/CSV-chunk-reader#csv-chunk-reader)

which has a CSV chunk reader by wrapping some of CSV.jl’s functionalities (so thanks to CSV.jl).

---

<div class="post-metadata">

**Author:** ![johnh](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/johnh/32/3615_2.png) [@johnh](https://discourse.julialang.org/u/johnh)\
**Post date:** [April 17, 2020, 8:56am UTC](https://discourse.julialang.org/t/repartitioning-2tb-of-csv-into-parquets/27716/9 "2020-04-17T08:56:14Z")

</div>

Would DatFramesDBs be of any help?

> [@\[ANN\] DataFrameDBs.jl](https://discourse.julialang.org/t/ann-dataframedbs-jl/35718):
>
> Hi all! Julia is my hobby and I had some free time for the last 4 weeks, so here is the first results of my experiments [GitHub - waralex/DataFrameDBs.jl: The DateFrameDBs is the prototype of persistent, space efficient columnar database on pure Julia](https://github.com/waralex/DataFrameDBs.jl) It is the prototype of columnar, persistent, type stable and space efficient database in pure Julia. Some examples on [this](https://www.kaggle.com/mkechinov/ecommerce-behavior-data-from-multi-category-store) dataset imported to DataFrameDBs julia\> using DataFrameDBs julia\> t = open\_table("ecommerce") DFTable path: ecommerce 10×6…

> **[GitHub - waralex/DataFrameDBs.jl: The DateFrameDBs is the prototype of...](https://github.com/waralex/DataFrameDBs.jl)**
>
> The DateFrameDBs is the prototype of persistent, space efficient columnar database on pure Julia - GitHub - waralex/DataFrameDBs.jl: The DateFrameDBs is the prototype of persistent, space efficient...

---

<div class="post-metadata">

**Author:** ![waralex](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/waralex/32/13415_2.png) [@waralex](https://discourse.julialang.org/u/waralex)\
**Post date:** [April 17, 2020, 9:21am UTC](https://discourse.julialang.org/t/repartitioning-2tb-of-csv-into-parquets/27716/10 "2020-04-17T09:21:29Z")

</div>

If you are considering using an sql layer, I would recommend ClickHouse ([https://clickhouse.tech/](https://clickhouse.tech/)). DataFrameDBs is still under development and I would not recommend it to a person who needs to analyze data right now. But yes, it was for such cases that I started to develop it

---

<div class="post-metadata">

**Author:** ![ValdarT](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/valdart/32/24146_2.png) [@ValdarT](https://discourse.julialang.org/u/ValdarT)\
**Post date:** [April 17, 2020, 11:11am UTC](https://discourse.julialang.org/t/repartitioning-2tb-of-csv-into-parquets/27716/11 "2020-04-17T11:11:45Z")

</div>

In case it might help someone: this would also be super simple to achieve with [Apache Spark](https://spark.apache.org/).

---

<div class="post-metadata">

**Author:** ![gabomgp4](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/gabomgp4/32/13366_2.png) [@gabomgp4](https://discourse.julialang.org/u/gabomgp4)\
**Post date:** [April 28, 2020, 4:22pm UTC](https://discourse.julialang.org/t/repartitioning-2tb-of-csv-into-parquets/27716/12 "2020-04-28T16:22:11Z")

</div>

Thanks you, I’m the original poster of the question… (my original account is lost). I tried with clickhouse months ago and i can confirm that is an excelent tool to play with medium data (some terabytes of info) in a easy way.

Really, pandas or vaex are not competition to clickhouse, agregating and playing with the data at this scale. In a single machine i can to aggregate, filter, etc at very fast performance. I’m tried with spark too, but spark is not a simple solution (from the point of view of deployment and operations), and Azure HDInsight is expensive (compared to a single machine with Clickhouse, that can be installed in the developer machine to develop the solution in a very simple way).

It’s sad that clickhouse is [not working in Google Colab](https://github.com/ClickHouse/ClickHouse/issues/10505)

---

<div class="post-metadata">

**Author:** ![gabomgp4](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/gabomgp4/32/13366_2.png) [@gabomgp4](https://discourse.julialang.org/u/gabomgp4)\
**Post date:** [April 28, 2020, 4:47pm UTC](https://discourse.julialang.org/t/repartitioning-2tb-of-csv-into-parquets/27716/13 "2020-04-28T16:47:32Z")

</div>

Thanks for the notification. I’m hoping Julia could offer a solution to medium data problems in some time. I think Clickhouse use a very good architectural approach in that sense. Maybe you are interested in this info:

> **[Overview of ClickHouse Architecture | ClickHouse Docs](https://clickhouse.com/docs/en/development/architecture)**
>
> ClickHouse is a true column-oriented DBMS. Data is stored by columns, and during the execution of arrays (vectors or chunks of columns).

and

> **[Overview of ClickHouse Architecture | ClickHouse Docs](https://clickhouse.com/docs/en/development/architecture)**
>
> ClickHouse is a true column-oriented DBMS. Data is stored by columns, and during the execution of arrays (vectors or chunks of columns).

Note that in clickhouse the size of the block (the chunk) is independent from the data source, so if i have 1 tb csv file, parquet o whatever, the chunks are of the same size always, and that size is selected to be very efficient with the vectorized operations that the rest of the execution pipeline is using, and to consume little memory.

---

<div class="post-metadata">

**Author:** ![xiaodai](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/xiaodai/32/15937_2.png) [@xiaodai](https://discourse.julialang.org/u/xiaodai)\
**Post date:** [April 28, 2020, 8:05pm UTC](https://discourse.julialang.org/t/repartitioning-2tb-of-csv-into-parquets/27716/14 "2020-04-28T20:05:32Z")

</div>

Have u tried [diskframe.com](http://diskframe.com)?

---

<div class="post-metadata">

**Author:** ![gabomgp4](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/gabomgp4/32/13366_2.png) [@gabomgp4](https://discourse.julialang.org/u/gabomgp4)\
**Post date:** [April 28, 2020, 11:24pm UTC](https://discourse.julialang.org/t/repartitioning-2tb-of-csv-into-parquets/27716/15 "2020-04-28T23:24:44Z")

</div>

Nop, i haven’t tried diskframe. diskframe has comparable performance and workflow to dask?

---

<div class="post-metadata">

**Author:** ![xiaodai](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/xiaodai/32/15937_2.png) [@xiaodai](https://discourse.julialang.org/u/xiaodai)\
**Post date:** [April 28, 2020, 11:32pm UTC](https://discourse.julialang.org/t/repartitioning-2tb-of-csv-into-parquets/27716/16 "2020-04-28T23:32:35Z")

</div>

I would say in samw ball park. But i am often surprised at how unoptimised dask is ag certain tasks.

---

<div class="post-metadata">

**Author:** ![anon92994695](https://avatars.discourse-cdn.com/v4/letter/a/ce7236/32.png) [@anon92994695](https://discourse.julialang.org/u/anon92994695)\
**Post date:** [April 29, 2020, 12:03am UTC](https://discourse.julialang.org/t/repartitioning-2tb-of-csv-into-parquets/27716/17 "2020-04-29T00:03:28Z")

</div>

A cheap way of doing this would be to chunk the data via linux’s `split` command such that each chunk fits into memory. You can write your own split, but the linux cli tools tend to be quite good and fast for these tasks!

For each chunk, read with csv, and break up the data into memory mapped column stores, Feather or JuliaDB file stores.

Pick a consistent naming convention for each file ie: “chunk\_00001”. Then voila. Now you can do 1 to 1 transforms/joins across cores on chunks and be 4-16x better then pandas :P.

But yea with Apache Spark you should be able to do something very similar to this with parquets. My advice is run like hell away from PySpark, and use real Spark. But Spark has it’s own issues…

This all of course really depends on your chunking strategy… Doing it simply by rows is not very handy for most cases. but! One can write very simple emitter functions to break up data by some other criteria in base Julia really easily…

I’ve been working on a simple package for something not directly related but maybe helpful if you are doing lookups: [https://github.com/caseykneale/LockandKeyLookups.jl](https://github.com/caseykneale/LockandKeyLookups.jl)

---

<div class="post-metadata">

**Author:** ![xiaodai](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/xiaodai/32/15937_2.png) [@xiaodai](https://discourse.julialang.org/u/xiaodai)\
**Post date:** [April 29, 2020, 12:21am UTC](https://discourse.julialang.org/t/repartitioning-2tb-of-csv-into-parquets/27716/18 "2020-04-29T00:21:36Z")

</div>

> [@anon92994695](#):
>
> A cheap way of doing this would be to chunk the data via linux’s `split` command

Yeah. I do that, the issue is windows users can’t use it.

Try [GitHub - xiaodaigh/DataConvenience.jl: Convenience functions missing in Julia](https://github.com/xiaodaigh/DataConvenience.jl)

CSV Chunk Reader  
You can read a CSV in chunks and apply logic to each chunk. The types of each column is inferred by CSV.read.

```julia
using DataConvenience: CsvChunkIterator
for chunk in CsvChunkIterator(filepath)
  # chunk is a DataFrame
  # do something to df
end

```

The chunk iterator uses CSV.read parameters. The user can pass in type and types to dictate the types of each column e.g.

```julia
# read all column as String
for chunk in CsvChunkIterator(filepath, type=String)
  # df is a DataFrame where each column is String
  # do something to df
end

```

---

<div class="post-metadata">

**Author:** ![anon92994695](https://avatars.discourse-cdn.com/v4/letter/a/ce7236/32.png) [@anon92994695](https://discourse.julialang.org/u/anon92994695)\
**Post date:** [April 29, 2020, 12:23am UTC](https://discourse.julialang.org/t/repartitioning-2tb-of-csv-into-parquets/27716/19 "2020-04-29T00:23:34Z")

</div>

So crazy I had to write one of these myself ~2 yrs ago. You’re all taking me back down julia memory lane and why I left python in the first place…

---

<div class="post-metadata">

**Author:** ![phelipe](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/phelipe/32/14643_2.png) [@phelipe](https://discourse.julialang.org/u/phelipe)\
**Post date:** [April 29, 2020, 4:20pm UTC](https://discourse.julialang.org/t/repartitioning-2tb-of-csv-into-parquets/27716/20 "2020-04-29T16:20:29Z")

</div>

There is a little package for [ClickHouse](https://github.com/maxmouchet/ClickHouse.jl) in julia.

[Next page](https://discourse.julialang.org/t/repartitioning-2tb-of-csv-into-parquets/27716.md?page=2)
