Distributed file systems and erasure coding
How huge files are stored across thousands of disks: the GFS and HDFS design with a metadata master and chunk servers, chunk sizes and replication, write pipelines, metadata scaling, object storage internals, and erasure coding versus replication for durability and cost.
Reading is half of it. See this used in a real interview: walk through Design Dropbox →
Storing petabytes means spreading data across thousands of machines whose disks fail every day. Google File System (GFS) and its open-source descendant HDFS established the core design; object stores such as S3 use related ideas with different trade-offs. Questions like "design Dropbox's storage layer" or "how does S3 achieve eleven nines of durability?" draw on this design, and on erasure coding, which provides durability far more cheaply than full copies.
The GFS and HDFS design
- Files are split into large chunks (64 or 128 MB).
- Chunk servers (data nodes) store chunks on local disks.
- A master (name node) holds metadata: the namespace (directories and files), which chunks make up each file, and where each chunk's replicas are.
- Clients ask the master for chunk locations, then read and write data directly with chunk servers. The master never handles file data, so it does not become a bandwidth bottleneck.
Why large chunks
- Less metadata: the master can keep all of it in memory.
- Fewer client-master interactions for large sequential reads and writes.
- Efficient for big files (logs, crawls, datasets, video), the workloads these systems were built for.
The downside is poor fit for millions of tiny files, each of which costs metadata. Pack small files into larger containers when needed.
Replication and placement
Each chunk is replicated (typically 3 times), placed across racks (and zones) so a rack power or switch failure does not lose all copies. The master monitors chunk servers through heartbeats and re-replicates chunks whose copies were lost, prioritising chunks with only one remaining copy.
Writes
Data is pushed along a pipeline of replicas (client to replica 1 to replica 2 to replica 3), so each machine's network bandwidth is used once. A primary replica orders concurrent mutations. GFS favoured appends (many producers appending records to a file) and relaxed consistency for concurrent writes; HDFS simplified to a single writer per file. Both suit write-once, read-many workloads.
Scaling the metadata
A single master in memory eventually limits the number of files and request rate. Approaches:
- Federation: several masters, each owning part of the namespace.
- Metadata in a distributed database: Google's Colossus stores metadata in Bigtable; many modern systems use a sharded, consensus-replicated metadata store.
- High availability: a standby master replaying the metadata log, with consensus-based failover. See how Raft works.
Object storage
Object stores (S3, GCS, Azure Blob) expose a flat key-to-object API over HTTP instead of a file system:
- No in-place modification; objects are written whole (or in multipart uploads) and replaced.
- Metadata is a massively sharded key index; data is spread across storage nodes, often erasure coded.
- Strong read-after-write consistency is now standard on major clouds.
- Practically unlimited scale; the default for media, backups and data lakes.
See object storage and files and data lakes and lakehouses.
Erasure coding
Three-way replication tolerates two lost copies at a 3x storage cost. Erasure coding splits data into k data fragments and computes m parity fragments; any k of the k + m fragments rebuild the data.
| Scheme | Storage overhead | Failures tolerated |
|---|---|---|
| 3x replication | 3.0x | 2 |
| Reed-Solomon 6 + 3 | 1.5x | 3 |
| Reed-Solomon 10 + 4 | 1.4x | 4 |
Spreading fragments across racks or zones gives high durability at half the cost of replication. Trade-offs:
- Reconstruction cost: rebuilding a lost fragment reads k others, using a lot of network and disk.
- Latency: reads may need several fragments; small objects are inefficient to code.
- CPU for encoding and decoding.
So systems often keep hot data replicated and convert colder data to erasure coding. See cost-aware system design.
Durability in numbers
Durability depends on how fast failures are detected and repaired relative to how often disks fail. More fragments across more failure domains, faster repair and checksums to detect silent corruption (scrubbing data regularly) all push durability toward the "eleven nines" that object stores advertise. Durability is not backup: a deletion or bug replicates instantly. See backups and disaster recovery.
In the interview
For Design Dropbox: metadata (files, versions, chunk lists) in a sharded database; file content as content-addressed chunks in an object store, erasure coded across zones; clients upload chunks directly. For Design a Web Crawler, raw pages appended into large files in a distributed file system or object store.
Checklist
- Metadata separated from data; clients read and write data directly.
- Large chunks or objects; small files packed.
- Rack- and zone-aware placement with automatic re-replication.
- Pipelined writes; append-friendly, write-once workloads.
- Scalable, highly available metadata.
- Erasure coding for colder data, replication for hot data.
- Checksums and scrubbing; backups separate from durability.