Is that accurate? It still has to deserialize from apache arrow format to whatever the cpu understands.
Is that accurate? It still has to deserialize from apache arrow format to whatever the cpu understands.
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.
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?