Why use Paxos instead of Raft?
neon.tech
neon.tech
I was fairly junior at the time, and Raft seemed much more approachable, but after really forcing myself to read and understand the Paxos paper, I see what he meant. I am pretty sure most of the love for Raft was that the original whitepaper was just a better presentation. The actual Paxos algorithm is quite simple.
If you go into Raft already having mastered Paxos (as this DE was doing), it's clear that Raft is basically isomorphic to a special case of Paxos.
This paper argues basically argues that raft has a different leadership election mechanism than paxos, but that if you tweak some terminology, and make a few relatively reasonable implementation choices for paxos they are otherwise pretty equivalent.
It even gives a raft style single page description of paxos (using raft style terminology), and shows how little it differs from the equivalent single page summary of raft.
The main implementation choices they use are: - combined roles into a single server role - enforce that log messages are decided in sequence (largely to avoid the having to specify the behavior of newly elected leader to propose operations for the gaps (possibly no-ops)) - numeric ballot number, rather than lexicographical pair (but this changes nothing except making the summary slightly easier to express)
We should start considering CRDT/OT/VCS/Diffsync approaches to distributed systems as well. They present a very nice alternative approach: whereas PAXOS/RAFT implement a consistent "distributed state machine", a CRDT, OT, VCS, or Diffsync system implements consistent "distributed state", upon which one can build a machine as a function of the state.
This latter approach is actually simpler, IMO, because it encapsulates all the challenge of distributed consistency within a smaller subset of the problem — state synchronization. This makes it more generally re-usable. When you create a system, you can just use an off-the-shelf library & algorithm to synchronize your data over a network, and then write synchronous functions on top of that to represent the system you want, however you want, without having to understand PAXOS/RAFT.
Paxos is complicated, but it’s well studied and proven.
"There are three types of consistent distributed systems: paxos, broken protocols, and single points of failures."
It's almost more impressive that people know who it is without saying the name.
knowing who said this allows readers to be able to look up more of the (unfortunately scant) content publicly available from him.
After reading the blogs about how good Erlang's concurrency model is and how we just just made a super implementation of it in XXX I have been led to formulate Virding's First Rule of Programming:
Any sufficiently complicated concurrent program in another language contains an ad hoc informally-specified bug-ridden slow implementation of half of Erlang.
This is, of course, a mild travesty of Greenspun (*) but I think it is fundamental enough to be my first rule, not the tenth.
[0] http://erlang.org/pipermail/erlang-questions/2008-January/03...
I used observers with Gluster previously and went from annoying split brain scenarios to flawless clusters just by adding a few, and their resource usage was basically nothing.
What's Neon's point of view about transient state in nodes? Is there a world where serverless client connections are stateless, or is the set up overhead not expected to be worth the cost?
However, read-only nodes require less coordination, and we have way more freedom there, so read-only Postgres as a function seems to be a more feasible concept.
We have some encouraging early results, but haven't committed to a particular technology (like cloud hypervisor) yet.
> Right now, such a change requires humans to be in the loop to ensure that the old safekeeper is actually down. It is on our roadmap to automate this procedure.
If you do implement this (which I don't recommend), be certain to also model it with TLA+. This level of automation, IMO, requires a human in the loop + a ton of visibility tracking on when it is happening.
A good way to roll it out is "semi-automation"—implement the automation but use it to ask a human to approve. The human will then do the normal (manual) verification. After you've run that successfully for a year, and your TLA+ model passes, you can then decide to fully automate without a human in the loop.
Otherwise, you're asking for an outage (caused by bad failover) IMO, and possibly data loss.
what happens when the compute server issues a read for something that has made it to the WAL servers but not the page servers?
I realize the answer depends on how big the cluster is, what state it is in at any given moment etc, but I'm happy to accept back of the envelope calculations/estimations!
Neon is more of an aurora approach detaching storage from the compute, you could scale up to more replicas and it could enable other functionality, though Postgres already can handle a pretty high replica count so you can scale out reads that way.
Citus (Cockroachdb, Yugabyte) has distributed compute which allows to engage multiple nodes per queries. This helps with analytical queries AND with scaling writes. But you lose out on compatibility and predictability of performance. Shared nothing systems are no longer Postgres.
"Observer" is not as well specified of a term [2], but from what I can find observer means non-voting replica which stores the log and a state machine replica.
Both of these make sense depending on the goals of the system. Observers make it easier to add new replicas without changing the size of the quorum, and witnesses make it cheaper to increase fault tolerance.
[1]: https://lamport.azurewebsites.net/pubs/web-dsn-submission.pd... [2]: https://cse.buffalo.edu/tech-reports/2016-02.orig.pdf
And Citus is shared nothing architecture plus columnstores. This means the use case is analytics or mixed workloads. Citus people should comment on this of course.