Things like eg. protobuffers support hierarchical data which seems like a superset of columns. Is there a benefit to a column based format? Is it an enforced simplification to ensure greater compatibility or is there some other reason?
Things like eg. protobuffers support hierarchical data which seems like a superset of columns. Is there a benefit to a column based format? Is it an enforced simplification to ensure greater compatibility or is there some other reason?
(A int, B int, C int, D int)
And I write: A + B
In a columnar representation, all the As are next to each other, and all the Bs are next to each other, so the process of (A and B in memory) => (A and B in CPU registers) => (addition) => (A + B result back to memory) will be a lot more efficient.In a row-oriented representation like protobuf, all your C and D values are going to get dragged into the CPU registers alongside the A and B values that you actually want.
Column-oriented representation is also more friendly to SIMD CPU instructions. You can still use SIMD with a row-oriented representation, but you have to use gather-scatter operations which makes the whole thing less efficient.
Columnar data is a struct of arrays, with each array representing a column.
1. AAAA BBBB CCCC DDDD
2. ABCD ABCD ABCD ABCD
One major heuristic in how CPUs make your code fast is to assume that if you access some memory, you're probably interested in the memory nearby. So when you access the first "A" bit of memory (common to both sequences above), depending on the memory layout you use, the CPU might also be smart and load the next bits into memory too -- maybe the next "AA", maybe "BC".
Depending on your workload, one or the other of those might be faster. If you're only interested in the first ABCD element because you're doing
SELECT * FROM users WHERE id=$1
then you'll likely want "row-oriented" data -- the #2 scheme above. But if you're interested in all of the A values and none of the values from B/C/D because you're doing SELECT AVG(age) FROM users
then you'll likely want something "column-oriented" -- the #1 scheme above.About the only thing protocol buffers has in common is that it's a standardized binary format. The use case is largely non-overlapping, though. Protobuf is meant for transmitting monolithic datagrams, where the entire thing will be transmitted and then decoded as a monolithic blob. It's also, out of the box, not the best for efficiently transmitting highly repetitive data. Column-oriented formats cut down on some repetition of metadata, and also tend to be more compressible because similar data tends to get clumped together.
Coincidentally, Arrow's format for transmitting data over a network, Arrow Flight, uses protocol buffers as its messaging format. Though the payload is still blocks of column-oriented data, for efficiency.
It also has the added benefit of eliminating serialization and deserialization of data between processes - a Python process can now write to memory which is read by a C++ process that's doing windowed aggregations, which are then written over the network to another Arrow compatible service that just copies the data as-is from the network into local memory and resumes working.
Is that accurate? It still has to deserialize from apache arrow format to whatever the cpu understands.
The advantage you describe is in the operations that can performed against the data. It would be nice to see what this API looks like and how it compares to flatbuffers / pq.
To help me understand this benefit, can you talk through what it's like to add 1 to each record and write it back to disk?
Consider Spark and PySpark. The Python bits of Spark are in a sidecar process to the JVM running Spark. If you ask PySpark to create a DataFrame from Parquet data, it'll instruct the Java process to load the data. Its in-memory form will be Arrow. Now, if you want to manipulate that data in PySpark using Python-only libraries, prior to the adoption of Arrow it used to serialize and deserialize the data between processes on the same host. With Arrow, this process is simplified -- however, I'm not sure if it's simplified by exchanging bytes that don't require serialization/deserialization between the processes or by literally sharing memory between the processes. The docs do mention zero-copied shared memory.
It's truly magical when you scope down a SELECT to the columns you need and see a query go blazing fast. Or maybe I'm easily impressed.
- Improved compression (e.g. a column of timestamps).
- Flexible schemas being easy to manage (e.g. adding more columns, or optional columns).
- Vectorization/SIMD-friendly.