Big Data and AnalyticsUnit 311 min read

Hadoop & HDFS: Architecture, HDFS Design, Data Storage & Fault Tolerance

Unit 3 of Big Data and Analytics explores Hadoop’s ecosystem, HDFS architecture (NameNode, DataNode, replication), data storage mechanics, and fault tolerance mechanisms like replication and rack awareness, with real-world applications in Nepal’s eSewa and Ncell systems.

TAKEAWAYS:

  • Hadoop is an open-source framework for distributed storage and processing of big data across clusters.
  • HDFS (Hadoop Distributed File System) splits files into blocks (default 128MB/256MB) and stores them across DataNodes for scalability.
  • NameNode manages metadata (file locations, permissions), while DataNode stores actual data blocks and reports to NameNode.
  • Fault tolerance in HDFS relies on replication (default 3x) and rack awareness to prevent data loss.
  • Data ingestion in HDFS uses write-once-read-many (WORM) model, optimized for batch processing (not low-latency updates).
  • Hadoop’s ecosystem includes YARN (resource manager), MapReduce (processing), and auxiliary tools like HBase, Hive, and Pig.

1. Introduction to Hadoop

Hadoop is an open-source framework designed to store and process large-scale datasets across clusters of commodity hardware. It was developed by Doug Cutting and Mike Cafarella at Yahoo! (2006) and later donated to the Apache Software Foundation. Hadoop’s core components include:

  • HDFS (Hadoop Distributed File System): For storage.
  • YARN (Yet Another Resource Negotiator): For resource management.
  • MapReduce: For parallel processing.

Why Hadoop?

  • Scalability: Handles petabytes of data by distributing storage and computation.
  • Fault Tolerance: Automatically recovers from node failures.
  • Cost-Effective: Runs on cheap commodity hardware (unlike expensive supercomputers).
  • Flexibility: Processes structured, semi-structured, and unstructured data (text, images, logs, etc.).

Real-World Example: eSewa’s Transaction Logs

eSewa processes millions of transactions daily. Storing and analyzing these logs in HDFS allows:

  • Scalable storage of transaction records (JSON/XML format).
  • Batch processing (e.g., fraud detection) using MapReduce.
  • Fault tolerance to prevent data loss during peak hours.

2. HDFS Architecture

HDFS is designed for high-throughput, batch-oriented workloads (not real-time analytics). Its architecture consists of:

A. Core Components

Component Role Example in Nepalese Context
NameNode Manages metadata (file locations, permissions, namespace). Like a library catalog tracking all books.
DataNode Stores actual data blocks (default 128MB/256MB). Like shelves storing books in a library.
Secondary NameNode Periodically checkpoints NameNode’s metadata (not a backup!). Like a backup librarian updating records.
Client Interacts with HDFS (reads/writes data via NameNode). Like a user searching for a book.

B. How HDFS Works: Data Storage Process

  1. File Split: A 1GB file is split into blocks (e.g., 8 blocks of 128MB each).
  2. Replication: Each block is replicated 3x (default) and stored on different DataNodes.
  3. Rack Awareness: Replicas are placed on different racks to avoid single-point failures.
  4. Metadata Update: NameNode updates its namespace (file → block locations).
DataNode 1 (Rack 1): Block 1 (128MB)DataNode 2 (Rack 2): Block 2 (128MB)DataNode 3 (Rack 3): Block 3 (Replica)NameNode Assigns Block LocationsClient Requests File Write
HDFS block storage hierarchy with rack awareness (1GB file split into 8 blocks)

C. Real Picture: HDFS Cluster Setup



3. HDFS Data Storage Mechanics

A. Block Storage

  • Files are split into fixed-size blocks (default: 128MB or 256MB).
  • Small files problem: HDFS is inefficient for millions of tiny files (overhead on NameNode).
    • Solution: Use Hadoop Archives (HAR) or combine small files into larger ones.

B. Replication Factor

  • Default replication = 3 (configurable via dfs.replication).
  • Trade-off:
    • Higher replication → better fault tolerance but more storage overhead.
    • Lower replication → saves space but higher risk of data loss.

C. Rack Awareness

  • HDFS places replicas on different racks to prevent data loss from rack failures.
  • Example: If Rack 1 fails, replicas in Rack 2/3 ensure data availability.
00.511.52Rack 11Rack 22Rack 32
Replica distribution across racks (1 primary, 2 replicas)

D. Write-Once-Read-Many (WORM) Model

  • Appends only: Once data is written, it cannot be modified (only new data can be appended).
  • Optimized for batch processing (not real-time updates like databases).
  • Use case: Log files, sensor data, historical records.

4. Fault Tolerance in HDFS

HDFS ensures data availability even if nodes fail:

T0Rack 1 failuredetected (Heartbeat tiT1NameNode triggersreplica creation on RaT2Data recoverycomplete (WORM model e
Fault tolerance recovery timeline in HDFS

A. Replication

  • If a DataNode fails, NameNode detects it and re-replicates missing blocks from other replicas.
  • Example: Ncell’s call detail records (CDRs) are stored in HDFS with replication=3 to prevent loss during network outages.

B. Heartbeat Mechanism

  • DataNodes send heartbeats to NameNode every 3 seconds.
  • If a DataNode stops responding, NameNode marks it dead and re-replicates its blocks.

C. Checkpointing (Secondary NameNode)

  • Secondary NameNode periodically merges NameNode’s edits log into FSImage (metadata snapshot).
  • Not a backup! It helps recover NameNode in case of crashes.

D. Real-World Example: Daraz’s Order Processing

Daraz processes thousands of orders per second. HDFS ensures:

  • Order logs are stored with replication=3 across multiple DataNodes.
  • If a DataNode fails (e.g., during a power outage), orders are not lost due to replication.

5. HDFS vs. Traditional File Systems

Feature HDFS Traditional FS (e.g., ext4, NTFS)
Storage Distributed across clusters Single machine
Scalability Petabytes Limited by disk size
Fault Tolerance Automatic replication Manual backups
Access Pattern Batch processing (WORM) Frequent reads/writes
Metadata Handling NameNode (bottleneck for small files) Local filesystem metadata
Use Case Big Data analytics General-purpose storage

6. HDFS Commands (Practical Example)

Students should know basic HDFS commands for exams and labs:

→ HDFS Block 1 (128MB)→ HDFS Block 2 (128MB)Local Filehdfs dfs -put command
File upload workflow using HDFS CLI
Command Description
hdfs dfs -put file.txt /input Uploads file.txt to HDFS /input directory.
hdfs dfs -cat /input/file.txt Displays contents of file.txt.
hdfs dfs -ls /input Lists files in /input directory.
hdfs dfs -mkdir /output Creates /output directory.
hdfs dfs -rm /input/file.txt Deletes file.txt (use -r for recursive delete).
hdfs dfsadmin -report Shows cluster health (live/dead DataNodes, storage stats).

Worked Example: Storing NEPSE Stock Data

Suppose NEPSE wants to analyze 10 years of stock market data (1TB).

  1. Split data into 128MB blocks.
  2. Replicate 3x across DataNodes.
  3. Store in HDFS for batch processing (e.g., trend analysis using MapReduce).
# Upload stock data to HDFS
hdfs dfs -put nepse_data.csv /stock_data/

# Check replication factor
hdfs dfs -ls /stock_data/nepse_data.csv
# Output: -rw-r--r-- 3 user supergroup 1073741824 2023-10-01 10:00 /stock_data/nepse_data.csv

7. Limitations of HDFS

While HDFS is powerful, it has key limitations:

  • Not for real-time processing: High latency for frequent updates (use HBase instead).
  • Small files problem: NameNode struggles with millions of tiny files (solution: HAR files or combining files).
  • No random writes: WORM model makes it inefficient for databases.
  • Single NameNode bottleneck: Metadata operations can slow down with millions of files.

8. In the Real World

HDFS powers critical systems in Nepal and globally:

  1. eSewa (Nepal)

    • Use Case: Storing transaction logs (JSON format) for fraud detection.
    • How HDFS Helps:
      • Scalable storage for millions of daily transactions.
      • Batch processing to analyze spending patterns (MapReduce).
      • Fault tolerance ensures logs are never lost during peak hours.
  2. Ncell’s Call Detail Records (CDRs)

    • Use Case: Analyzing call patterns for network optimization.
    • How HDFS Helps:
      • Stores terabytes of CDR data with replication=3.
      • Enables batch analytics (e.g., identifying busy hours).
  3. Google’s Web Crawling (Global)

    • Use Case: Storing petabytes of web pages.
    • How HDFS Helps:
      • Distributes web crawl data across thousands of DataNodes.
      • Uses MapReduce to index pages for Google Search.

9. Exam Tip

For TU/PU exams, focus on: ✅ HDFS architecture (NameNode vs. DataNode roles). ✅ Block storage (default size, replication factor). ✅ Fault tolerance (replication, rack awareness, heartbeat). ✅ Limitations (WORM model, small files problem). ✅ Real-world applications (eSewa, Ncell, NEPSE). ✅ Commands (hdfs dfs -put, -ls, -report).

Common Exam Questions:

  1. "Explain how HDFS ensures fault tolerance." → Replication + rack awareness + heartbeat.
  2. "What is the default block size in HDFS? How does it affect small files?" → 128MB/256MB; inefficient for small files.
  3. "Compare HDFS with a traditional file system." → Use the comparison table above.
  4. "How would you store NEPSE’s stock data in HDFS?" → Split into blocks, replicate 3x, process with MapReduce.

Avoid: ❌ Memorizing exact port numbers (e.g., NameNode runs on 8020 by default, but exams may not ask). ❌ Confusing Secondary NameNode with a backup (it’s a checkpointing mechanism). ❌ Forgetting rack awareness in fault tolerance explanations.


Final Note: HDFS is the backbone of Hadoop’s storage. Master its architecture, replication, and limitations to ace this unit! 🚀

Based on the TU BITM syllabus for Big Data and Analytics (IT278), unit 3.

Discussion

Loading…