Big Data and AnalyticsUnit 311 min read

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

Unit 3 of Big Data and Analytics explores Hadoop’s distributed ecosystem and HDFS (Hadoop Distributed File System), covering its architecture, components (NameNode, DataNode), replication, data flow, and real-world applications like eSewa’s transaction logs and Ncell’s call detail records.

TAKEAWAYS:

  • HDFS splits files into blocks (default 128MB) and distributes them across DataNodes for parallel processing.
  • Hadoop’s NameNode manages metadata (file locations, permissions) while DataNodes store actual data blocks.
  • Replication (default 3x) ensures fault tolerance—if one DataNode fails, data remains available.
  • HDFS is not a database: it’s optimized for batch processing (MapReduce) and not low-latency queries.
  • Small files problem: HDFS struggles with millions of tiny files; solutions include HAR (Hadoop Archives) or SequenceFiles.
  • Rack awareness improves data locality by storing replicas on different racks to survive rack failures.

1. What is Hadoop?

Hadoop is an open-source framework for distributed storage and processing of large datasets across clusters of commodity hardware. It follows the "write-once, read-many" paradigm and is designed for fault tolerance and scalability.

Why Hadoop?

  • Cost-effective: Uses cheap hardware instead of expensive supercomputers.
  • Scalability: Can grow from a single server to thousands of nodes.
  • Fault tolerance: Automatically recovers from hardware failures.
  • Flexibility: Supports structured, semi-structured, and unstructured data.

Hadoop Ecosystem

Hadoop is not just HDFS—it includes:

Block Storage (128MB/256MB)Replication (3x default)HDFS: StorageDivide → Map → Shuffle → ReduceMapReduce: Batch ProcessingResource NegotiationJob SchedulingYARN: Resource ManagementColumn-Family StoreHBase: NoSQL DatabaseQL (Hive Query Language)Hive: SQL-like QueryingPig Latin ScriptsPig: ETL ScriptingHadoop Ecosystem
Hierarchical breakdown of Hadoop components with key features

2. HDFS: Hadoop Distributed File System

HDFS is the storage layer of Hadoop, designed for high-throughput access to large files. It divides files into blocks (default 128MB or 256MB) and distributes them across a cluster.

Key Features of HDFS

Feature Description
Block Storage Files split into fixed-size blocks (default 128MB).
Replication Each block stored on 3 DataNodes by default (configurable).
Master-Slave NameNode (master) manages metadata; DataNodes (slaves) store data.
Write-Once Files can be written once and read many times (immutable).
Data Locality Processing happens where data is stored to minimize network I/O.

How HDFS Works: Data Flow

  1. Client writes a file → Splits into blocks → Sends to NameNode.
  2. NameNode selects DataNodes (based on rack awareness and load balancing).
  3. DataNodes acknowledge receipt → NameNode updates metadata.
  4. Replication begins: NameNode sends block copies to other DataNodes.
  5. Client reads file → NameNode directs to nearest DataNode (data locality).
1. Client RequestUser submits filewrite operation2. NameNode AllocationNameNode assignsblocks to DataNodes3. DataNode StorageData written in128MB/256MB blocks4. ReplicationBlocks replicatedacross 3+ nodes5. Client ConfirmationWrite completionacknowledged
Step-by-step HDFS write operation workflow

3. HDFS Components

A. NameNode (Master)

  • Role: Manages metadata (file locations, permissions, block mappings).
  • Storage: Keeps metadata in RAM (for fast access) and fsimage/edits log (on disk).
  • Single Point of Failure (SPOF): If NameNode crashes, the entire cluster goes down.
    • Solution: Secondary NameNode (not a backup!) and High Availability (HA) modes.

B. DataNode (Slave)

  • Role: Stores actual data blocks and communicates with NameNode.
  • Functions:
    • Receives blocks from clients.
    • Sends heartbeats to NameNode (every 3 seconds by default).
    • Performs block reports (lists stored blocks).
    • Handles read/write requests from clients.

C. Secondary NameNode

  • Misconception: It is not a backup for NameNode!
  • Actual Role:
    • Periodically merges fsimage and edits log to prevent metadata bloat.
    • Helps in failover recovery (in HA setups).

4. HDFS Replication & Fault Tolerance

HDFS ensures data availability by replicating each block 3 times (configurable via dfs.replication in hdfs-site.xml).

0128256384512Default Replication3Minimum Replication1Maximum Replication512
HDFS replication configuration limits (default=3)

How Replication Works

  1. Default Replication Factor = 3.
  2. If a DataNode fails:
    • NameNode detects missing heartbeats.
    • Triggers re-replication of blocks from other DataNodes.
  3. Rack Awareness:
    • Replicas stored on different racks to survive rack failures.
    • Example: If Rack 1 fails, data remains available on Racks 2 and 3.

Replication Strategies

Strategy Description
Default (3x) Balances availability and storage overhead.
Under-Replication Fewer replicas (e.g., 2x) for cost savings (but higher risk).
Over-Replication More replicas (e.g., 4x) for critical data (higher storage cost).

Worked Example: Ncell’s Call Detail Records (CDR) Storage

  • Problem: Ncell processes millions of call logs daily (unstructured data).
  • Solution:
    • Store CDRs in HDFS with 3x replication.
    • Use MapReduce to analyze call patterns (e.g., peak hours, fraud detection).
    • Rack awareness ensures data survives even if a server rack fails.

5. HDFS Limitations & Challenges

Limitation Cause Solution
Small Files Problem HDFS overhead per file (~150KB metadata). Use HAR (Hadoop Archives) or SequenceFiles.
Not for Low-Latency Optimized for batch processing. Use HBase for real-time queries.
High Latency for Random Writes Write-once model. Use append-only or HBase.
Single NameNode Bottleneck Metadata operations slow down. Enable HDFS HA (High Availability).

6. HDFS vs. Traditional File Systems

Feature HDFS Traditional FS (e.g., ext4, NTFS)
File Size Handles petabytes of data. Limited by filesystem size.
Scalability Scales to thousands of nodes. Limited by single machine.
Fault Tolerance Automatic replication. Manual backups required.
Access Pattern Batch processing (MapReduce). Random reads/writes.
Metadata Handling NameNode manages metadata. Local filesystem handles metadata.

7. Real-World Applications of HDFS

A. eSewa: Transaction Logs

  • Problem: eSewa processes millions of transactions/day (mobile payments, bill payments).
  • HDFS Role:
    • Stores transaction logs in HDFS with 3x replication.
    • Uses MapReduce to analyze spending patterns (e.g., which districts use eSewa most?).
    • Rack awareness ensures logs survive server failures.

B. Daraz: Order Processing

  • Problem: Daraz handles millions of orders/day (e-commerce).
  • HDFS Role:
    • Stores order data (JSON/XML) in HDFS.
    • Uses Hive to run SQL-like queries (e.g., "Top 10 selling products in Kathmandu").
    • Small files issue: Daraz uses Parquet/ORC formats to reduce metadata overhead.

C. NTC: Network Traffic Analysis

  • Problem: Nepal Telecom (NTC) needs to analyze terabytes of call logs.
  • HDFS Role:
    • Stores CDR (Call Detail Records) in HDFS.
    • Uses Spark for real-time fraud detection (e.g., unusual call patterns).
    • Replication ensures logs survive hardware failures.

8. Configuring HDFS

Key configuration files in $HADOOP_HOME/etc/hadoop/:

  • core-site.xml:
    <property>
      <name>fs.defaultFS</name>
      <value>hdfs://namenode:8020</value>
    </property>
    
  • hdfs-site.xml:
    <property>
      <name>dfs.replication</name>
      <value>3</value> <!-- Default replication factor -->
    </property>
    <property>
      <name>dfs.blocksize</name>
      <value>268435456</value> <!-- 256MB block size -->
    </property>
    

Worked Example: Changing Block Size for NEPSE Stock Data

  • Problem: NEPSE stores daily stock data (small files, ~1MB each).
  • Solution:
    • Increase block size to 256MB (dfs.blocksize=268435456).
    • Use SequenceFile to combine small files into larger ones.
    • Result: Reduces NameNode metadata overhead by 90%.

9. HDFS Commands (Cheat Sheet)

Command Description
hdfs dfs -ls / List files in HDFS root.
hdfs dfs -put localfile /hdfs/ Upload a file to HDFS.
hdfs dfs -get /hdfs/file local Download a file from HDFS.
hdfs dfs -cat /hdfs/file Display file contents.
hdfs dfs -du -h / Check disk usage in HDFS.
hdfs dfsadmin -report Check cluster health (live/dead nodes).

10. Exam Tip: How to Score Full Marks

  1. Define HDFS clearly:
    • "HDFS is a distributed file system that splits files into blocks and stores them across DataNodes with replication for fault tolerance."
  2. Explain NameNode vs. DataNode:
    • NameNode = metadata manager (fsimage, edits log).
    • DataNode = data storage (blocks, heartbeats).
  3. Replication is key:
    • Default = 3x, but explain rack awareness for fault tolerance.
  4. Limitations matter:
    • Small files problem, not for real-time queries, NameNode bottleneck.
  5. Real-world tie-ins:
    • eSewa (transaction logs), Daraz (order data), NTC (CDR analysis).
  6. Diagrams save marks:
    • Draw HDFS architecture (NameNode, DataNodes, replication).
    • Show data flow (client → NameNode → DataNodes).

Common Mistakes to Avoid:

  • ❌ Saying "Secondary NameNode is a backup" (it’s not!).
  • ❌ Confusing HDFS with a database (it’s for batch processing, not OLTP).
  • ❌ Ignoring rack awareness in replication discussions.

11. Quick Revision Summary

mindmap
  root((HDFS))
    Architecture
      NameNode["Master (Metadata)"]
      DataNode["Slave (Data Storage)"]
      SecondaryNameNode["Metadata Merge"]
    Features
      BlockStorage["128MB/256MB blocks"]
      Replication["3x by default"]
      DataLocality["Process near data"]
    Limitations
      SmallFiles["Metadata overhead"]
      NoRandomWrites["Write-once model"]
      NameNodeBottleneck["Single point of failure"]
    Commands
      put/get/cat["File operations"]
      dfsadmin["Cluster health"]
    RealWorld
      eSewa["Transaction logs"]
      Daraz["Order data"]
      NTC["CDR analysis"]

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

Discussion

Loading…