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?
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?
* 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/
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.