Roshi: a CRDT system for timestamped events
developers.soundcloud.com
developers.soundcloud.com
Interesting. We have an event stream database implemented in Haskell, and this looks like an excellent way to index it. Especially associative, commutative and idempotent can (probably) all be encoded in the type system!
http://christophermeiklejohn.com/coq/2013/06/11/distributed-...
Also, see related papers from PaPEC '14:
Basically a lot of this:
But using Haskell's type system is an interesting idea.
You can start with one level higher than just looking for associative, commutative and idempotent, because that actually comes the definition of a semilattice
http://en.wikipedia.org/wiki/Semilattice
The crucial part is that you have an appropriate meet or join operator (and that it satisfies the above requirements).
Informally think of a functions like max() over natural numbers or a union() over sets. Those are some examples.
I don't know Haskell but found this interesting gist of someone who has tried this, and frankly a lot of stuff there is above my head, but it might help you:
I think some of the theory there is starting to get into operational transformation territory [1]. If that's the case then a really interesting application might be applying the same semantics indexing a stream of events to expressing events as a diff against state and propagating them to clients...
http://www.cs.indiana.edu/~lkuper/
Has tons of papers and
https://hackage.haskell.org/package/lvish
is the lib.
With CRDTs the merge operation is made such that the simplify operation is id.
It's also interesting to note that the merge operation is just a pullback in the appropriate category.
http://hackage.haskell.org/package/lattices-1.2.1/docs/Algeb...
Twitter's Algebird[1] has a wide range of monoid implementations (in Scala). It's a surprisingly useful type class for something so simple.
I have a talk[2] and some code [3] that goes into more detail on CRDTs and the connection to monoids.
[1:] https://github.com/twitter/algebird
[2:] http://noelwelsh.com/programming/2013/12/20/crdts-for-fun-an...
The read view just has to be caught up from when it was last accessed. You don't do this all the time, only when a user requests their timeline. So still an individual CDRT set still has to be persisted for each individual user a bit like an out-of-date inbox. You don't create the whole inbox from scratch each request or do you?
maybe you do because you only need recent events??
---------------- EDIT.
I think I worked it out ... the inboxes are created dynamically. This is what was meant by stateless. All the time series data is able to fit into one server's memory, so the IO overhead of assemblage is low
Weirdly enough I am implementing a CRDT list today (not a set). The TreeDoc example in the main paper has a problem I think, if two writes with same key go in, the disambiguator keeps an ordering. But if that occurs, you now can't insert between those two edits because the disambiguators are not dense.
The set stuff looks good though.
After glossing over the readme of github.com/soundcloud/roshi I came away with the impression that Roshi's LWW-element-set implementation does not strictly adhere to the qualities of a CRDT, mainly that operations must be commutative.
Roshi documents two uses-cases:
If first we apply an add operation to the set, then apply a remove operation with the same tuple, the resulting set will ignore the remove operation: A(a,1) R() + remove(a,1) = A(a,1) R()
If we apply the remove operation, then next the add operation, the resulting set will ignore the add operation: A() R(a,1) + add(a,1) = A() R(a,1)
This means the resulting state of the set depends on the order of operations, which violates the commutative property of CRDTs.
As a consequence, if we assume this CRDT is replicated on each redis node in a Roshi cluster, then there is a case (admittedly rare and short-lived) where a Roshi cluster cannot be eventually consistent:
Given no more future operations on a set with key K, if, node 1 contains a key K with set A(a,1) R(), and node 2 contains a key K with set A() R(a,1), then the system cannot be eventually consistent.
Roshi seems to have implemented a weak form of CRDT, which under rare and short-lived situations impedes eventual consistent. But given the operational realities, this is totally acceptable.
> Of course, reads are difficult. If you follow thousands of users, making thousands of simultaneous reads, time-sorting, merging, and cutting within a typical request-response deadline isn't trivial.
Even within a "typical request-response deadline" blasting through some thousands of "events" does sound kind of trivial? (Which I suppose is reflected in the relative simplicity and elegance of this solution :-)
Now, if you absolutely need the inboxes of user A and user B who follow the same thousand users to end up identical if they are accessed in a similar time-frame, then things do indeed sound a bit hairy -- but do you really need that? (Or even demand that both A and B see X re-sharing track N, before user Y re-shared track M -- if the sharing events are within the threshold of clock-skew internal to your systems)?
I suppose it's great to be able to say that these guarantees will hold, but I'm not sure I see them as needed for (all of) soundcloud's events?
> Even within a "typical request-response deadline"
> blasting through some thousands of "events" does sound
> kind of trivial?
If they were all on the same physical machine, sure. But the events are spread across multiple nodes, so a single stream read can potentially touch every Redis instance. > I suppose it's great to be able to say that these
> guarantees will hold, but I'm not sure I see them as
> needed for (all of) soundcloud's events?
I'm all about cheating, if it nets me a big win. How would you cheat, here? Honest question; it's not apparent to me what corners we could cut, given the whole point, more-or-less, is to avoid pre-materializing views.For example, (I'm told that) Twitter has effectively two separate systems for distributing updates for users, one for normal users (you and I), and one for celebrities. Propagating Lady Gaga's tweets to all 30 million followers can take something like five minutes to accomplish. It's even more fun when these celebrities start tweeting back and forth with each other.
Maybe SoundCloud's CRDT approach would work better for natural networks? At some point, network communication could start to be an issue with a compile-the-list-at-read design, but I could be wrong.
I personally don't see why a CRDT is necessary if you used something like Twitter's snowflake approach to assign ids. Pull in events is sorted order from each machine, merge, then push to the user. It's not clear to me what adding the CRDT abstraction adds to that besides overhead. Maybe there are other benefits I'm not grasping?
> At some point, network communication could start to be
> an issue with a compile-the-list-at-read design, but I
> could be wrong.
You're absolutely right. Optimizing for that bottleneck occupied most of the development time.At a very superficial level it seems to me that you can't implement a reorder operation within the model, but perhaps you can circumvent the problem with some tricks...