Apache Arrow: A new open source in-memory columnar data format
blogs.apache.org
blogs.apache.org
https://git-wip-us.apache.org/repos/asf?p=arrow.git;a=blob;f...
(edited post because I fail reading git and didn't notice the java implementation)
This isn't actually true. The java implementation has been complete and used in Apache Drill, a distributed SQL engine, for the past few years. While we anticipate a few small changes to make sure the standard works well across new systems, this is by no means an announcement without tested code.
https://git-wip-us.apache.org/repos/asf?p=arrow.git;a=commit...
"Analytical queries" are faster in a column-based format, because you're doing things like computing the correlation between 2 variables. Instead of looking at many columns in a specific row, you're looking at most of the data in a few columns. So instead of reading a whole row at once, it would be nice to skip the columns you don't care about and just grab big chunks of two or three columns.
Does that make sense?
edit: For a long time I was confused about why HBase was described as a columnar data store when access was still pretty row-based. I think the reason is because you group columns into column families which can be stored and retrieved separately, so you still get some of the benefit of a true column-oriented store.
The term was overloaded by column-family stores, which were often referred to as just 'column stores', probably by people who were not aware of systems like Vertica and MonetDB.
http://dbmsmusings.blogspot.co.uk/2010/03/distinguishing-two...
Since the column values are stored together, more of them can be crammed into a data page. Querying the data of a column can process much more data per page load than row-based, whereas a row-based table's data page has all other columns.
The values of a column can be stored in a sorted order, which is highly compressible, enabling cramming even more data into a data page. You can get much more data with each data page loading. E.g. The column Name is stored as sorted along with the associated row id.
Name, (row id)
--------------
Joan, (101)
Joanna, (307)
John, (15)
John, (32)
John, (6)
Johnson, (31)
Johnson, (44)
The duplicate names can be stored as one: Name, (row id)
--------------
Joan, (101)
Joanna, (307)
John, (15,32,6)
Johnson, (31,44)
Prefix encoding can further compress the sorted data: Name, (row id)
--------------
Joan, (101)
4+na, (307)
2+hn, (15,32,6)
4+son, (31,44)
Now imagine the query: select count(Name) from Table. Scanning the Name column values only touches 4 records right next to each other in a data page and adds up each record's reference row ids.Select Name from Table where Name like "John%" would do a binary search down the sorted names, load two records 2+hn and 4+son, and expand them into 5 names.
The Github mirror is now live, you can see it here: https://github.com/apache/arrow
Here is a page with the clone URLs: https://git.apache.org/
Have you folks considered Supersonic engine from Google, which was designed with similar (but not as extensive as Arrow) goals in mind?
"Arrow's cross platform and cross system strengths will enable Python and R to become first-class languages across the entire Big Data stack," said Wes McKinney, creator of Pandas.
Code committers to Apache Arrow include developers from Apache Big Data projects Calcite, Cassandra, Drill, Hadoop, HBase, Impala, Kudu (incubating), Parquet, Phoenix, Spark, and Storm as well as established and emerging Open Source projects such as Pandas and Ibis.
this is pretty much nuke-from-orbit.
That analogy might imply overkill, thus highlighting the tactical advantages of the SFrame approach in processing a month's worth of 1-10GB daily-generated SQLite files, for instance.
http://dask.pydata.org/en/latest/dataframe.html
"Dask dataframes look and feel like pandas dataframes, but operate on datasets larger than memory using multiple threads."
Is this an antipattern for Pandas?
If Pandas and other high profile products are endorsing it (and may adopt it), it's going to be very hard for 99% of people to choose something else.
I'm somewhat ignorant as I don't run any Apache projects but I'm curious as to why people choose Apache to back their project these days. I guess why choose a committee instead of just leaving it on Github. I suppose its the whole voting and board stuff.
As for why projects choose Apache over Github, there are lots of reasons. Some of them apply to choosing any foundation (Apache, Eclipse, etc.) over Github: legal rigor, known quantity for enterprisey consumers, and so on. Probably the biggest reason many company sponsored projects end up at the ASF is because the ASF has a reputation as a good place for competing companies to collaborate on a common code base.
Personally, I have chosen to donate a great deal of my own time and energy to the ASF because I greatly treasure its emphasis on governance by individual contributors rather than corporations. (That the ASF is a 501(c)(3) non-profit rather than a 501(c)(6) like some of the more slick, consortium-like foundations is related.)
I don't know as much about the internals of Avro, but I know it is a bit different from Parquet, in that it can be used to serialize and deserialize smaller amounts of data. It is used to store large datasets in files, although it will in most cases be less space efficient than Parquet. It has also been used as a way of embedding complex structures into other systems (similarly to how JSON can be embedded in a database), or for serializing individual structures between systems. The binary representation of Avro needs to be read into a system-specific format like a C/C++ struct/object, Java object, etc. for consumption.
In contrast, Arrow is designed to represent a list of objects/records efficiently. It is designed to allow for a chunk of memory to handed to a lightweight language-specific container that can immediately reference into the memory to grab a specific value, without reading each of the records into it's own individual object or structure.
The existing columnar data formats such as Parquet and ORC aim to be space-efficient since the data is stored in disk and IO operations are usually the bottleneck. The columnar data formats shine in big-data area so the amount of data will be huge. Given that columnar data formats can be compressed efficiently and that's of the main points of columnar data formats such as Parquet and ORC, I'm not sure that I understand the main point of in-memory columnar data formats.
Once the data is in-memory and we can access any column of a row in constant-time what's the difference between a row-oriented data format and columnar data format?
My understanding is that insert will suffer badly for things like compressed columns. If my relational engine is used alike-Datasets where is common to iterate by rows AND do joins, unions, projections and filters, so how balance the thing?
ORC has had its own in-memory columnar data-format for the last couple of years - VectorizedRowBatch.
This doesn't use any Unsafe access, but uses pure JVM arrays for layout since it allows the Java JIT to unroll a lot of loops internally.
Here's an example of a patch from Intel for unrolling the 64 bit operations into 256 bit operations and allowing the Java code to be auto-vectorized to use ymm<n> registers
https://issues.apache.org/jira/browse/HIVE-10180
Secondly, we maintain isRepeating=true independently for each column, which means that operations for 1024 rows can sometimes be as fast as operations for 1 row.
While in a row-based model, it would have no way of maintaining duplicate information on a column directly.
I have looked at the ValueVectors impl in Drill and Tungsten in Spark - which are more cache efficient than Hive's internal columnar structure. However erasing type information into a byte[] structure prevents the JIT from auto-vectorizing inner loops.
My understanding is that serialization became a thing because in-memory representations tend to use pointers to shared data structures that may thus be referenced multiple times while being stored only once. This would not translate 1:1 to serialized representations (where memory offsets would no longer hold meaning) -- much less in any language-agnostic way.
So I have this suspicion that Apache Arrow would not support reusing duplicate data while storing it only once. Would anyone mind clarifying on this point?
How will Arrow use vectorized instructions on the JVM? That seems to be only available to the JIT and JNI, which is a frustrating limitation.
Will Cassandra Java drivers support this ?
http://blog.cloudera.com/blog/2016/02/introducing-apache-arr...