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
- File Split: A 1GB file is split into blocks (e.g., 8 blocks of 128MB each).
- Replication: Each block is replicated 3x (default) and stored on different DataNodes.
- Rack Awareness: Replicas are placed on different racks to avoid single-point failures.
- Metadata Update: NameNode updates its namespace (file → block locations).
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.
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:
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:
| 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).
- Split data into 128MB blocks.
- Replicate 3x across DataNodes.
- 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:
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.
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).
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:
- "Explain how HDFS ensures fault tolerance." → Replication + rack awareness + heartbeat.
- "What is the default block size in HDFS? How does it affect small files?" → 128MB/256MB; inefficient for small files.
- "Compare HDFS with a traditional file system." → Use the comparison table above.
- "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…