5 ms·
Using Promise Theory to solve the distributed consensus problem
- andras_gerlits 3y agoHow to decentralise and scale distributed consistency beyond the commonly accepted limits
- mrkeen 3y agoAs far as I can read into it, this is 95% metaphor and 5% "use a master database but call it a distributed system". > 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".
- andras_gerlits 3y agoThe CAP argument falls apart as soon as you decouple consistency from wall-clocks. Consistent systems only suffer from CAP limitations if they need strict serializability. 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.
- heavenlyblue 3y agoThere's 0 relativity involved in distributed databases
- andras_gerlits 3y agoYes, that's exactly the problem with the existing models and why CAP was formalised.
- heavenlyblue 3y agoNope, that's not true, and you neither solved nor explained it in this message
- mrkeen 3y ago> The CAP argument falls apart as soon as you decouple consistency from wall-clocks 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.
- andras_gerlits 3y agoCan I replicate people deterministically and control all their sensory inputs in this scenario?
- danbruc 3y agoCompared to the mathematical rigor that gets usually thrown at such problems, this is all pretty vague, but maybe this is just the informal version. I admittedly did not completely understand the idea, but I of course have thoughts about it. 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.
- andras_gerlits 3y agoI don't think I ever heard people accuse our paper of not being rigorous enough, but more than happy to listen to specific problems with it: https://www.researchgate.net/publication/359578461_Continuous_Integration_of_Data_Histories_into_Consistent_Namespaces https://www.researchgate.net/publication/359578461_Continuou...
- danbruc 3y agoI completely read the first half of it, hoping that everything would eventually make sense, that the pieces would fall into place, but that never happened. I only skimmed the second half and everything seems to just become more and more incoherent. There are some recognizable underlying themes but nothing of it makes really any sense. If I would have to guess, I would guess that ChatGPT generated that gibberish. For large section I could at least imagine that the ideas are just way over my head, especially since I have never heard anything about promise theory and did not read the references. But page 16 really convinced me, that it all is just nonsense - how on earth do we suddenly and out of nowhere end up with differentials, Fourier series, and Heisenberg's uncertainty relation? And while the paper superficially looks sophisticated and scientific with all the symbols and notation, there is no substance behind it. The symbols are really only used for providing short labels for all kinds of things but they are [essentially] never used to relate anythings, let alone to derive or prove something.
- sausagefeet 3y agoThis post never actually delivers on its claims. What about CAP or FLP does this resolve in a new way? 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.
- andras_gerlits 3y agoPACELC is an extension of CAP, so it suffers from the same problem of trying to apply a universal clock to the whole system. If you do that, you will suffer these limitations. With client-centric consistency, you can work around these problems and global order "unfolds" in a "just in time" manner. The consistency-levels can go all the way to SNAPSHOT.
- sausagefeet 3y agoCAP says nothing about a "universal clock over the whole system". CAP is about the decision that has to be made in some unit of the system, it could be the whole system or it could be a bit, at the point of an operation. It's physics, there is no way around it. You can make different decisions on the semantics your system needs, but if you have two nodes that physically cannot communicate but need to be consistent for a client to move forward, the client cannot move forward. Full stop. Could you please show a failure mode that this system can handle that CAP says is not possible?
- andras_gerlits 3y agoIt's ironic that so many of you are missing the rigour, as that is exactly the thing that "undoes" the CAP arguments. 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_Continuous_Integration_of_Data_Histories_into_Consistent_Namespaces 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 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.
- sausagefeet 3y agoMath isn't what's missing, but Mark's post is just a bunch of metaphor and no rigor. At the very least it could go over failure modes and shows how it alleviates them but other databases fail.
- andras_gerlits 3y agoWe do you one better. We show how all information can be made redundant via determinism and how that means you can supply multiple copies of them across parallel, redundant channels. 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-fulfil-the-promise-of-microservices-5e397cb12e63#7df1 https://itnext.io/how-simple-can-scale-your-sql-beat-cap-and...
- sausagefeet 3y agoThis post has no rigor. "What happens if a node dies? Well another takes its place". Ok, show me. This is hand waving.
- andras_gerlits 3y agoI also have a demo here where I show transactionality between two SQL-databases, a MySQL and a Postgres instance: https://www.youtube.com/watch?v=XJSSjY4szZE https://www.youtube.com/watch?v=XJSSjY4szZE 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 https://www.youtube.com/watch?v=R4_phLs4d_M
- sausagefeet 3y agoNo-one doubts if you put a message queue in front of your database, you can do what these videos show. The doubt is if this says anything interesting about distributed systems. At least these videos don't demonstrate that.
- andras_gerlits 3y agoSure, there are other ways of doing this, the demos don't prove what I say, they only show that the theory works on a practical level. The science-paper however, does show how this mechanism can scale consistency. 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.
- sausagefeet 3y agoWe already know this technique works, it's how any asynchronous database replication is implemented. 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.
- dustingetz 3y agoAs I understand it, Omniledger retrofits highly available multi-master writes, across your existing SQL tables on existing databases, in a way that is transparent to application devs, by intercepting JDBC API calls. This approach of intercepting JDBC calls allows an enterprise to unify 100s of siloed SQL instances across the enterprise into a single globally consistent enterpriser view, without changing any application, simply by intercepting JDBC. To do this without requiring all txns to go through a single master authority, which would obviously be too slow, with OmniLedger the sysadmin must supply a configuration that maps a single authority (i.e. single database), for each topic (i.e. table). Therefore, "wide" transactions (e.g. "change email for user", where user table is repeated 100s of times for each database in your enterprise) may need to coordinate with many databases before committing and the txn will be delayed as it negotiates with the user table authority (which OmniLedger coordinates transparently to the application via JDBC interception). And "narrow" transactions, say to a single application-specific topic/table, are routed to that single authority, which is likely co-located with the application, giving local performance for this narrow transaction (i.e. no routing it through a central master db). The difference between this and say Spanner, is that you don't have to migrate/rewrite your entire enterprise of applications to integrate Spanner – OmniLedger intercepts the JDBC calls your applications are already making, with no application source code changes needed!! Source: I've met with the founder. Andras, how did I do?
- andras_gerlits 3y agoAlmost there Dustin, thank you. This article is about the mental model behind scaling consistency across arbitrary geographical distances and how this model allows us to communicate time-information the same way we now communicate data. It's an intro to the science-paper, not the implementation. 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.
- tsimionescu 3y agoI've tried to follow this blog post, but when it gets to the nitty gritty of the solution, it seems self-contradictory. There are constant claims that this presents a solution to have all three items of CAP - exciting. But then the solution seems to only have consistency. 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.
- andras_gerlits 3y agoIt's clear for everyone that if you define Consistency via Linearizability, CAP-like problems will apply, as you're necessarily creating original information on a potentially remote node. That's not the issue. 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 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.