Using Promise Theory to solve the distributed consensus problem
mark-burgess-oslo-mb.medium.com
mark-burgess-oslo-mb.medium.com
The author, additionally, seems to not be well-versed in the existing distributed database literature.
Essentially they have added a queue in-front of all database operations. That queue is totally ordered, so you can't have consistent issues. That queue apparently does all operations via a two-phase commit (without calling it that and being unclear on semantics, so I am not 100% sure).
Ok. So you've moved the question of availability and consistency to your queue. Is that a single queue? If so, then it's a liability. Why not just use a single database at that point. Is it multiple queues? Then you still have a consensus problem to solve. Are you using two-phase commit? Well, now your availability is seriously impacted.
There is nothing there.
It's a shame because there are models that are much more interesting that provide a useful mental model for this. PACELC is my favorite. The essence is that when everything is going fine, your decision is around latency and consistency. When there is a partition, your question is between availability and consistency.
https://mark-burgess-oslo-mb.medium.com/using-promise-theory...
I haven’t done heavyweight distributed systems theory in a few years, and so who knows, maybe @aphyr shows up and Jepsen-stamps this thing, but after reading the paper and going over the commercial website personally I’m a pass. I’m one of about a zillion people on HN who have dealt in petabyte and trillion event days on logs as mundane subsystems for years on end and it’s just almost always: “there’s the lock, didn’t see it right away because of the bespoke nomenclature”.
History is replete with misunderstood geniuses who weren’t appreciated because they were so far ahead of everyone else (Galois is my favorite example personally), and we might all end up eating crow over it, but this has a wink and a grin vibe that I’m going to regard as great satire rather than hyperventilating. I suspect the object of the satire would be blockchain.
The consistency-levels can go all the way to SNAPSHOT.
Say we have a database of customers, with two replicas - one in the USA, one in Europe. Say a customer in Europe wants to update their shipping address. We ship products every month to this customer from the USA to their current shipping address, on the 10th of that month.
The customer is updating their address on the 9th at 10 AM PST. However, we are in the midst of a massive network partition that started on the 9th at 02 AM PST and is expected to last until the 11th.
Do we perform the update in the European replica and tell the customer it succeeded (giving up consistency)? If we do, the US side of the business will ship the products to the old address, even though the update happened a full day earlier.
Alternatively, do we tell the customer the update failed (giving up availability)? If we do, then the customer can't even let the European side of the business know of their new address.
This is the CAP theorem in a nutshell, and it is obviously inescapable. It doesn't require appeals to immediacy that are anywhere close to the bounds of relativity. And while 3-day long network partitions are quite rare, partitions that last for many hours are not.
That's exactly the point of this section in the essay: https://medium.com/p/5e397cb12e63#7df1
Networks will have latency-spikes, but if you can stream time-information the same way you can stream others, you can use redundancies to mitigate the disruption of any single channel.
CAP theorem beaten by declaring P to be unlikely.
https://medium.com/p/5e397cb12e63#7df1
Considering that this is maybe the 12th time I'm linking the explanation, I now think you're not looking for a discussion, you're here for a fight. Please let me know if you find any problems with the article in its reasoning. It was vetted by a lot of people before, so I really need to notify them too if you do.
Redundancy doesn't fix any problems, it just makes them less likely to occur. Again this not only does nothing to address the CAP theorem but is done for a wide array of problems already, there is nothing new here.
https://medium.com/p/5e397cb12e63#373c
There are three separate mitigation-solutions that go into how total-order strong consistency can keep marching on even if a specific child-group is intermittently isolated.
The first one is redundancy by deterministic replication, so there will always be many replicas which aren't just shards, but full copies of the _consistency ledger_. It's not a database, not a cache, just the thing that establishes the consistency between nodes. These instances all "race each other" to provide the first value of the outputs to other nodes.
The second one is the latency-mitigation we talked about earlier, I don't think we need to waste more breath on that.
The third one is that since the consistency-mechanism requires an explicit association-instruction to interleave the children's versions into its parent's (so that these versions can be exposed to nodes observing it from afar). If the child goes AWOL, it won't be able to associate its versions to its parent, so it won't keep up everyone else either. In this case the total-order won't be affected by the child-group's local-order, which is still allowed to make progress, as long as it's not trying to update any record that is distant to it.
2. Redundancy in network connection. Yes this can help. Again, it doesn't resolve anything CAP related, it just reduces the likelihood of a certain class of distributed systems fails. Note, there are LOTS of ways to have unbounded latency that this does not resolve. Anything from misconfigured routers to disk drives dying. Again, not resolving anything CAP states, just attempting to reduce the some probabilities.
3. If I understand this correctly, you are saying some data is "homed" in certain regions, and if that region becomes partitioned from the rest of the world, they can still modify the "homed" data there. They own it. This doesn't address anything CAP related.
Assuming I understand the architecture correctly, then yes, some architecture's will benefit from this. Some might not. But all of these decisions align with existing understanding and decisions architects make in the face of CAP.
It feels to me that you have believe CAP defines a very particular database architecture (something like a PostgreSQL) database, and your architecture addresses limitations in that, and thus your architecture solves limitations of CAP. But that's just not true. Take Riak, Cassandra, BigTable, Spanner, CockroachDB, all of these represent architectures defined in the face of CAP and they all have different trade-offs. They don't look a lot like PostgreSQL. But they cannot get around the simple laws of physics for how information is communicated.
But this is about the solution, not the science. The science is basically unlimited, its only real restriction is around our budget.
If you want to understand how the client-centric consistency mechanism takes care of these things, I write about it on medium and Twitter all the time.
But again, I don't feel like you owe me your time or attention
To summarise: we have a consistent system that works via one-way streaming, via redundant channels. Best possible cache-invalidation. Reads are both consistent and locally available. Upper-bound time of writes is predictable and doesn't suffer from blocking. I'm not sure what other improvements anyone could want from such a system, this is the best such systems can ever be. These improvements go well beyond the limits CAP sets out and it basically makes the whole argument moot.
Even if there were two different US nodes, and only one node lost connectivity to the others, the problem wouldn't change at all - the branch which only has access to the disconnected node would still either send shipments to the wrong address or stop being able to read the status at all.
https://blog.dtornow.com/the-cap-theorem.-the-bad-the-bad-th...
"Note that Gilbert and Lynch’s definition requires any non-failed nodenot just some non-failed nodeto be able to generate valid responses, even if that node is isolated from other nodes. That is an unnecessary and unrealistic restriction. In essence, that restriction renders consensus protocols not highly available, because only some nodes, the nodes in the majority, are able to respond."
Partition with a capital P doesn't really move me as our experiences contradict its assertions, as the model it bases its statements on is fundamentally broken.
So let's talk about partition in the real world. Practically speaking, partitions mean that some nodes experience high latency when communicating. If you need to update some record in a specific data-centre and you can only to talk to that place via a single channel and that channel is disrupted; yes, you're going to experience a slowdown or even halt. Nowadays however, even widely used consensus-algorithms can mitigate that, and the delinquent node will eventually be dropped from the group if it causes enough problems. We don't do anything very different in this regard. Since our deterministic ledger can be replicated across multiple data-centres (as no nodes in it create original information), and since observers will only rely on records the ledger already sent out and since the only reason you would need to talk to this ledger is if you need to modify a shared record with other observers, you can always pick a different ledger-instance, there are no "master copies" anywhere. Remember that we can stream time-information with the data, so the client can always calculate its "point in time" and reconstruct a (even globally) consistent view of the data it has.
Sure, if you choose to centralise all your data behind a single flaky connection, you're gonna have a bad time. The point is that the setup we built allows you to not need to do that and it does that transparently, behind SQL semantics.
For the specific quote you gave, that is an obvious assumption. A client only has access to some of the nodes in the distributed system. Of course we want any node to give the correct answer - the whole purpose is to reduce the burden on the client. The client is not responsible for searching all of the nodes. And note that the proof doesn't actually require that all nodes return the right answer - the contradiction is reached as long as all the nodes that the client has access to return the wrong answer.
Another bad claim in the article is that the proof of CAP requires that the partition is permanent. Maybe it's written like that for simplicity, but it obviously only requires the partition to be longer than the client's bound on response time. If the client is willing to wait an hour for a response, then any partition event that's two hours long will lead to the same conclusion. Since clients never have unbounded time to wait, and since partition duration is unbounded even if not permanent, then the argument still holds.
Also, major network outages that disconnect whole regions of the internet for hours from the rest of the world happen somewhat regularly (more than once a year). Whole AWS regions have become disconnected, ~half of Japan was disconnected for a few hours, Ukraine has been disconnected several times, etc. If you run a large distributed system for a significant amount of time, you will hit such events with probability approaching 1.
CAP is a bad model for more reasons than the ones listed in that article. My favourite one is that it requires Linearizability, which nobody does in SQL. The disconnect when saying "SQL is not consistent" to me is just too much. CAP is based on a badly defined idea that comes from a presentation that was wrong in what it said.
That you need to tolerate outages of entire regions is a good argument to make in itself, there's no need to point at CAP. My answer to that is that as there's a way to define consistency in a way that allows for it to manage partition problems more gracefully, and that is the model we show. If you require communication to establish consistency and stream the changes associated with the specific timeslot at the same time, partition means that the global state will move on without the changes from the partitioned areas and that they will show up once the region rejoins the global network. While separated, they can still do (SQL-) consistent reads and writes for (intra-region) records they can modify.
In contrast, any distributed system has this property (or it may refuse the query entirely) in the face of partitions.
> https://blog.dtornow.com/the-cap-theorem.-the-bad-the-bad-th...
Especially:
> The “Pick 2 out of 3” interpretation implies that network partitions are optional, something you can opt-in or opt-out of. However, network partitions are inevitable in a realistic system model. Therefore, the ability to tolerate partitions is a non-negotiable requirement, rather than an optional attribute.
> CAP requirements are absurd
Yes! Literally. One would roundly ridicule someone who claimed to have met (or exceeded) those requirements.
Could you please show a failure mode that this system can handle that CAP says is not possible?
This section discusses latency-spike mitigation (which is how Brewer defines CAP colloquially): https://itnext.io/how-simple-can-scale-your-sql-beat-cap-and...
This section dissects the problem when trying to apply CAP to non-linearizable systems like SQL: https://itnext.io/how-simple-can-scale-your-sql-beat-cap-and...
Again, if you're not happy with the lack of scientific rigour in this technical article (ie: not science-paper), you can connect the dots in this one: https://www.researchgate.net/publication/359578461_Continuou...
> This section discusses some failure modes (the first one is about the failure of a specific node): https://itnext.io/how-simple-can-scale-your-sql-beat-cap-and...
> In this setup, a copy of the node is replaceable by taking a (potentially days old) backup copy of it and replaying the events that happened since the time the backup was established.
The replaced node does NOT have access to the events that happened after the partition.
* If the replaced node serves up stale data anyway, you built an AP system.
* If the replaced node refuses to serve up stale data, you built a CP system.
* If you're pretending it has access to the latest events, there's no partition and you built a CA system.
https://medium.com/p/5e397cb12e63#04a5
This is the summary: "In other words: any distributed solution that fits the SQL standard can rightly claim that it scales SQL databases, and Brewer’s model can certainly accommodate a framework for that. His model however, is not the only kind of distributed SQL database that can exist, therefore his assertion that all distributed consistent systems must pick where they position themselves on his famous triangle is wrong. The system we explain here for example, is an exception. Formally: because our consistency model stays within the bounds of what the SQL standard allows and includes network communication; and informally because we can fine-tune latency variability according to the use-case of the specific datastore within the system and can even be reduced to only be a theoretical concern."
And you can also give the benefit of the doubt when allowing some number (less than the quorum) of nodes to fail (and letting them restart and catch up, etc.) while the system still makes progress ("CA mode"). After all, that's the point of distributing a system in the first place - there's no one master which can die and bring down the system.
But yeah, at this point I think OP is just going to keep implying that partitions don't happen or something...
> What replica sets using Paxos and Raft propose is that you will literally deal with one master server and try to keep backup copies of an entire database aligned independently. When a master dies or fails, the clients try to decide on a new leader so that they are all talking to the same server. This leads to a delay in which no one can write.
No comment on Raft, but "a master" only exists in Paxos if you invent one. Any client can talk to any Proposer (which will then try to reach a quorum of Acceptors & Learners). Losing just one Paxos node will not result in a partition, so you can still operate happily with availability and consistency.
> There's no particular need to replicate a whole database if we only want to share a few records. Granularity is the answer to scalability and reliability.
Lots of data, lots of CAP tradeoff. Not much data, not much CAP tradeoff.
> In IT, correct values are assumed to be the latest values. It’s a race to be last, because the last value wins by overwriting and obliterating what came before. So if you have an evil demon flooding the system with nonsense, you’re in trouble.
Paxos itself does not allow for overwriting of values (Learners do not change their mind about the value of a given key.) In order for 'obliteration' to occur, you need to augment your key with a version number or timestamp - something that a downstream system could interpret to mean as "happened after".
The point is that if you only look at the order of the data and not the wall-clock of some actor in the system, you can "calibrate" these separate ordering mechanisms together into a coherent whole and from that, you can build up all of the SQL guarantees.
Linearizability makes systems suffer because it tries to enforce a Newtonian model in a relativistic world. Order will naturally emerge faster when you're closer to the data than if you're farther. Measuring these with the same clock is what causes CAP, not some inherent property of distributed consistency.
Hogwash
Five friends are seated at a restaurant, about to agree on a flavour of pizza via a simple majority quorum. Before they do so, the three girls excuse themselves to the restroom while the two boys remain at the table. The waiter arrives and asks what flavour pizza they'll have.
Do the two boys answer on behalf of all five friends, or do they wait for the girls to return?
No-one checked their own watch.
First, there is no Availability, because the solution requires a centralized Interloper service. There is some handwaving that the Interloper is distributed, but all of the arguments about maintaining the three databases consistent with each other only work if the Interloper can be assumed to have (distributed) consistency. Which is of course obvious - if you've solved distributed consistency, you can use it to make other things consistent. But this doesn't explain how to solve distributed consistency in the slightest.
Secondly, the described system doesn't actually offer partition tolerance:
> If one of the databases becomes unavailable (e.g. if it loses power or its network connection creating a “partition”) no harmful misalignment can be observed by any client, because the interloper disallows reads until everything is reported to be back in sync.
So if one database goes down/is disconnected, the whole system grinds to a halt. So, no partition tolerance.
The paper has a different problem. It seems you are redefining the notion of consistency because you don't like the definition used in the CAP theorem. But you don't actually disagree that no system can display what the CAP theorem calls "consistency" at the same time as availability and partition tolerance, you just don't think it's necessary. This is a defensible position (and one often taken by many distributed databases), but there is no reason to undermine others' work instead of simply saying so.
As more of a side note, the paper you wrote keeps referring to the CAP theorem as a conjecture, but it has in fact been formally proven in 2002 by Gilbert and Lynch. You don't seem to have a refutation of their proof, which you don't even cite.
The issue is that people in practice almost never use Linearizability for their Consistency, for example I don't know of a single SQL-implementation that does full linearizability (or Strict Serializable, same difference). So the industry already means something totally different from CAP's definition, and for those consistency-levels CAP doesn't apply. In fact, I present a technical series of arguments here how you can overcome these limits in practice: https://medium.com/p/5e397cb12e63
There's a specific section about CAP in the end, but I talk about node-loss. replication, strong consistency and all the others also.
Instead of “equality at all times” we should be asking “alignment whenever someone actually looks”, because this is all we can promise about someone else’s state.
It seems to me that we generally do not know when and where someone will look, so we will still have to be consistent at all times and all places because all times and all places is when and where someone could look. If the idea is to delay consensus from the write to reads, then I do not see how this makes things any easier. Instead of ensuring that everyone receives a write when it happens, you now have to ensure that you can gather all the writes that happened when a read occurs.
Source: I've met with the founder. Andras, how did I do?
With regards to the implementation (which is omniledger.io): We basically make SQL scale via Kafka by piggybacking on the semantics of both technologies.
There's no central authority for any of the tables. The version-ledger is totally schema-agnostic, only the clients understand it. In fact, it federates schemas the same way it does records. Tables are not topics either, in fact, we scale by importing namespaces via a command-line interface, which specifies a schema. Any database can register to this namespace after which they become as much of the "master" to the namespace as any other. There's no hierarchy between them, the ledger does everything for them. Since the ledger itself is deterministically replicated, a separate instance of it can be co-located with each instance that runs the JDBC-connections to the database, raising availability of it to the availability of the whole system. Each component in the setup scales with the number of partitions created for it.
And another, where I show loose coupling, ie: that the system continues even when an instance goes down and that it catches up once restarted. https://www.youtube.com/watch?v=R4_phLs4d_M
If you're not constrained by CAP, I'd hardly expect you to be constrained by Two Generals either.
> And another, where I show loose coupling, ie: that the system continues even when an instance goes down and that it catches up once restarted
I've built one of those, it's fun! But I digress. Back to CAP.
1) How many instances do you have?
2) What is your quorum (simple majority?)
3) How many instances can go down before the system stops serving consumers?
Think about it this way: Our system simplifies running inter-system consistency and makes it much faster. Considering that Spanner exists, what proof would you accept to validate our claims?
I understand that reading all the material and putting them together isn't a small ask, but no new tech is easy to understand at first. Anyway, the material is there for anyone who cares to look and the system does what we claim it does.
I don't think anyone owes us their time to check our claims, but I don't think the fact that not lot of people will do this changes anything either.
I've read all of your content and I can see nothing even close to, for example, the Spanner or Amazon Dynamo paper which go through the operation details of how the systems work. Literally your articles are just a bunch of metaphors followed by hand waving. No operational details.
I don't know if you know you're selling snake oil or just don't understand what you're implementing, but even in this thread you've gotten plenty of feedback that how you describe what you're doing is not coherent, you might want to address that. Or not, I don't know, if you're selling like hot cakes then keep doing it.
Anyway, if math is what you guys are missing, it's in this science-paper linked in the article: https://www.researchgate.net/publication/359578461_Continuou...
This is a gentler intro to the concepts. You can also read my essay on why this setup works better than the often used semantics: https://medium.com/p/5e397cb12e63
There's a specific section at the end on why CAP only applies to a very specific subset of SQL databases.
My essay talks about this in detail around the latency-mitigation section and the failure-modes part, but these are questions much closer to the actual implementation. Mark discusses the new mental model behind it, I talk about technology.
https://itnext.io/how-simple-can-scale-your-sql-beat-cap-and...