5 ms·
Pandas on Ray – Early Lessons from Parallelizing Pandas
- kgos 8y agoCode Repo: https://github.com/modin-project/modin https://github.com/modin-project/modin
- smittywerben 8y agoFor those of you confused like me, "Pandas on Ray has moved into the Modin project" http://ray.readthedocs.io/en/latest/pandas_on_ray.html http://ray.readthedocs.io/en/latest/pandas_on_ray.html
- innagadadavida 8y agoDoes anyone here know if Ray is some sort of Yarn competitor? If not what problem space is it in?
- maccam912 8y agoIt feels to me like more of a spark competitor, aiming for the people who want more flexibility than Spark and not needing as much cluster setup (a redis cluster is the hardest part). It might even be able to use Yarn to do some of the resource management.
- mehrdadn 8y agoI don't know what Yarn is, but Ray is a distributed heterogeneous computing framework. It's meant to make it easy to take fairly arbitrary computational programs (with machine learning being an important test/use case) and run/debug them in parallel across lots of machines with high performance in a natural fashion and without drastic changes. [1] The advantage compared to (say) Hadoop is that it allows for heterogeneous programming models and isn't limited to or designed for (say) MapReduce; you can invoke functions in parallel in a dynamic fashion pretty easily. [1] You can get an idea of that here: https://ray.readthedocs.io/en/latest/ https://ray.readthedocs.io/en/latest/
- jamesblonde 8y agoRay is, like YARN, a resource manager - but it's lightweight and fast. Kind of like Erlang is for a lightweight process model, Ray can start tens of thousands of containers with much lower latency than YARN or Mesos or Kubernetes. It has found a niche in distributed Reinforcement Learning (deep RL), and is establishing a beach-head there. Now it faces the chasm....
- rmbeard 8y agoUnclear what this is good for.
- AdamM12 8y agoIt is a WIP distributed Pandas implementation. Allows you to spread massive dataframes, their data and computations, over multiple machines vs. Pandas which is only local to your machine. It's not quite there yet [1] [1] http://modin.readthedocs.io/en/latest/pandas_on_ray.html#using-pandas-on-ray-on-a-cluster http://modin.readthedocs.io/en/latest/pandas_on_ray.html#usi...
- nerdponx 8y agoPandas replacing PySpark? Sign me up.
- AdamM12 8y agoIt looks like that would be the long term goal. Haven't used Spark enough to understand it's downsides.
- makmanalp 8y agoDask is already this! They have a dataframe replacement, a numpy array replacement, and some lower level primitives like dask.delayed too. Plus, the nice thing is that it's already being used with large amounts of data, and the warts (which were plentiful two years ago) are rapidly reducing. http://matthewrocklin.com/blog/work/2018/06/26/dask-scaling-limits http://matthewrocklin.com/blog/work/2018/06/26/dask-scaling-...
- axiom92 8y agoThis could be really helpful for implementations that were written with relatively smaller datasets in mind but now need to be scaled up. However, for someone starting from scratch, it is not clear what advantages do they plan to offer against Spark used with the Dataframe API.
- miggyrozay 8y agoHow does this compare to dask.distributed? Dask dataframes are also a wrapper on pandas API. edit- They explain differences in a section of this blog post: https://rise.cs.berkeley.edu/blog/pandas-on-ray/ https://rise.cs.berkeley.edu/blog/pandas-on-ray/
- guard0g 8y agoThis looks interesting. Thanks for sharing and will have my DS team try it out.