5 ms·
What people think of as Paxos was not originally intended to be called Paxos by Lamport, but the leader election phase for what is now called Multi-Paxos (and h
by _benedict 5y ago
What people think of as Paxos was not originally intended to be called Paxos by Lamport, but the leader election phase for what is now called Multi-Paxos (and he intended to call Paxos). So I don't think there's a difference for algorithmic heritage between Paxos and Multi-Paxos (perhaps you know this, but for those reading your comment it might help clarify).
I think it is anyway reasonable to say that most consensus algorithms derive concepts from Paxos, since (despite how terribly it was presented by Lamport) it is the consensus protocol that captured most attention. Most recent advances in distributed consensus derive from Paxos, not Viewstamped Replication. As far as I know all leaderless distributed consensus protocols derive from Paxos, and most protocol optimisations that have been developed apply to Paxos or one of its derivatives.
I also happen to think Paxos is pretty easy to implement correctly, particularly by comparison to the other protocols, in large part due to its active replication semantics, permitting that commands may be processed by replicas in any order. This means failover is much less complicated to negotiate, as the new leader does not expect to have a complete view of the log. Though of course membership changes remain complicated, and it may be beneficial for the leader to be able to assume it has a complete view of the log - but this is an optimisation rather than an inherent property for correct implementation.
> correctness unfortunately rapidly breaks down when physical disks are used, since these may misdirect reads and writes
What do you mean by "misdirect" here?
- _vvhw 5y agoI would say that the majority of leader-based state machine replication protocols today are based on the more popular RAFT, which is itself a derivative of Viewstamped Replication. At least most leader-based consensus implementations that I know of seem to be derived from RAFT these days? I would also guess that leaderless distributed protocols derived from Paxos tend to be more niche than leader-based state machine replication? While Viewstamped Replication as contributed by Brian M. Oki certainly established all the foundational elements of consensus that Paxos would later reiterate and generalize, Viewstamped Replication also immediately showed how to use consensus to do leader-based state machine replication, and to do this simply and practically, as seen in RAFT's massive success in industry. > What do you mean by "misdirect" here? A misdirected read/write I/O. This is a rare kind of storage fault in the literature where the disk device redirects the I/O to a different sector. It's rare, but it happens. > I also happen to think Paxos is pretty easy to implement correctly What was especially interesting about the PAR paper from UW-Madison was how even fairly common single disk sector faults such as latent sector errors or corruption could cause global cluster data loss for implementations based on Paxos and RAFT as specified, where they depend on pristine fault-free stable storage for correctness. The PAR paper further motivated why simple checksums are not enough and why local storage and global protocol need to be aware of each other if cluster availability is to be maximized. Implementing Paxos correctly in the presence of a realistic storage fault model is certainly non-trivial. I believe Viewstamped Replication (VSR), and in particular the 2012 revision by Liskov and Cowling, is a much better place to start when implementing a state machine replication protocol, since it places absolutely no demands on the disk for correctness of the protocol. In fact, this was the key reason we picked VSR for TigerBeetle, it's not only elegant, but it's one of the very few protocols that actually make sense in a production context where a realistic storage fault model is at play. For example, the VSR view change is entirely correct without any reliance on disk, whereas most other protocols would need significant deviations (and proofs) to make the same guarantee.
- _benedict 5y ago> I would say that the majority of leader-based state machine replication protocols today are based on the more popular RAFT, which is itself a derivative of Viewstamped Replication. In production systems, most software appears to use Raft. In academia much more significant advances have been made for Paxos, and most of the research I am aware of is invested into Paxos. I am not really aware of much protocol advancement in academia that is built upon Raft? > Implementing Paxos correctly in the presence of a realistic storage fault model is certainly non-trivial. I disagree. The failure model for Paxos as specified assumes that processes fail-stop. So all that is required for correctness is that a process fail-stops in the presence of a detectable disk fault. To detect the kind of fault you report here, we only require checksums - either that are in part computed upon data that is not stored in the sector in question, or that are computed over multiple sectors (or both). This latter criterion holds pretty commonly, I think. Of course, all distributed consensus protocols must handle reconfiguration in order to replace faulty processes, which would need to be invoked here. So this does not really seem any more challenging? Certainly the software I maintain (Apache Cassandra), that currently employs Paxos is protected against this kind of fault by virtue of the chunk-level checksums that span many sectors, and the property that processes will not respond to operations that encounter any kind of local fault.
- _vvhw 5y ago"So all that is required for correctness is that a process fail-stops in the presence of a detectable disk fault." Yes, except the PAR paper goes into many more storage faults. Recovering from these is not trivial and needs to be integrated into the consensus protocol, something not specified by Paxos. For example, how can a single replica even detect that it had a disk fault if the last I/O it acknowledged externally was simply never written to disk at all, despite an fsync barrier? To recover from this correctly at the consensus layer, you would need your consensus protocol to also include and run a recovery protocol at startup, akin to Viewstamped Replication's 2012 Recovery Protocol. Of course there are details I am leaving out here, but the point should hopefully be clear, it's not obvious or trivial when a storage fault model is assumed. I would highly recommend the PAR paper to you, if you have not incorporated it already into Cassandra. And because this is not only about correctness, you also want your system to be able to recover and to be highly available in the presence of storage faults. Otherwise, the redundancy afforded by the consensus protocol and replication is not really being fully utilized and resources are being wasted.