3 ms·
For people interested in this topic I can recommend "Can Applications Recover from fsync Failures?" by A. Rebello et al[1]. The paper analyzes how file systems
by simonz05 4y ago
For people interested in this topic I can recommend "Can Applications Recover from fsync Failures?" by A. Rebello et al[1]. The paper analyzes how file systems and PostgreSQL, LMDB, LevelDB, SQLite, and Redis react to fsync failures. It shows that although applications use many failure-handling strategies, none are sufficient: fsync failures can cause catastrophic outcomes such as data loss and corruption.
- [1] https://www.usenix.org/conference/atc20/presentation/rebello https://www.usenix.org/conference/atc20/presentation/rebello
- AdamProut 4y agoA nice outcome here for distributed SQL databases that store multiple copies of the data by default is that they can just failover over on an fsync failure and not try to "handle it". If fsync starts failing with EIO there is a good chance the disk is going to die soon anyways.
- simonz05 4y agoThat depends. Replicated state machines need to distinguish between a crash and corruption. Most systems don't do that. It can be disastrous to truncate the journal when encountering a checksum mismatch for instance. See "Protocol-Aware Recovery for Consensus-Based Storage" for more research on that topic. [1][2] - [1] https://www.usenix.org/system/files/conference/fast18/fast18-alagappan.pdf https://www.usenix.org/system/files/conference/fast18/fast18... - [2] https://www.youtube.com/watch?v=fDY6Wi0GcPs https://www.youtube.com/watch?v=fDY6Wi0GcPs
- AdamProut 4y agoAn fsync() failing doesn't necessarily mean there is a disk corruption. I agree all logging and recovery protocols should have different handling for a corruption vs a torn tail of the log for example, but I view that as mostly orthogonal. I'm talking about exiting the process if fsync() fails and letting the distributed databases normal failover processing do its thing. This is a normal scenario for a failover (i.e, its the same as process crashing or getting OOM killed by linux, etc).
- ryanworl 4y agoYou're probably aware of this, but for the sake of others reading: Crashing after an fsync failure isn't sufficient if you're using buffered IO either. Dirty pages in the page cache could cause your consensus implementation to e.g. allow voting for two different leaders for the same term if at some future point that machine crashes and the dirty page contained the log record for that vote. Your process would restart after the machine restarts and no longer have access to the dirty page and potentially vote again if asked. Edit: Thought I'd add the series of steps for this to happen to make it clearer: 1. Process X receives an RPC from Process Y to vote for Y as leader in term 1. 2. Process X writes log record containing vote information to page cache and calls fsync. 3. Process X receives EIO in response to fsync (if it is lucky...) and crashes. 4. Process X restarts and receives a retry of the same RPC from Process Y for term 1 and this time it responds affirmatively because it can see it already voted in term 1 for Process Y, which _should_ be safe. 5. Process X crashes because the machine hosting X experiences a temporary hardware failure. 6. Process X restarts after the hardware failure and receives an RPC from Process Z and responds affirmatively because the dirty page that was available back at T4 never actually made it to disk. Consensus is now (potentially) broken from electing two leaders for the same term.
- remram 4y agoIn distributed system terms, this is a byzantine failure which most consensus mechanisms are not equipped to handle.
- _vvhw 4y agoAnd yet many non-Byzantine consensus protocols are equipped to handle the network fault model, which could be seen as equally Byzantine under this definition. The problem is really that many formal proofs of consensus have focused only on the network fault model, and neglected the storage fault model. Both network/storage fault models require practical engineering efforts to get right. I think a better term for this is “near-Byzantine” fault tolerance. It's what non-Byzantine fault tolerance looks like when implemented correctly in the real world—the GP comment is a great example of how to approach and think about this from an engineering perspective. I dive into this in detail also here: https://www.youtube.com/watch?v=rNmZZLant9o https://www.youtube.com/watch?v=rNmZZLant9o
- _vvhw 4y agoIndeed! I was so glad to see you cite not only the Rebello paper but also Protocol-Aware Recovery for Consensus-Based Storage. When I read your first comment, I was about to reply to mention PAR, and then saw you had saved me the trouble! UW-Madison are truly the vanguard where consensus hits the disk. We implemented Protocol-Aware Recovery for TigerBeetle [1], and I did a talk recently at the Recurse Center diving into PAR, talking about the intersection of global consensus protocol and local storage engine. It's called Let's Remix Distributed Database Design! [2] and owes the big ideas to UW-Madison. [1] https://github.com/coilhq/tigerbeetle https://github.com/coilhq/tigerbeetle [2] https://www.youtube.com/watch?v=rNmZZLant9o https://www.youtube.com/watch?v=rNmZZLant9o
- AdamProut 4y agoTo be clear if fsync() on linux had well defined behavior on errors (no matter the file system used) I wouldn't suggest failing over as a reasonable solution for a distributed database. Its mentioned in the parent video, but fsyncgate[1] was a recent reminder of this. [1] https://news.ycombinator.com/item?id=20491965 https://news.ycombinator.com/item?id=20491965