1. What It Is
When data is too big for one disk, distributed file systems split files into chunks across commodity machines. GFS and HDFS are the textbook models β master for metadata, chunkservers for bytes.
What:
Master-worker distributed filesystems (GFS, HDFS) designed to store and stream massive, Multi-Terabyte files across clusters of commodity hardware servers.
Primary purpose:
Providing high-throughput sequential data access for analytics engines while surviving frequent physical server disk failures.
Usually used for:
Analytical data lakes, MapReduce/Spark data sources, and LSM-tree database storage backings (HBase, Cassandra SSTables).
2. Core Mental Model
GFS/HDFS solved "store multi-TB files on cheap disks" for the MapReduce era. Modern cloud data lakes often replace HDFS with S3 + Spark/Presto (concept #21) β same sequential read pattern, no NameNode RAM ceiling. Still name GFS/HDFS when the interviewer asks about Hadoop-era architecture or block replication mechanics.
π Decoupled Master Control
The NameNode manages directory mappings *only*. Clients query the Master for block addresses, then read data bytes directly from DataNodes.
π§± Massive Block Segmentation
Files are cut into massive 64 MB or 128 MB blocks (rather than standard 4 KB OS filesystem sectors) to minimize NameNode metadata RAM footprints.
π§± Pipeline replication
During file writes, clients stream block chunks to DataNode A in a pipeline. Node A forwards bytes to B, which forwards to C, minimizing network bottlenecks.
In the room
You rarely design a new GFS β but the pattern appears in S3 (object + metadata separation) and HBase. Mention single master bottleneck, large chunk size (64MB), and append-only writes. Good for batch analytics, not low-latency OLTP.
3. Why It Matters in HLD
Distributed file systems split large files across commodity nodes β GFS and HDFS are the textbook models. Three lenses:
Needed When:
Designing massive big-data batch analytics platforms, or configuring file-level storage systems for Exabyte-scale log pipelines.
Avoids:
Master metadata storage exhaustion, data pipeline network congestions, and critical data loss from frequent hardware server disk crashes.
Optimizes For:
Data pipeline sequential read throughput, cluster storage cost boundaries, rack-aware durability, and fault recovery speeds.
4. Architecture & Data Flow
Walk GFS architecture as interview steps. Step 1 β Client: library splits file into 64 MB chunks. Step 2 β Master: metadata server maps chunk β chunkserver locations. Step 3 β Write: pipeline chain replication to chunkservers. Step 4 β Read: client contacts master for locations, reads from nearest replica. Step 5 β Fault: master re-replicates under-replicated chunks.
In the room
Note the single master metadata bottleneck β interviewers often ask how you'd scale metadata (sharded namespace, ZFS-style, or object store migration).
5. Key Characteristics
Large chunk size, single master, and append-only writes β design characteristics we compare:
- Cluster roles β who holds metadata vs data bytes:
| Node Type | Operational Mechanic | Architectural Role |
|---|---|---|
| Master / NameNode | Maintains directory tree hierarchy and maps file paths to block arrays in RAM. |
|
| Worker / DataNode | Stores raw 64 MB / 128 MB data block segments on local Linux file disks. |
|
6. Strategic Tradeoffs
Petabyte scale on cheap hardware trades master bottleneck and latency β we state both:
| Benefit | Cost |
|---|---|
| Sequential Throughput Economy (streaming massive 128 MB blocks directly from DataNodes achieves bare-metal disk throughput speeds) |
|
| Commodity Fleet Durability (built-in 3x replication across distinct hardware server racks shields against routine host failures) | Master Metadata RAM Limits (because the NameNode maps every block in RAM, small-file proliferation exhausts Master memory) |
7. Failure / Bottleneck Awareness
Master SPOF, small-file overhead, and hotspot chunks β we volunteer mitigations:
Problem: Millions of small files (10 KB each) each consume block metadata in NameNode RAM (~150 bytes per block). The master runs out of memory and the cluster stalls.
Mitigation: Pack small records into larger archive files (SequenceFile, HAR) so one logical file maps to fewer blocks.
Problem: If the single active Master NameNode crashes, metadata mappings are inaccessible, completely freezing all cluster write/read operations.
Mitigation: Run a standby NameNode with ZooKeeper-managed failover so a passive node can take over when the active master fails.
8. Common HLD Usage
MapReduce, log storage, and ML training datasets use HDFS-style systems:
| Big Data Workload | Distributed File System Use | Architectural Rationale |
|---|---|---|
| Hadoop MapReduce / Spark | HDFS (Distributed File System) | High-volume sequential data-analytics jobs stream blocks locally from neighboring DataNodes, conserving cross-rack network bandwidth. |
| BigTable / HBase | GFS / HDFS Storage Layer | Appends database WAL records and SSTable segments sequentially to immutable DFS blocks, utilizing native 3x replica durability. |
| ML Training Data Lake | HDFS or cloud object store (S3 + Spark) | Terabyte Parquet datasets land as large sequential files. Many teams now skip on-prem HDFS and read S3 directly β same append-only access pattern, managed durability. |
9. Decision Signals
Reach for distributed FS when files are huge, append-heavy, and sequential read dominates:
- You are architecting massive MapReduce batch analytics, clickstream log pools, or ML data lake houses.
- Files are extremely massive (Multi-Gigabyte to Terabyte scales) and strictly append-only (no random modifications).
- You want to build highly-durable data storage tiers using cheap commodity hardware servers rather than expensive SAN storage arrays.
11. Deep Dive (Optional)
Rack-Aware HDFS Replication Policy (Data Safety Economy)
To shield data against physical switch failures and electrical outages, GFS/HDFS implements a highly specific **Rack-Aware Placement Policy** during write pipelines:
- Replica 1: Placed on a local DataNode node in the same rack as the client (minimizes initial network latency).
- Replica 2: Placed on a DataNode located in a *different* remote physical server rack (shields against complete local rack power outages).
- Replica 3: Placed on a secondary DataNode located in the *same* remote physical server rack as Replica 2 (avoids paying high inter-rack network routing overheads for the third copy).
This optimizes the write pipeline path, balancing local rack writing speed with complete rack-failure survivability.
Review
How helpful was this walkthrough?
Click a star to rate. We actively use this feedback to refine and update our system design content.
Discussion
Share your thoughts, ask questions, or help others.