3 ms·
If you read through the comments you'll find some contradictory explanations of what Dask does, so as some background: At its bottom layer, Dask takes a graph
by itamarst 5y ago
If you read through the comments you'll find some contradictory explanations of what Dask does, so as some background:
At its bottom layer, Dask takes a graph of tasks and dispatches them to a scheduler. Schedulers can range from a few threads, a few processes (using multiprocessing), a few process on local machine (with the Distributed scheduler), or tens of thousands of machines.
That latter example is not made up, I've seen geoscience demos doing that in 3 lines of code on one of their giant clusters.
The task graph can be generic Python functions, using e.g. the bag API, which gives you a map-reduce-y like API: https://docs.dask.org/en/latest/bag.html https://docs.dask.org/en/latest/bag.html
Separate, and optionally, Dask also has emulation layers for NumPy, Pandas, and scikit-learn I believe as separate project. Here you write code that looks like normal Pandas code, say, but instead of executing immediately it creates a execution graph, and then you submit the graph to your scheduler and it runs it for you. These emulation APIs are partial, by their nature, but also optional.
One interesting use case for this emulation API is processing larger-than-memory datasets. In ordrer to support multiple processes, and even more so multiple machines, Dask reimplements the NumPy and Pandas APIs using batching underneath. And so you can run a single machine and take advantage of that batching to process data that wouldn't fit in memory when using normal Pandas APIs, while also getting access (optionally) to multiple CPUs: https://pythonspeed.com/articles/faster-pandas-dask/ https://pythonspeed.com/articles/faster-pandas-dask/
(If you specifically want Pandas APIs, there are other alternatives mentioned here in there in comments, e.g. Modin.)
My main personal experience was using the lower-level API to run image processing code in parallel, on a single machine with multiple worker processes. It worked great.