Roshi: A large-scale CRDT set implementation for timestamped events
github.com
github.com
Having worked with Cassandra & Redis extensively, this might work better with Cassandra. Cassandra already has sets, maps, list as datatypes and can be leveraged to build a CRDT without having to write a new database. Oh, well I'm guessing we're past that conversation? ;)
- a time series CRDT is... just an append-only log, sure you can call it a CRDT, but it doesn't really have any meaningful properties. So what's the point of calling this a CRDT?
- LWW Last Write Wins is also an odd choice, because the notion of "last" doesn't exist in a distributed system - last according to who? When? What about clock drift? If who, then do you rely upon an authority? Then you aren't a CRDT.
Finally, regardless of CRDT stuff, what benefits does this give over just using timestamps series into Redis directly?
I applaud and have called for more people to build databases and distributed systems tools, so please keep it up. But I'm a tad worried this is just trying to cash in on the emerging hype (finally!) around CRDTs. Could you explain yourself more?
> a time series CRDT is... just an append-only log, sure you can call it a CRDT, but it doesn't really have any meaningful properties. So what's the point of calling this a CRDT?
I'm confused by this question. I don't think it's true that "a time series CRDT is just an append-only log". In Roshi, each object is identified by a key, and LWWSet semantics are provided via the timestamp. It's a CRDT because writing the same key repeatedly with different timestamps results in a single "winner" being kept in storage; if it were an append-only log, that wouldn't be true. Right?
> LWW Last Write Wins is also an odd choice, because the notion of "last" doesn't exist in a distributed system - last according to who? When? What about clock drift? If who, then do you rely upon an authority? Then you aren't a CRDT.
I'm also confused by this question. LWWSet requires a concept of "last", it's even right there in the name. And the documentation clearly indicates that Roshi's "last" is clock time of the node that processes the request. As long as you ensure that time is unique and monotonic (usually by including things like node ID and a per-node monotonic counter) then it's absolutely sufficient to act as the timestamp component of a LWWSet. If there is clock drift between nodes, then a write W1 against some node may actually "beat" a later write W2 against another node with a slower clock, but in terms of correctness, that doesn't actually matter: the write conflict still resolves deterministically.
> What benefits does this give over just using timestamps series into Redis directly?
I'm also confused by this question, and I'm not quite sure how to answer it. The operation semantics, and the benefits they confer, are quite exhaustively described in the README. You might also check the READMEs in the cluster and farm subdirectories, which go into more detail. What specifically don't you understand?
That sounds like a great way to produce a bad UX where a computer with clock drift doesn't successfully write. With real world clock drift, this can be on the order of years especially on older hardware with a drained button cell on the mainboard or badly implemented GPS clocks. You should also be watching the newest timestamp you got for a key in each data node, and if it receives a user (not peer) set operation, it should take the max between the timestamps. Conflict resolution between peers is unchanged, but you avoid an entire class of bugs that are near-impossible to debug and yet incredibly frustrating.
From a quick glance at the readme, it doesn't seem that you do this. If the client timestamps come from your web servers, and they're all on NTP, then it won't be too bad, and the effect will only happen for the maximum NTP offset.
EDIT: did some thinking and reread what I wrote
That's true, but any reasonable datacenter operation is going to detect clock skew beyond some tolerance (order of seconds) and either correct it, or decomm the nodes until it can be corrected. Pragmatically, data systems like Roshi have to operate within that tolerance, but don't need to accommodate huge differentials. Or, maybe more accurately stated: it's OK if they have weird apparent user behavior in those circumstances, as long as it doesn't throw the whole system into a logically nondeterministic state. Which is the case here.
> If the client timestamps come from your web servers, and they're all on NTP, then it won't be too bad, and the effect will only happen for the maximum NTP offset.
That's right. And it's important to zoom out and look at the actual consequences of missed writes at the application layer. This stuff only comes into play when a specific user makes different actions on a specific key, i.e. if a SoundCloud user likes and then immediately un-likes a track. It could be that the like "beats" the un-like, due to clock skew, and when the user reloads the stream, they still see the track with a heart next to it. "That's annoying," they say, and un-like it again, which almost certainly works. No big deal.
It takes time to recognize such failures. If the clock skew is only off by a few seconds or even minutes, the impact isn't so horrible, but if the clock jumps ahead by 20 years, you've effectively prevented any writes to that key for 20 years. Even if you only consider a few days, that's long enough that it can result in lost business.
> logically nondeterministic state
timestamps should only be used to tiebreak between logically observed events. You've somehow combined a timestamp field as a tiebreaker, as client data to be stored, and as a request id. If you separate them then adjusting the timestamp as a tiebreaker (along with node and counter) will remain logically consistent.
> And it's important to zoom out and look at the actual consequences of missed writes at the application layer
Lets. You wrote a distributed data store that can lose writes under certain circumstances. If the UX impact is negligible for you, that's great. However, it may be catastrophic for another use (say storing the latest price read for a stock option). For developers to use your database, it means they need to do an impact analysis for every single use. Just wait until someone decides the API is easy enough to use that they try to put payment data in there.
That's true, but (luckily) it's easy to remediate those objects in the Roshi architecture.
> If the UX impact [of lost writes] is negligible for you, that's great.
One interesting thing that I've observed about this kind of data system work is that there seems to be a as-of-yet unsolvable tension between "general purpose applicability" and "specific purpose usability". That is, there isn't a way to design a CRDT-based causally-consistent data system which has a general purpose API that consumers can use for a suitably broad set of use cases. Instead, success stories like Roshi seem to require specializing both data semantics and APIs to the specific use cases they target.
So, yes, if someone can't tolerate/remediate lost writes due to clock drift, then Roshi is not a suitable system. But, that's fine: data systems don't have to be general purpose to be useful, or to advance the state of the art.
Dealing with this UX stuff is very important, and we have distributed systems correctness tests (https://github.com/gundb/panic-server) that does split brain simulation and recovery, and code to handle timestamp skew even if NTP is running.
If you pass this off to developers to be responsible for, they won't know the difference between mistakes in your system and mistakes in theirs. They shouldn't need to be a distributed systems expert to use the database.
As a noob, some kind of Starting Assumptions and Context would be cool for all these distributed infra projects. For instance, I wouldn't have known about mitigating time skews in data centers.
Or perhaps even a reference deployment, image(s), playbook, or whatever the devops people are now calling it.
You are correct. Append-only log as an implementation detail doesn't converge and cannot be a CRDT. Although semantically logging can and should be a CRDT, as it doesn't need to be append-only at all.
> As long as you ensure that time is unique and monotonic (usually by including things like node ID and a per-node monotonic counter)
As other comment noted, it's not enough. You need some causality here, at least a lamport timestamp (which is a rule of thumb on timestamps anyway). Because even with synced clocks some will always be faster and always win writes.
That's true, but, in the use cases that Roshi targets, it's also not a problem. Users will notice and re-issue writes.
>Client timestamps are assumed to correctly represent the physical order of events coming into the system. Incorrect client timestamps might lead to values of a client either never appearing or always overriding other values in a set.
So basically, the clients provide timestamps that are treated to be true. There's probably a lot of use cases where this is still useful, though I agree regarding wanting to hear more about what advantages from CRDTs this is meant to gain.
I thought it was interesting and wanted to hear what people on HN thought. Specifically, why this (on its face) very cool looking piece of software didn't get so much attention when it was actively developed some time ago.