Apache Arrow 4.0
arrow.apache.org
arrow.apache.org
I read the overview, and I'm not sure yet, but is this like an in-memory database that runs inside your process? Like, sqlite without disk persistence, or Erlang ETS, but then columnar?
I can't completely tell from the overview whether it's about the data format or the querying capability. A columnar ETS alternative would be splendid indeed!
Generally people start off by doing these computations in a series of batch jobs (an "ETL pipeline", orchestrated with something like Airflow), to transform data into whatever shape they ultimately want it in; streaming technologies like Spark Streaming and Kafka can help with incrementally adding new rows to your data, rather than recomputing the whole thing every batch-job run.
Whenever you want to involve multiple systems or multiple libraries in your dataframe transformations, there's potentially a lot of computational overhead in serializing the dataframes or just converting them between memory representations. Arrow is a standardized format, spearheaded by the person who wrote Pandas, that attempts to match the in-memory representation, so that whether you're passing the data between libraries in-memory or writing a file for some other system to read, no unnecessary transformations need to happen to work on the data.
> To do this, you generally want your data in column-major format
I'd argue that the basic element of linear algebra is matrix vector multiplication, which I figured was best done row-major. Column major is great in other data use cases, but 'linear-algebra-like, therefore column major' doesn't feel right.
* Dictionary encoding: US,US,US,US,FR -> US:0,FR:1;0,0,0,0,1
* Run-length encoding: 0,0,0,0,1 -> 4x0,1x1
* Delta encoding: 0,1,2,3,4 -> 5x'+1'
* Storing the min and max for a chunk
Basically: exploit the data type to compress it.
Which enables very fast filtering and projections. (And now that the IO bottleneck has been managed you can do your gigantic logistic regression)
But it's also possible to think of it as "Grab one element of the vector, use it to scale the corresponding col of the matrix, and repeat, summing results." Both are efficient means of finding the result, and both have block-level versions that play nicely with the machine cache.
Meanwhile, linear algebra also often involves finding vector norms, and scaling vectors, and so on, and the way we usually set up tables means that the vectors of interest are generally columns of the data tables.
As others have mentioned, for some operations it can also save you from loading whole columns that aren't relevant for your transformation. The compression point in the sibling comment is definitely also relevant, especially for serialization. A whole lot of reasons to use column vectors.
Using "column-major" here might've been terminology abuse; sorry for the confusion.
if you're doing a linear algebra like transformation, you want to do it on all the prices or all the quantities, and a linear algebra library expects a big array of numbers, which is why you have to transform your records into an array of prices and an array of quantities.
"column" here refers to properties of objects, and not rows vs columns with in an array of number
ETL is a design pattern.
Kafka and Cassandra are tools.
Arrow is a data format.
Sorry, will cop to being completely defensive here, but don't call my baby ugly.
>Sorry, will cop to being completely defensive here, but don't call my baby ugly.
Not trying to be antagonistic here for what it's worth - I see your statement as far more guilty of casting your understanding of this technology as a culture jamming point of pride than the GP.
Furthermore the implication that "people who unironically refer to this process as ETL" is the dominant culture (and statements expressing the irony of giving the most common form of computation a specific meaningless acronym are necessarily "culture jamming") is not correct in my experience but YMMV
Dominant culture has nothing to do with it, it's about respecting boundaries. You not knowing something isn't an indictment on the something.
It's a chunked columnar storage format for large data sets. Think of storage layout for a table of data. With arrow, the table is first split row-wize into chunks (sometimes/often just one chunk), then each column of each chunk is stored as an array. The underlying arrays layout supports vectorized operations.
Querying is largely independent of arrow itself which is just the memory format. But the format was designed to support efficient querying. If used as a disk format, for example, you can efficiently load a subset of columns without touching the entire file. And if you're lucky and/or you've chunked the data appropriately you may be able to skip entire chunks.
The format is also language agnostic - as long as the language is python or C++ :) - and allows zero-copy passing of data across the language barrier assuming a shared memory model. This zero-copy feature is important when dealing with large in-memory data-sets.
Unfortunately, until the entire Python data-science ecosystem is re-written from scratch, the application for arrow will largely for library writers and plumbing.
Yes, as a data-scientist, you can easily turn your arrow table into a Pandas Dataframe or a numpy array but you run the risk of expensive copies occurring (actually a copy is guaranteed to happen if the table has more than one chunk) which sort-of defeats the whole purpose. And since to do anything useful with the data you're going to have to perform this conversion - as most of the Python data-science ecosystem is built on numpy and pandas - the format is not particularly interesting to data-science users, I feel.
it their storage formats, I believe, will continue to dominate for columnar storage in the Python world for the foreseeable future.
I use PyCall to use Python from Julia, and PyJulia to do the reverse. PyCall & PyJulia have functionality to easily share arrays. PyCall is pretty seamless, PyJulia is a bit more work but still solid.
To share a DataFrame, I convert it to Arrow, get a bytearray representation of that, use PyCall/PyJulia to access the array from a different process, and reinterpret it as Arrow data within that process.
There is feather for persistence, but you don't need it: just as how you can stream binary arrow buffers to processes, you can write raw arrow to disk. In theory it might give some teams in some setups parallel read/write speedups, but we've been exploring other paths there, e.g., 90+GB/s per node via GDS https://pavilion.io/nvidia . I'm not aware of feather efforts targeting that kind of perf but would be curious!
To utilize w/ spark.. it already does underneath ;-) an increasing flow is something like spark filter -> gpu compute+ai, where the transfer is spark cpu rdd -> arrow (spark-native) -> rapids/tensorflow
Edit: Arrow dev does seem more active than parquet/orc (and a lot of their dev is _by_ arrow devs!), so give it another couple of years, and I can see arrow being stable enough that you can persist data with less fear of having to reprocess older files and having most of the compression features you'd want!
Can you elaborate why Arrow is not a good format for storing to disk? If you’re using it for in-memory querying, why would you not want to also serialize it directly to disk instead of using some intermediary format?
Performance: Arrow does not do significant compression. Feather started adding it, but that adds even more change risk. Parquet/ORC/Arrow are all fairly similar, so until Arrow catches up and stablizes, I'd stick w/ Parquet/ORC. We do GPU stuff, and get in-GPU decompression already, so that's been a win/win.
Think, “CapnProto” or “Protobufs” but for querying data rather than only data transfer.
In theory Arrow will enable users to “trivially” create high performance SQLite type querying in any programming language. Arrow doesn’t help with any other features of SQLite tho.
Maybe one day Arrow will also target on-disc analytics/persistence but for now it delegates to Parquet.
Arrow has compute libs, but we don't use them (yet). More as interop for compute engines that are fast (ex: rapids for gpu) or feature rich (ex: pandas for chunks). Likewise, for I/O, interop for good formats there: on-the-fly arrow -> persistent parquet|orc.
It's enabling us to do cool stuff like fast interop between our cpu<>gpu code, and in the latest initiative, crossing even process & language runtime boundaries w/ zero-copy.
So I will extract it (get it out of those other places), transform it (put it in my format), and load it (store it here).
The term has been around for decades and traditionally is used by database people, but it can apply to any process that does this.
A concrete example might be if you run a business where you sell used books through Amazon and eBay. They each have their own format for info about the status of a product you have listed (whether it has sold, which shipping option the buyer chose, whether payment was received, etc.), but you might want to have that data from both sources in one place so you can see a dashboard or analyze it.
Kafka is a distributed append-only log widely used in data streaming.
Cassandra is a distributed NoSQL datastore.
Not sure why you raised that, though, given they're not mentioned in the linked article.
Parquet is not Arrow, but they work well together, in that one can easily be (de)serialized to the other.
Feather uses the Arrow IPC format internally.
[0] https://github.com/apache/arrow-datafusion
[1] https://github.com/apache/arrow-rs
[2] https://arrow.apache.org/blog/2021/05/04/rust-dev-workflow/
https://uwekorn.com/2021/01/11/apache-arrow-on-the-apple-m1....
Special kudos to the Rust team for Parquet predicates pushdown feature.
I mean it has more complex data structures (FST) than just columnar, but yeah for doc values and such I agree, that is exactly why I'm curious if Arrow is targeting that usecase, and will be competitively "highly optimized".
> and won't see any benefit from changing on disk formats.
I'm less interested in Lucene actually migrating to Arrow (although if it reduces tech-debt they should look into it), I'm most interested in if Arrow will help future Lucene-like libraries get implemented with competitive performance.
Also since Arrow version=X is cross-language compatible, it would be amazing to be able to create "Lucene" indexes (or segments) in Java (perhaps for legacy reasons), then use Rust or Go to query the data.
No persistent storage. Arrow is meant to be used for in-memory queries.
Column-major storage. This enables more efficient data-science-like queries, such as univariate statistics on columns.