Core Concept

Distributed File Systems: GFS & HDFS

GFS and HDFS spread large files across commodity servers as fixed-size blocks. A NameNode holds metadata; DataNodes store bytes and replicate for rack-aware durability.


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.

Loading...

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 TypeOperational MechanicArchitectural Role
Master / NameNodeMaintains directory tree hierarchy and maps file paths to block arrays in RAM.
  • Single Coordinator (saves metadata edit logs to disk
  • does NOT touch file data bytes).
Worker / DataNodeStores raw 64 MB / 128 MB data block segments on local Linux file disks.
  • Data Store (streams block bytes directly to clients
  • executes 3x replication loops).

6. Strategic Tradeoffs

Petabyte scale on cheap hardware trades master bottleneck and latency β€” we state both:

BenefitCost
Sequential Throughput Economy (streaming massive 128 MB blocks directly from DataNodes achieves bare-metal disk throughput speeds)
  • No Random Modifies (files are strictly append-only
  • editing bytes inside the middle of a file is completely unsupported)
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:

Small-file NameNode pressure

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.

🐒 NameNode Single Point of Failure (SPOF)

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 WorkloadDistributed File System UseArchitectural Rationale
Hadoop MapReduce / SparkHDFS (Distributed File System)High-volume sequential data-analytics jobs stream blocks locally from neighboring DataNodes, conserving cross-rack network bandwidth.
BigTable / HBaseGFS / HDFS Storage LayerAppends database WAL records and SSTable segments sequentially to immutable DFS blocks, utilizing native 3x replica durability.
ML Training Data LakeHDFS 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:

🎯 Think Distributed File Systems When:
  • 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:

  1. Replica 1: Placed on a local DataNode node in the same rack as the client (minimizes initial network latency).
  2. Replica 2: Placed on a DataNode located in a *different* remote physical server rack (shields against complete local rack power outages).
  3. 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

Help Us Improve

How helpful was this walkthrough?

Click a star to rate. We actively use this feedback to refine and update our system design content.

Placeholder
Optional but highly appreciated!

Discussion

Share your thoughts, ask questions, or help others.

Loading comments...