Replication cannot fix corruption: Why distributed systems drop acked data
Three copies of your data offer an illusion of safety if the underlying system cannot detect, diagnose, and repair silent errors.
- Published
- Reading
- 9 min
- Tags
- #storage#durability#object-storage
In a distributed system, three replicas feel safe for durability and redundancy, but a replica is not a recovery plan. A copy of the data only protects you if the system can detect the bad data, diagnose the fault correctly, and repair it using a healthy copy. Three research papers demonstrate each of those steps failing across eight widely used systems, including Kafka, Cassandra, and ZooKeeper.
On AWS, an EC2 instance’s local disk does not fail solely by crashing. It can return errors for specific blocks or silently return incorrect bytes, which the underlying replicated data store simply treats as valid. Two field studies cited in [1] measured this frequency. One tracked over a million drives for 32 months, finding that 8.5% of high-capacity disks and 1.9% of enterprise-class disks developed at least one latent sector error, a block that becomes permanently unreadable. Another study of 1.53 million drives over 41 months found more than 400,000 blocks with checksum mismatches. The cloud does not eliminate this reality. As the paper warns, as deployments migrate to the cloud, reliable storage hardware, firmware, and software can no longer be assumed; data storage systems require rigorous end-to-end integrity checks.
The first paper [1] states this bluntly:
Despite the presence of checksums, redundancy, and other resiliency methods prevalent in distributed storage, a single untimely file-system fault can lead to data loss, corruption, unavailability, and, in some cases, the spread of corruption to other healthy replicas.
The researchers built a fault-injecting file system that sits between the database and the disk. They used it to introduce a single fault on one node across Redis, ZooKeeper, Cassandra, Kafka, RethinkDB, MongoDB, LogCabin, and CockroachDB. Each system ran as a three-node cluster with checksums and synchronous writes enabled. They deliberately injected only one fault to give each system the maximum opportunity to recover.
1. Detection fails: The node doesn’t know its data is bad.
Consider a single node. Occasionally, the disk’s firmware writes data to the incorrect location or silently drops a write. The corrupted bytes remain on disk, and common Linux file systems like ext4 and XFS pass them directly up to the application. Because neither maintains data checksums, applications consume corrupted data without generating any errors [1].
Beneath this layer sits fsync, the system call that instructs the operating system to persist data to disk and reports success or failure. Before fsync completes, writes reside in the page cache the operating system’s in-memory copy of a file as dirty pages waiting to be flushed. When fsync fails, ext4, XFS, and Btrfs incorrectly mark those dirty pages as clean. Consequently, calling fsync again writes nothing but reports success [3]. The lost write appears successful until the operating system evicts the page from memory. Applications continue reading the new contents from cache, but once the page is read back from disk, they receive the old contents. Redis, notably, does not verify whether fsync failed at all. Furthermore:
PostgreSQL had been using fsync incorrectly for 20 years.
The mitigation adopted by PostgreSQL and MySQL was crashing on an fsync failure and replaying the log on restart is insufficient. The paper demonstrates that this approach fails because the restarted application may read from a page cache that no longer reflects the true disk state. Ultimately, a node unaware of its own data corruption will never request a healthy copy from a replica.
2. Diagnosis fails: A noticed fault is misread, or the node crashes.
Databases append changes to a log. Following a power loss or crash, they execute crash recovery routines to clean up half-finished writes at the tail of the log. A torn write one truncated by a crash fails its checksum, exactly as a subsequently corrupted record would. Because the two scenarios look identical, most log-based systems treat every checksum mismatch as a crash [2]. The first paper observed this consequence across every tested system [1]:
On detecting a checksum mismatch due to corruption, all systems invariably run the crash recovery code (even if the corruption was not actually due to crash but rather due to a real corruption in the storage stack), ultimately leading to undesirable effects such as data loss.
Kafka’s recovery code misinterpreted corruption as evidence of an earlier crash. It truncated the log at the corrupted entry, discarding all subsequent valid data instead of repairing the single faulty record.
Misdiagnosing the fault is only one failure mode and failing to handle it entirely is another. Crashing was the most common system response to a fault. However, restarting provides no relief when the underlying fault remains on the disk. The node enters a crash loop on every restart until an operator manually intervenes [1]. ZooKeeper’s developers admitted to the researchers that crashing on corruption “was not a conscious design decision.”
The file system can also take the entire node offline. If a write to the file system’s internal journal fails during fsync, ext4 remounts the disk as read-only (halting all new writes), while XFS shuts down the file system entirely [3].
A small corrupted region can destroy exponentially more data than it occupies. Table 1 details the blast radius of a single localized fault in each system [1]:
| System | Where the fault hit | Fault | Data lost |
|---|---|---|---|
| Redis | Append-only log file (metadata section) | Any fault | Entire dataset unreadable |
| Redis | Append-only log file (user-data section) | Read or write error | Entire dataset unreadable |
| Cassandra | Table data file (first block) | Corruption | First entry lost |
| Cassandra | Table index file | Corruption | Whole SSTable unreadable |
| Cassandra | Schema metadata files (compression info, filter, statistics) | Corruption or read error | Whole table unreadable |
| Kafka | Log file (header section) | Corruption | Entire log lost |
| Kafka | Log file (any other block) | Corruption or read error | Log from that entry onward lost |
| Kafka | Replication checkpoint file | Corruption or read error | All data lost |
| Kafka | Replication checkpoint temp file | Write error | All data unreadable |
| RethinkDB | Database file (transaction head, or metablock) | Corruption | Whole transaction lost |
3. Repair fails: The replicas don’t save you, and can actively make things worse.
Almost all systems in many cases do not use redundancy as a source of recovery and miss opportunities of using other intact replicas for recovering.
While healthy copies exist on peer nodes, the software frequently fails to request them. Worse, replication protocols often propagate the damage [1]. When Cassandra replicas disagree during a read, the read repair mechanism selects one value and synchronizes it across all nodes. In the event of a timestamp tie, it selects the value that sorts later lexically. Consequently, a corrupted value can win the tiebreaker and silently overwrite healthy copies. When a Redis follower resynchronizes, it blindly copies the leader’s state. A corrupted append-only file (Redis’s transaction log) on the leader will seamlessly replicate to all followers. Kafka maintains a list of in-sync replicas (ISR) eligible for leadership. A node that incorrectly truncated its log remains in the ISR, meaning it can subsequently be elected leader and silently erase acknowledged data.
Consensus systems like ZooKeeper and LogCabin exhibit similar failure modes [2]. If they maintain a log across a five-node cluster, they consider a change committed once a quorum (three nodes) acknowledges it. Now suppose one node encounters a corrupted entry and truncates its log. If two other nodes happen to be lagging, and the two fully updated nodes temporarily go offline, the node that truncated its data can form a quorum with the two lagging nodes. Together, they will overwrite the state of the healthy nodes when they return. Committed, acknowledged data is permanently lost.
The common operational practice of wiping the faulty node and letting it resynchronize fails via the exact same mechanism. When confronted with corrupted log entries where a pristine copy existed elsewhere in the cluster, systems recovered correctly in only 46 out of 2,401 tested scenarios. As the paper concludes: “Recovering from storage faults in distributed systems is surprisingly hard.”
If your architecture relies on local disks
This is not an absolute argument against local disks. If your workload demands single-digit-millisecond latency, local NVMe remains a valid architectural choice. However, the system design must assume the local disk can fail at any moment, and it must execute detection, diagnosis, and repair flawlessly. The research outlines the necessary steps:
- Test meticulously: Inject faults block by block on real disks to trigger the file system’s actual error-handling paths, rather than mocking error codes in unit tests [3].
- Diagnose separately: Handle data corruption with distinct logic from crash recovery, isolating and repairing only the corrupted or unreadable sectors [1].
- Recover from disk: Rebuild application state explicitly from persistent storage, never trusting the page cache during recovery [3].
Otherwise, as the first paper concludes, “the general expectation that redundancy can help availability of functionality and data is not a reality” [1].
Or let a cloud object store handle the undifferentiated heavy lifting
Object stores like Amazon S3, Google Cloud Storage, Azure Blob Storage, and Cloudflare R2 are purpose-built to manage these three phases as a managed service. For detection and repair, S3 continuously verifies data using checksums and rapidly repairs degraded redundancy [4]. Google Cloud Storage periodically revalidates checksums and corrects errors using redundant copies [5]. Azure Storage similarly validates stored data via CRCs and repairs corruption autonomously [6]. Writes are acknowledged only after the data is stored redundantly across multiple facilities (as in Cloud Storage [5]) or safely persisted to disks designed to rebuild objects during hardware failures (as in R2 [7]). Diagnosis remains trivial because there are no torn writes to mistake for corruption: systems like S3 provide read-after-write consistency and never commit partial objects [8].
It is undifferentiated heavy lifting that receives little credit, yet it perfectly encapsulates the detection, diagnosis, and repair cycles that the aforementioned research proves most distributed systems get wrong.
Note: Papers [1], [2], and [3] tested versions from 2016 to 2019. Several bugs have been patched since, including a ZooKeeper bug that halted all writes [1] and PostgreSQL’s fsync retry logic [3]. However, firmware and hardware disks continue to fail.
References
- Ganesan, Alagappan, Arpaci-Dusseau, Arpaci-Dusseau. Redundancy Does Not Imply Fault Tolerance: Analysis of Distributed Storage Reactions to Single Errors and Corruptions. FAST ’17.
- Alagappan, Ganesan, Lee, Albarghouthi, Chidambaram, Arpaci-Dusseau, Arpaci-Dusseau. Protocol-Aware Recovery for Consensus-Based Storage. FAST ’18.
- Rebello, Patel, Alagappan, Arpaci-Dusseau, Arpaci-Dusseau. Can Applications Recover from fsync Failures? USENIX ATC ’20.
- Amazon Web Services. Data protection in Amazon S3. Amazon S3 User Guide.
- Google Cloud. Data availability and durability. Cloud Storage documentation.
- Microsoft. Data redundancy - Azure Storage. Microsoft Learn.
- Cloudflare. Durability. Cloudflare R2 documentation.
- Amazon Web Services. PutObject. Amazon S3 API Reference.
Arun Lakshman Ravichandran