Apache Arrow 1.0
arrow.apache.org
arrow.apache.org
Even after issues with its inability to map datetime64 values properly, I was reasonably happy with my design choice.
I became less happy on discovering that it's very weak as an interchange format for cross-language work.
In my case I wanted to use some existing JVM-based tooling. This caused huge pain. The Jvm/Java library/API is a complete mess, sorry, and if people are complaining about the Python/C++ documentation, there's basically nothing for the Java library. It's barely useable and the dependencies are horrific - the whole thing is mingled with hadoop dependencies - even the API itself.
And the API is barely above exposing the file-format. Nothing like "load this parquet file" into some object which you can then query for it's contents - you're dealing with blocks and sections and other file-format level entities.
The other issue is caused its flexibility - for example Panda's dataframes are written with what's effectively a bunch of "extension metadata" which means it works great for reading and writing pandas from Python but don't expect anything to be able to work with the files out-of-the-box in other languages.
In the end, the only way I could get reliable reading and writing from the JVM was to only store numeric and string data from the Python side. Even then it feels flakey - with a bunch of hadoop warnings and deprecation warnings. I know the JVM has little appreciation in the data science world which is maybe a reason for the sorry state of the Java library.
Edit: to be specific, I am talking about my experiences with Arrow/Parquet.
Andy Grove is building a "distributed compute platform implemented in Rust, using the Apache Arrow memory model": https://github.com/ballista-compute/ballista. Seems possible.
Build a spark killer in Julia. Everyone can read the code.
Rust, to me, seems like a natural enough choice. It is easy to mate it to other languages, including all the major data science ones, so it would theoretically work well as the basis for a distributed compute engine that has good support for all of them as client languages. Would the same work for Julia? IIRC, it's a bytecode compiled language, which I imagine would make it difficult to link Julia libraries from other technology stacks.
Databricks rewrote the Spark engine in C++ (called the Delta Engine, see here: https://databricks.com/blog/2020/06/24/introducing-delta-eng...).
Hopefully there will be some good alternatives that don't consume so much memory soon ;)
The alternative to Spark will be Spark, as JVM slowly improves their memory APIs from release to release.
People tend to hand wave how much it costs to reboot an ecosystem.
> It's barely useable and the dependencies are horrific - the whole thing is mingled with hadoop dependencies - even the API itself.
These are comments about http://github.com/apache/parquet-mr which is a different open source project.
For C++ / Python / R many of the developers for both Apache Arrow and Apache Parquet are the same and we currently develop the Parquet codebase out of the Arrow source tree.
So, I'm not sure what to tell you, we Arrow developers cannot take it upon ourselves to fix up the whole JVM data ecosystem.
That could be improved without fixing the whole JVM data ecosystem, but that's mostly up to JVM developers. It's unfortunate if the Spark developers using Arrow aren't contributing in this area (especially since many of them are being paid), but it's all open source and undoubtedly pull requests are welcome.
Congratulations on the 1.0 release, it's only going to keep getting better! Really exciting to be able to share data in-memory across languages.
But at your landing page, it's claimed "Apache Arrow defines a language-independent columnar memory format for flat and hierarchical data, organized for efficient analytic operations on modern hardware like CPUs and GPUs. " and that "Libraries are available for C, C++, C#, Go, Java, JavaScript, MATLAB, Python, R, Ruby, and Rust.". This certainly gave me the impression that more than just Python, C++ and R would be well supported.
The JVM isn't complete irrelevant in data-science given the position of Spark/Scala. This also raised my expectations of arrow/parquet because it seems to be the de-facto standard for table storage for this JVM platform. And I experienced no issues on that platform.
To be clear, I'm not blaming you for my design decision (I'm a software engineer not a data-scientist btw), and I still think parquet/arrow rocks for Python but in my experience it doesn't really deliver a useable "cross-language" file format at the moment.
Similarly, the reference implementation of Parquet may be in Java, but consuming it from a Java language, outside of a Spark cluster, is still a royal pain. Whereas doing it from Python isn't too bad.
Long story short, I think that expecting a project that's just trying to implement a columnar memory format to also muck out the world's filthiest elephant pen is perhaps asking too much. Though perhaps a project like Arrow could serve as the cornerstone of an effort to douse it all with kerosene and make a fresh start.
Stuff like Arrow doesn't come even into the radar of IT.
Yes, it's just a building block. The easiest way to use Parquet on Java is to use Spark's integration with it, because it provides the query engine for you. But, if I'm not mistaken, it's much bigger than Pandas.
I also strongly dislike the Hadoop coupling (hey, it's called parquet-mr for a reason...), but it's more or less an invisible annoyance if you use Spark in an environment like Amazon EMR.
I think you'll get a lot of confused responses from people if you want to try to build an application that reads Parquet directly. It's the disk format for a bunch of distributed database engines. They'll wonder how you plan on querying it.
* Easy type inference, even for nested maps and structs
* lz4 compression support
* SQL and directory partitioning out of the box
If you use pyarrow, you get: * Write support for nested types, but read support is broken / incomplete (it throws a TODO error)
* no lz4 support and a load of Jira politics blocking it
* You can query using Pandas (not SQL), but that querying can be much slower versus Spark (Spark is naturally parallelized)
pyarrow might need to hit 2.0.0 to be really viable. It’s definitely easier to use than parquet-mr though.
I've been meaning to take a closer look at ORC. As a Spark user, I just sort of defaulted into Parquet. ORC is very similar, though, and seemingly gives every indication of being the more mature product.
Also is Parquet and Arrow the same ? df.to_parquet('df.parquet.gzip', compression='gzip') will not use arrow i presume ? i have to use a separate library to save to parquet using arrow. a bit confused.
RE:Feather -- Arrow itself isn't necessarily a full file format -- you can imagine memory buffers being all over the heap with giant gaps inbetween b/c diff cols generated at diff times -- but in practice folks will indeed serialize to disk consolidated buffers (pa.Table -> write stream -> file). If we couldn't do that, RPC wouldn't work ;-) My understanding of Feather is it standardizes ideas around this consolidation, but we are able save to disk (within versions) without it. We found it more predictable to stick to ~Parquet for storage and Arrow buffer passing for streaming, but now that Feather networking APIs for accelerated bulk transfers may be stablizing, there may be speed advantages to using it over manual buffer streaming (and still stick w/ Parquet for persistent files).
Arrow<>Parquet conversion is super fast b/c of the co-design around similar concepts: both using record batches of dense binary column buffers means implementations can pointer-copy, memory map, use bulk copy primitives, etc. for zero-copy or at least highly accelerated interop. Python RAPIDS GPU kernels can therefore selectively stream in a few parquet columns across many parquet files through a single 900GB/s GPU, compute over them, and write back out to arrow or a new parquet.
http://arrow.apache.org/faq/index.html#what-about-arrow-file...
You can store them long-term if you want (and you'll still be able to read them 5 years from now) but we aren't optimizing the Arrow IPC format for the _needs_ of long-term storage.
I'm a big fan of Pandas, and didn't know about Arrow. I've been considering do a talk advocating for a consistent data-frame api across languages since IMHO, it's the next fundamental data structure that should have baked in support everywhere. So it appears you've at least somewhat beaten me to the punch.
Since Arrow is more than an API to tabular data structures, what would you think about a Promises/A+-like specification for dataframes?
How much of the Arrow API do you think end users will wind up using, as opposed to being a lower-level framework that projects like pandas and dplyr wind up using behind the scenes?
Finally, do you think that Arrow has the potential to be the logical successor to pandas? If not, what is your long term strategy to address the shortcomings that you see in pandas?
I use pandas every day, thank you for that. Just watched this keynote and I really like the vision; I currently work with a bunch of guys who prefer DPLYR.
Is Arrow just for in-memory analytics, or are there plans to support in-database analytics too?
We've been on quite the journey here: https://www.graphistry.com/blog/graphistry-2-29-5-upload-100... . Think json -> protobuf -> arrow, and paralleled in our work on parallel js -> opencl+js -> python rapids/cuda for accelerated native compute over it. The blogpost demos the ~100X bigger & faster dataset result of supporting Arrow for our uploader & RAPIDS for our parser when you can't and want us to convert for you.
Something folks miss with Arrow, IMO, is it is like google protobufs for everyone else. Arrow is not just having a nice binary format, but is also ready out-of-the-box for streaming, rich datatypes, and for larger/longer-term projects, standardized & auto self-describing schemas. If you've had to manually decipher, maintain, and update generated protobuf schemas (because you aren't google with all the internal ~protobuf tooling / integrations / infra / etc.), that should sound pretty good ;-) A lot more to do, but already way ahead of most other things, esp. in aggregate.
Having a single high-performance in-memory format means different programs can read/write from the same source without serializing/copying/deserializing. For instance, if you wanted to pass a huge table of data from R to Java to Python (because your tools span different languages), normally you'd have to copy and serde (Protobuf? JSON?) to pass the data, which equals huge overheads. With Arrow, each of those languages can directly interact with the same copy of in-memory data, in-process -- with the highest possible performance.
You also get the performance of columnar databases without implementing your own columnar data structure.
But of course, no harm adding a short description to the title to broaden its audience. Arrow is truly something amazing and the more people know about it the better.[1] Folks who program against traditional databases might not know about it, and I think they should, especially if they need to generate analytics (i.e. fast filtering/aggregation for dashboarding or for data pipeline tasks).
[1] Overview: https://arrow.apache.org/overview/
"Libraries are available for C, C++, C#, Go, Java, JavaScript, MATLAB, Python, R, Ruby, and Rust."
In fact, if my coworkers (educators) are any indication, few people do.
Even 2 sentences on this page would've helped a lot.
A lot of people just stay in their lane. Having a solid description, like the one you provided, would be super useful. Perhaps a tool-tip feature of HN, even.
In general, where's the best place to learn more about Arrow? I've approached it several times, and can find a lot about how to integrate it into other products, but none of the tools like the query engines that I would find very useful.
Are you talking about there being support for multiple language libraries like PyArrow or about there being multiple Apache projects that utilize Arrow like Parquet and Spark?
If not, I'm not following what sub-projects you are speaking about. As far as I know, Arrow is principally the Arrow Columnar Format and Arrow Flight with some other potentially interesting interfaces for compute kernels and CUDA devices.
Am I missing something?
There are some Jira issues on this, but there doesn't seem to be a consensus on the way forward. Does someone have more information, is the general idea to wait for specialization to stabilise, or is there a plan, or even an intention, to stop relying on it?
As for stablizing packed_simd, It's completely unclear to me when that will land in stable rust. I recently had a project where I just ended up calling out to C code to handle vectorization.
EDIT: According to, https://github.com/rust-lang/lang-team/issues/29, the effort looks abandonded/deprioritized, so it may be a long time before it sees stable rust
packed_simd provides a convenient platform independent API to some subset of common SIMD operations. Rust's standard library does have pretty much everything up through AVX2 on x86 stabilized though: https://doc.rust-lang.org/core/arch/index.html --- So if you need vectorization on x86, Rust should hopefully have you covered.
If you need other platforms or AVX-512 though, then yeah, using either unstable Rust, C or Assembly is required.
We've been using Arrow for a few years in our startup (link in profile) - it's been great as a common format for passing tabular data between Python and Javascript.
Haven't used the in-memory/zero-copy features as much, but as a binary, high-performance, typed format for the data analytics world it can't be beat. And now that the Feather format is basically the Arrow format, I expect to see it really take off as a common interchange and even storage format for medium- to even long-term projects.
(Also nice to see the reduced Python wheel size - a little bonus for the 1.0 release :) )
This can particularly be useful for low memory systems like ARM SBC when conducting long duration research. If you want to build Apache Arrow from source for ARM, I've written a How-To here[1].
[1]https://gist.github.com/heavyinfo/04e1326bb9bed9cecb19c2d603...
As I try to move to Arrow/Flight over using JSON or binary formats like Protobuf, one thing I see missing is tooling or schema -> code generation. It would be nice to see some progress or roadmap in that direction as well.
I like to imagine a future where data is freed from the languages/tools used to operate on them. In-memory objects would then be views into data stored in shared memory (or disk), with the data easily read/manipulated from multiple languages.
Oh look, yet another Apache real time/batch/big data/stream processing/ingestion/workflow/whatever product.
Apache Druid
Apache Spark
Apache Storm
Apache Flink
Apache Beam
Apache Apex
Apache Airavata
Apache Samza
Apache TEZ
Apache Hama
It's basically a terrible joke at this point. There's no single Apache page helping you to decide which one you want, and they all seem to have such large overlap. Most of them seem to have bad documentation, and give the appearence of not really being maintained.
This puts me off even trying to use them. If there's this much scope creep/NIH/reinventing the wheel happening across the board, I can't imagine how bad each product is individually.Apache Kafka seems to be the only exception...
If I never use multiple languages in the same process can I safely ignore it or are there other benefits?
edit - this looks a good explanation, parquet is more widely used and smaller, but slower. https://stackoverflow.com/questions/48083405/what-are-the-di...
On the nodejs services + frontend JS side, there was no real equiv tool for interop -- traditional soln is slowly round-tripping through SQL/ORM or manually doing protobuf through some sort of pubsub -- so Arrow is part of how we have been taming that mess too. It'll be a longer journey for Arrow or something like it to get adopted in JS land, but I can see folks doing TypeScript and serverless wanting it for a cleaner solution to typed data, faster serialization & streaming, etc. (There is no true JS equiv of pandas nor a typed variant.) We were a bit early here b/c we wanted streaming through WebGL + OpenCL/CUDA, and while we are fans of typed data, found protobuf tooling to be too unintegrated and manual in practice.
As you mention in your second paragraph, arrow is perfect for getting typed query results from a database server directly to a web frontend with minimal overhead. With the Transferable interface, it's even possible to zero-copy transfer the arrow data from the network buffer to a webworker.
In comparison, CSVs are pretty much unstructured text files with added suggestions.
More pipelines should natively support Parquet files, or something like it, and have a thin CSV to Parquet conversion layer on top.
P.S. I'd love to create a Parquet file visualization tool (something like https://github.com/in3rsha/sha256-animation) in the near future.
It should get a lot more press.
That seems pretty straightforward, what part was confusing for you?
The passive aggressive isn't necessary.
So by itself it doesn't necessarily introduce new use cases that weren't there before, except where performance was a bottleneck.
But we're talking like 500x performance increase on workloads vs ODBC, so night and day for a lot of data scientist pipelines.
It boils down to a standard columnar format usable from disk or memory that is useful for data interchange between languages or frameworks, as well as any use case where a columnar format alone has benefits.
I was originally using a simple binary format for fast decoding, and switched to Arrow to be able to select only a few columns at a time. I was impressed by the speed gains, and the size benefits of storing data in column order.
The one thing I wish Arrow had was a way to attach some metadata to its files (if there is, I haven't found it). I originally tried to write a small header before starting to write the Arrow data, but that made it impossible to read back the data: as far as I can tell Arrow looks at the whole file and stores its column definition at the end of the file and computes data offsets based on the start of the file, meaning that there's nowhere left for me to store anything.
It's still a very promising library and I'll definitely be checking out the 1.0 release.
I do that kind of stuff all the time with Go and it's pretty fast with 20~40 million records, averaging 100 KB each. Are those tools oriented to billions, instead of millions? What are the benefits?
Best example is probably pyspark.
Languages that use Arrow can interact with the source data directly in-memory, in-process. No need to move any data around. The fastest serde operation is one that you don't have to do at all.
How does Arrow relate to Flatbuffers?
Flatbuffers is a low-level building block for binary data serialization. It is not adapted to the representation of large, structured, homogenous data, and does not sit at the right abstraction layer for data analysis tasks.
Arrow is a data layer aimed directly at the needs of data analysis, providing a comprehensive collection of data types required to analytics, built-in support for “null” values (representing missing data), and an expanding toolbox of I/O and computing facilities.
The Arrow file format does use Flatbuffers under the hood to serialize schemas and other metadata needed to implement the Arrow binary IPC protocol, but the Arrow data format uses its own representation for optimal access and computation.
* Standardizes binary interop and "serialization" of large structured data, removing all conversions / serialization at ingest and export boundaries. This alone can mean > 2-100x performance improvement in an application that processes a lot of data
* The Arrow in-memory format is an ideal data structure to code analytical algorithms against.
Check out my 18min talk from a few years ago about the vision for the project https://www.youtube.com/watch?v=wdmf1msbtVs
This might answer your question on what is the significance of arrow, given today pandas are kind of basic ingredient in scientific computing, AI and ML.
[1] https://wesmckinney.com/blog/apache-arrow-pandas-internals/
Avro and Parquet can't really do timestamp with timezone, which is a major pain when processing data from geographically distributed locations.
But if Arrow can do it, then maybe there's hope for Parquet as well.
Maybe I'm just getting old, but back when we had to melt our own sand to build a computer, we called these things arrays.