CharmPy – A high-level parallel and distributed programming framework
charmpy.readthedocs.io
charmpy.readthedocs.io
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/
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.
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.
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.
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?
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.
Please, for the love of God, import names explicitly or use e.g. `import charmpy as cp` and subsequently `cp.foo` so that reading example code we get a better sense of your API without having to guess which names were possibly overwritten.
> Have you used CharmPy?
> You mean PyCharm?
> No, CharmPy!
> Are we talking about the same thing?
> No
> Oh, that's just confusing, then.
Either way, it's done.
What determines the number of processes used is the launcher (e.g. charmrun, or something like aprun or ibrun on other systems). During initialization, the charmpy runtime will figure out internally how many charmpy processes are active in the job.
With charmrun, you can launch multiple processes in one host, but also across multiple hosts (by ssh'ing into each one and spawning the processes). This is done automatically by charmrun assuming you specify a list of hosts (called nodelist, see http://charm.cs.illinois.edu/manuals/html/charm++/C.html). Again, the application code is not affected by this.
Similarly, on other systems you can launch charmpy applications with the system job launcher (e.g. aprun, sbatch, ibrun…). We have done so for example on Cray supercomputers. It is simple enough but we have to update the documentation to at least show an example of this.
We don't offer an API yet in charmpy to explicitly do things like parallel apply (but will soon). You can however look at `examples/parallel-map/par-map.py` in the source code which shows a simple example of how to do it with the current API and might be what you are looking for.