A Case Study in LASP and Distribution at Scale [video]
youtube.com
youtube.com
Working a lot in IoT fields, or even with onsite deployments of services, this pattern becomes a blocker. It's more difficult than traditional "setup a server, turn it on, install software" approach, especially for small deployments. There's a lot of need (IMHO) for various orchestration systems which are self-boostrapping and fully P2P for smaller teams deploying on clouds, IoT devices in a given locale, research labs etc.
I've been watching your work for a while, and it's great you're moving toward usable systems! Is there a good place to watch your work and/or engage with the community? I'm starting work on some new IoT devices, and while this won't be directly applicable I'd like to play around with some of the core ideas.
A public slack would be nice. It seemed like LASP might be dead for a while when looking at some of the project pages a few months ago. Understandable when you're focused on solving hard problems, but an FYI from an outside perspective. This talk is really helpful in seeing what you've learned and where it's going. Thanks again!
What's the right way to build an open community? What do you recommend?
PPDP '17 paper on scaling: http://christophermeiklejohn.com/publications/ppdp-2017-prep...
Ensuring monotonicity via types: http://prl.ccs.neu.edu/blog/2017/10/22/monotonicity-types-to...
From the video also enjoyed the details about how things ran on variety of cloud environments. The google one was funny: kubernetes in kubernetes on borg on hw virtualtization on hardware.
Is there more work that generalizes this idea? E.g. attach metadata to more generic 'messages' that operate on 'processes' and then have the processes define how (or if) the merge occurs, given the history of all messages at different cloned replicas? Can LASP support something like this?
Passing through a causal context seems like a good idea. I'm not that familiar with Erlang but perhaps it doesn't have to be fully transparent? I presume there is a subset of processes that would participate in the distributed implementation and each could be encapsulated by another LASP process that wraps/unwraps messages as they go out/in and handles the metadata?
Btw, some related ideas exist in Croquet/TeaTime as well (http://www.vpri.org/pdf/tr2003001_croq_collab.pdf)
Did you take a look at Scaleway's "Dedicated ARM cores"? They start at ~3EUR for 4c/2GB. https://www.scaleway.com/pricing/
Their instances come up in under a minute. They have a nice client https://github.com/scaleway/scaleway-cli and they are from the EU which is a great place to spend EU grants ;)
Not affiliated. Just a happy customer.
I've done some work on gossip systems in the past, http://gossiperl.com is the result of my research. Gossiperl was based on work I've done at Technicolor Virdata (shut down nearly couple of years ago). We've built a distributed device management / IoT / data ingestion platform consisting of over 70 VMs (EC2, OpenStack, SoftLayer). That was before Docker became popular, virtually everyone was thinking in terms of instances back then. These machines would hold different components of the platform: ZooKeeper, Kafka, Cassandra, some web servers, some hadoop with Samza jobs, load balancers, Spark and such. Our problem was the following: each of these components have certain dependencies. For example, to launch Samza jobs, one needs Yarn (the Hadoop one) and Kafka, to have Kafka, one needs ZooKeeper. If we were to launch these sequentially, that would take significant amount of time considering that each node would've get bootstrapped every single time from zero (base image with some common packages installed) using Chef and installing deb / rpm packages from the repos. What we put in production was a gossip layer written in ruby, 300 lines or so. Each node would announce just a minimum set of information: what role it belongs to, what id within the role it has, the address. Each component would know the count of the dependency it requires within the overlay to bootstrap itself. For example, in EC2, we would request all these different machines at once. Kafka would be bootstrapping at the same time as ZooKeeper, Hadoop would be bootstrapping alongside. Each machine, when bootstrapped, would advertise itself in the overlay and the overlay would trigger a Chef run with a hand crafted run list for the specific role it belonged to. So each node would effectively receive a notification about every new member and decide to take an action, or not. Once 5 ZKs are up, Kafka nodes would configure themselves for ZooKeeper and launch. Eventually Kafka cluster was up. Similar process would've happen on all other systems, eventually leading to a complete cluster of over 70 VMs running (from memory) about 30 different systems being completely operational. Databases, dashboards, MQTT brokers, TLS, whatnot. We used to launch this thing at least once a day. The system would usually become operational within under half an hour, unless EC2 was slacking off. Our gossip layer was trivial. In this sort of platform there are always certain nodes that should reachable from outside: web server, load balancer, mqtt broker. Each of those would become a seed, any other node would contact one of those public nodes and start participating.
From the capabilities perspective, the closest thing resembling that kind infrastructure today, is HashiCorp Consul. Our gossip from Virdata is essentially what the service catalog in Consul is, our Chef triggers is what watches in Consul are. With these two things, anybody can put up a distributed platform like what you are describing in your talk and what we've built at Virdata. There are obviously dirty details like, one needs to have a clear separation of installation, configuration and run of the systems within the deployment. The packages can be installed concurrently on different machines, application of the configuration triggers the start (or restart), system becomes operational.
Or do I completely miss the point of the talk. I'd like to hear more about your experiences with Mesos. You're not the first person claiming that it doesn't really scale as far as the maintainers suggest.
By the way, HyParView, good to know, I've missed this in my own research. Maybe it's time to dust off gossiperl.
* edit: wording