4 ms·
> 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 Vi
by _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.
- _benedict 5y ago> 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? Could you explain how the paper addresses this particular kind of fault, as I cannot find it? This isn't one of the fault categories listed in Table 2, and the paper assumes that the storage layer reliably detects faults. I cannot see how a local system by itself can ever reliably detect this class of fault? I can imagine coordinators detecting that a response from a replica is not up-to-date, but that happens anyway in most consensus protocols and this doesn't help correctness if only f+1 replicas had persisted the record (and another f have persisted a record with a conflicting outcome) >I would highly recommend the PAR paper to you, if you have not incorporated it already into Cassandra. I have read (well, skimmed) the PAR paper a few times, and while it is a nice paper I am unconvinced it is particularly helpful, beyond reiterating the fact that you cannot trust your disks - but this problem has to anyway be addressed by databases. I think that the kind of active recovery it outlines could even be counter-productive in some cases. Once a disk fault is detected, it is unclear that it is safe or helpful to try to replicate correct data to the node. It may have less storage available to it now, if it has sensibly isolated the faulty disk, or may be unable to persist the data and so the recovery process may simply result in a lot of recurring work for the system. Probably the most general purpose course of action is to replace the node, since nodes are plentiful. The faulty node can then have its disk replaced before being returned to the pool of available nodes.
- _vvhw 5y ago> Probably the most general purpose course of action is to replace the node, since nodes are plentiful. The faulty node can then have its disk replaced before being returned to the pool of available nodes. This is expensive and wasteful, and could potentially lead to cascading failure. Worse, what if all nodes have a single disk sector error in different places on disk? The cluster would be lost.
- _benedict 5y ago> This is expensive and wasteful. Recoverable disk faults of course are handled at the storage layer, so we're talking here about unrecoverable faults. My risk model prefers replacing any disk that is encountering unrecoverable errors, to avoid the disk causing additional problems as it degrades (no matter how good your error detection is, it is preferable not to lean on it more often than necessary). To ensure only the affected data is concerned, it is possible to relocate only that disk's data to a new replica, to minimise the burden. I'll concede that a better model might permit a replica some number of unrecoverable errors (few, but more than one in some time interval) before a total replacement is performed, but in this case there are still repair mechanisms to bring the replica up-to-date, and simply performing normal recovery for consensus operations whose state it cannot recover before it participates in further decisions is sufficient for ensuring correctness. > could potentially lead to cascading failure Trying to correct faults on broadly failing disks can lead to cascading failure, as the cluster becomes preoccupied with correcting faults that cannot be resolved, and user queries become unserviceable. Bootstrapping a new replica is also a very modest CPU burden. > Worse, what if all nodes have a single disk sector error in different places on disk? This isn't a concern in practice from experience, and theoretically the risk is also minimal. Unrecoverable disk errors occur at a rate of fewer than one per many petabytes. This provides plenty of time to replace a replica - and if this isn't sufficient, you can increase your replication factor until you can survive the requisite number of coincident faults. > Could you explain how the paper addresses this particular kind of fault? I'm genuinely curious to hear your answer here if you have a moment. I can imagine some mechanisms to help recover from a single fault like this in a distributed fashion, but would love to understand what I might be missing from the paper.