4 ms·
Can you please lay out the differences between this and Dask? https://dask.pydata.org/en/latest/ https://dask.pydata.org/en/latest/ I work on a parallel progr
by devxpy 8y ago
Can you please lay out the differences between this and Dask?
https://dask.pydata.org/en/latest/ https://dask.pydata.org/en/latest/
I work on a parallel programming framework for python myself. Not geared towards performance, but the ease of use.
http://zproc.readthedocs.io/en/latest/ http://zproc.readthedocs.io/en/latest/
- eslaught 8y agoAs far as I know, Dask is at its core a tasking model (i.e. tasks have input and outputs, and run automatically when inputs are ready in a dataflow-like model). Charm++ (on which CharmPy is based) is an actor model. Think Erlang for HPC. You've got a set of objects that are all nominally running concurrently, and objects can send messages to one another. Personally, I prefer the task-based model (but of course I'm biased since I work on one myself). In a proper task-based model, you can't have races or deadlocks, everything looks to a first approximation to be sequential. In actor models there's pretty much no way to hide the conncurrency, and all the traditional pitfalls of parallel programming are exposed to the user.
- juanjgalvez 8y agoThere are quite a few differences between them. Disclaimer: I work on CharmPy, and I'm not an expert on Dask, so my comments might be biased and not entirely accurate with respect to Dask. Obvious difference between the two is programming style. CharmPy (its current core API) is based on asynchronous method execution between distributed objects. Being objects they can have state and data which allows for a lot of flexibility. In Dask, you express a workflow as a series of dependent tasks (which as far as I know are stateless so it's more like functional programming) and dask schedules it for you. The scheduling is centralized (done in only one place, so it's like a master-worker pattern) even if you use the "distributed" scheduler (which is needed for multi-node runs). With CharmPy you can have truly distributed applications. Another thing I observed with the dask model is that, since everything needs to be translated into a task graph before execution, there seems to be poor support for mutable distributed numpy arrays. A mutation operation like modifying a single element of a distributed array is not allowed as far as I know (I have tried), and other mutation operations that are supported actually generate a completely different task graph as a result, with the overhead this entails. In charmpy, this restriction does not exist since you can just invoke a method on the object that holds the data that you want to modify, and do it in-place. In terms of performance, our initial tests have shown huge performance difference, with CharmPy being up to 200x faster (this is comparing with dask distributed scheduler for a very simple BSP-style program). Of course, difference will vary by workload, but one thing to note is that Dask is pure-python, while CharmPy's core runs in C/C++. The current level of task granularity that we can comfortably support is a few hundred microseconds, and we expect to improve it further. In contrast, the Dask documentation for the distributed scheduler explicitly warns against using small task granularity, recommending tasks larger than 100 ms duration. And something like Jug recommends tasks longer than 20 seconds. We are planning on adding other APIs on top of the core charmpy API, to accommodate other programming styles. For example, offer better support for the functional parallel programming style (there is a small example of parallel map in the codebase using charmpy), or task scheduling.
- p1esk 8y agoI had a task recently where I needed to convert several million audio files from one format to another, and I did it with python's multiprocessing module (similar to this: https://stackoverflow.com/questions/50662610/using-multiprocessing-to-batch-convert-wav-to-flac-python-pydub https://stackoverflow.com/questions/50662610/using-multiproc... ) Just like the poster of that question on SO, I'm wondering if that's the best way (in terms of speed or ease of use). Do any of the third party libraries (like yours) offer any advantages for this use case? To clarify, I'm only talking about doing work on a single workstation.
- juanjgalvez 8y agoFor a single workstation and the task you describe, the pool.map() functionality of the multiprocessing module should be perfectly adequate. Not sure how scheduling overhead would compare between charmpy and multiprocessing, but for this task it shouldn't matter (I assume you need at least a second to convert one file, and even if the conversion is faster, you can chunk the tasks anyway to mask overhead). I would say the big difference for this task is if you want to run it in parallel on multiple hosts, which pool.map can't do. With charmpy we can provide a distributed parallel map offering the same or similar API as pool.map. There is a simple example in 'examples/parallel-map/par-map.py', but we are working on offering a library on top of charmpy with more features and a solid API.
- p1esk 8y agoOh, good point about batching - my files were really small (audio samples for speech recognition), so a conversion of a single file took a lot less than a second. I looked at the par-map.py example, however I can't quite understand where do I enter a server IP or something like that. The whole process is fuzzy to be honest. What do I need to do if I want to run my conversion task on two local workstations? E.g. I install CharmPy on both, then what?
- juanjgalvez 8y agoYou don't actually have to specify hosts or addresses in your application code. When the application starts, the runtime will know how many processes there are and on which hosts. The key is to use a job launcher. For the par-map.py example, suppose you want to run it on 4 hosts and 8 processes per host. One way to do this is by launching the application with "charmrun". First, install charmpy on all hosts like you said. Then you would create a nodelist file with the names or addresses of the 4 hosts. Finally, launch like this: `$ charmrun +p32 par-map.py ++nodelist mynodelist.txt` I have updated the "Running" section of the docs to try to explain this better, also pointing to the charmrun manual. Hopefully things are clearer now.