3 ms·
I've been centering Arrow in my data workflows, which has led to a system that uses Ray for distributed computation, Polars for single node data frame computati
by liminal 4y ago
I've been centering Arrow in my data workflows, which has led to a system that uses Ray for distributed computation, Polars for single node data frame computation on Arrow data sets, and Ray Data for distributed data frames using Arrow. I love the speed and Polars' expressive API, and I love that I can use any other Arrow tool on the data without any serialization overhead.
That said, using Ray places limits on that, since it can distribute data across machines, breaking the Arrow model of in-memory data sharing between tools. I know that Ray Data is starting to use Polars internally, but I wish I could use the Polars API over distributed Ray Data sets so I could write code once and not worry about whether Ray has split it across machines.