Keeping CALM: When Distributed Consistency Is Easy
m-cacm.acm.org
m-cacm.acm.org
This seems to be saying that any algorithm with a monotonic output with respect to input information "has a consistent, coordination-free distributed implementation".
As I understand it, for data to be monotonic requires that the data is partially orderable.
CRDTs require partial ordering, as well as a merge() function so as to create a lattice.
This seems then that CRDT's have stronger requirements - this seems to make sense, since CRDT's are about sharing data, whereas this CALM theory is only talking about making a local decision.
Both CRDT and CALM are about sharing data to make a local decision which is globally consistent.
Both use an order relation over datasets to modelise the concept of "adding more info to some partial input or output".
I would say that the difference is on their focus. CALM determines the frontier where no coordination is required and provides general criteria (no need either to retract former output nor to hear everything there is to hear nor to know all the participants). CDRT provides a mean to meet this criteria. By taking the least upper bound of former partial results, CDRT ensures that the outcome is growing.
Good thing, that that's not every database ever cough
A 'delete' is then the removal of all edges that lead to the forgotten item.
My point is that 'deletion' in every dbms - be it SQL or NoSQL - is equivalent to logical negation of a statement, which in turn is equivalent to the non monotonic query example that they use throughout the article `¬∃ ("not exists")`.
There is a major logical difference between an explicit delete (negation in a closed world interpretation) and a forget (temporary inconsistency in an open world interpretation), in the sense that the latter can be repaired by subsequent messaging.
This means that you can't have easy distributed consistency with our current databases and data models.
This is partially addressed by temporal databases like datomic and juxt, however they too have delete/retract as logical negation, which makes their monotonicity properties limited to the set of 'change operations' instead of the data itself, which in turn makes it hard to work with. If their monotonicity were on the data-model layer the entire database would essentially become a CRDT.
My experience has also shown that the effort for getting this right is vastly underestimated, initially and that eventually most of these many many applications end up with a set of bugs triggered by edge cases nobody thought about.
There's a whole ecosystem of minor catastrophes that don't quite warrant a blog post. Do we need a name for the Chesterton house on the corner of Dunning Way and Kruger Lane?
At the same time, it's really hard to learn why that fence is there until you at least see someone else get injured by ignoring the signs, so in that respect TDWTF probably taught a lot of people about quite a few fences.
Coordination as used in this article basically means "needs to be decided by an authority". But because of CALM we know that in order to elect an authority/leader, you need a consensus algorithm that actually has monotonicity at it's core, like raft.
There's a really great talk on this by Heidi Howard: https://www.youtube.com/watch?v=KTHOwgpMIiU
I actually don't think that database isolation levels are a good way to talk about this stuff, because they are practical guarantees, not theoretical ones. Snapshot isolation for example also requires a leader.
From a theoretical point of view there is no in between (at least none that I can see), either you are eventually consistent or you're not, either you are monotonic or you're not.
In which case yes, I am conflating the two. I’m saying the paper is making a distinction that doesn’t exist. Concurrent reachability is much the same problem everywhere, even between mark and sweep and reference counting. It’s all thread safe graph theory.
They talks about a more general concept of consistency that is per se unrelated to graph algorithms.
The examples they have use to illustrate the idea are deadlock detection and garbage collection in a distributed setting, which are two very distinct problems, in that the former will eventually converge, while the other does not necessarily due to the 2 generals problem (you could try to move all the information to a single node, and you can do so in practice, but because of the 2 generals problem you can't once your connections have a chance of being lossy).