6 ms·
> 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,
by _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.
- _vvhw 5y ago> 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. Your risk model is of course not exclusive of the PAR protocol. However, PAR would mean that a cluster can at least recover itself and then replace the disk in the background, without this being a showstopper event. The paper has some very simple examples of where a single disk sector failure can render a cluster unable to complete a view change and thus completely unavailable, but where PAR's distinction between uncommitted/committed ops can mean that the cluster can make forward progress and recover. There are also further optimizations that you can do on top of PAR once you get started, for example to support constant-time acks even from lagging followers which don't have a complete log, to help them catch up faster and reduce ack latency spikes, while still sticking to the VSR invariant of not allowing gaps in the committed log (as distinct from the uncommitted log). I won't go into further details, but you can see how we do this exactly in TigerBeetle (src/vsr/replica.zig). > 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. Theoretically, but in practice a bad batch could easily have higher failure rates. You would also need to make sure that your scrubbing rates are surfacing bad sectors soon enough and that you can then reconfigure faster than additional sectors fail. It's also rare to see test suites actually injecting storage faults and testing that the scrubbing system can detect and repair. This aspect of recovery is often simply not tested. However, as I mentioned before, PAR also shows cases where even a single local disk sector is enough to bring down a cluster that doesn't implement protocol-aware recovery for consensus storage. The principle is that local storage and global consensus protocol really do need to be integrated, unless maximizing availability or network efficiency are not primary goals of the system. > Could you explain how the paper addresses this particular kind of fault? Sure, this kind of fault is addressed by Viewstamped Replication's 2012 revision, not by PAR. See the Recovery Protocol in the paper by Liskov and Cowling. In the paper, it's used to recover RAM contents after a crash but the technique is also helpful for storage, and indeed the paper always has stable storage in mind, it's just that it doesn't require it for correctness as other consensus protocols tend to do. VSR is more of a "write-back cache" style protocol and I believe this makes it a better place for engineers to start from, because it's easier to work with disks, when you don't place much reliance on them.
- MatteoFrigo 5y ago> what if all nodes have a single disk sector error in different places on disk? Well, there are many ways to skin this particular cat, but pretty much all large-scale storage systems (think S3) use some kind of erasure coding. That is, the replication protocol replicates a "log" whose execution generates a "state" of the "state machine". The log is fully replicated, but the state is erasure coded. This scheme tolerates single-sector errors in the state. If you have an error in the log you throw away the entire replica, but since the state is 1000x larger than the log, this is not a big deal. At the end of the day the disk must guarantee something, or else you are screwed no matter what. For example, if an acceptor acknowledges a phase-1 (view change) message but the disk lies about storing the new ballot/view/term, then all these protocols are incorrect.
- _vvhw 5y ago> Well, there are many ways to skin this particular cat Sure, of course, but why shouldn't we implement the easy wins as well? PAR is pretty simple to implement. It's not exclusive to other techniques, and it increases availability significantly. It's not either/or but both I believe. > At the end of the day the disk must guarantee something, or else you are screwed no matter what. For example, if an acceptor acknowledges a phase-1 (view change) message but the disk lies about storing the new ballot/view/term, then all these protocols are incorrect. No, that's actually exactly my point about Viewstamped Replication as per the 2012 revision (and your example is precisely what I've had back in mind throughout this thread). VSR's view change is different from all the others, in that it does not require any guarantee from disk. The whole view change protocol is entirely in-memory. It's the one protocol (at least that I know of) that remains correct, even where all the others would be incorrect (since, contrary to VSR, they unfortunately require pristine stable storage for correctness). I think we are making the same point, you just missed this aspect of VSR, in that it places no reliance on disk (at all) for correctness. This sets it apart from all the others. VSR is closer to a "near-byzantine" model as I like to call it. You can have completely byzantine storage (which is probably not a bad way to think of physical disks and firmwares and kernel caches) and VSR will remain correct, despite requiring only the same resources as an otherwise non-byzantine protocol.
- 5y ago