A distributed data processing library for CL
Monomyth is a distributed data processing system built using Common Lisp and based heavily on Broadway. It is designed to split the messaging systems into two, one defined and controlled by Monomyth, and one defined and controlled by the user. The messaging controlled by the user pertains only to data, which moves between persisted data streams and node threads. The structure and manipulation of this data is largely defined by the user (though there are data stream specific aspects certain node types might handle). Monomyth itself handles all aspects of system orchestration via the Monomyth Orchestration Protocol (MMOP). The work itself is done on a group of distributed workers that use concurrent, user defined nodes to process the data and are controlled by a single master server.
Design documentation can be found at
To run the tests you will need to have
qlot installed, the tests
have been verified on SBCL 2.0.3.
Then perform the following:
source test_env.sh docker-compose up -d qlot install qlot exec ./bin/test.ros
The example right now is pretty arbitrary, both in set up and the computations
In one process start
qlot exec ./bin/example-master.ros, and then in a few
qlot exec ./bin/example-worker.ros.
Once all the workers have connected to the master (you will see log entries),
press the return button in the first process.
This project has its basic architecture set up, but lacks most of the functionality needed for a 1.0 release.
The features likely to be targeted for a 1.0 release:
- Support for sending the results of a node to multiple queues.
- DSL support for sending multiple nodes to a single queue.
- DSL support filtering which message goes to which queue.
- DSL support for iteration.
- Support for nodes not sending messages anywhere at all (for loading nodes).
- Support for nodes not picking messages off queue (for generators).
- Basic heartbeat and restart for workers.
- Basic monitoring of workers and nodes, possibly data flow.
- More fine grained control of node threads.
- Startup CLI.
- Possible (but very unlikely) Kafka support.