5 ms·
> 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
by _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.
- _benedict 5y ago> 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. This isn’t a showstopper event, it’s a process stopper event. I also outlined non-PAR approaches to recovering just fine without replacing the node. > Theoretically, but in practice I also have extensive practical confirmation that this approach works well. > This aspect of recovery is often simply not tested. I agree that test suites do not cover this well enough, and I commend you for working this into Tiger Beetle. This is something I hope to expand Cassandra's new deterministic simulation framework to incorporate in future as well, but Cassandra does benefit from a great deal of real world exposure to this kind of fault. > 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. I’m not sure I agree. The paper discounts reconfiguration because there could only be f+1 live processes and one could have corruption in a relevant sector, but under normal models this is simply f live processes. It seems that to accommodate this scenario we must duplicate the record identifier, and store both separately from the record itself. It's not clear to me this scenario warrants the additional complexity, storage and bandwidth, as I’m not sure guarding against this is enough to reduce your replication factor. But it's worth considering, and this particular recovery enhancement is quite simple (we just need to duplicate the ballot in Paxos, so we can arbitrate between a split decision of other replicas that retain an intact record), so thanks for highlighting it more clearly for me. > this kind of fault is addressed by Viewstamped Replication's 2012 revision Could you point me to the place in the paper? I cannot see how VR solves this problem without introducing a risk of inconsistency, without necessitating an additional round-trip before responding to a client. Specifically, I think there are only three ways to address this particular problem, and I don’t see them discussed in VR: 1. The coordinator may record the responses of each replica before acknowledging to the client, so that the loss of any a write to any one disk may be recoverable. 2. The coordinator may require k>f+1 responses before answering a client, so that we may tolerate the loss of k-(f+1) disks losing a write 3. The coordinator may require an additional round-trip to record the consensus decision before answering a client I don’t see any of these approaches discussed in VR. I see discussion of recovering the log from other replicas, but this cannot be done safely if the replica is not itself aware that it has lost data.