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:
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
- Client writes a file → Splits into blocks → Sends to NameNode.
- NameNode selects DataNodes (based on rack awareness and load balancing).
- DataNodes acknowledge receipt → NameNode updates metadata.
- Replication begins: NameNode sends block copies to other DataNodes.
- Client reads file → NameNode directs to nearest DataNode (data locality).
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).
How Replication Works
- Default Replication Factor = 3.
- If a DataNode fails:
- NameNode detects missing heartbeats.
- Triggers re-replication of blocks from other DataNodes.
- 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%.
- Increase block size to 256MB (
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
- Define HDFS clearly:
- "HDFS is a distributed file system that splits files into blocks and stores them across DataNodes with replication for fault tolerance."
- Explain NameNode vs. DataNode:
- NameNode = metadata manager (fsimage, edits log).
- DataNode = data storage (blocks, heartbeats).
- Replication is key:
- Default = 3x, but explain rack awareness for fault tolerance.
- Limitations matter:
- Small files problem, not for real-time queries, NameNode bottleneck.
- Real-world tie-ins:
- eSewa (transaction logs), Daraz (order data), NTC (CDR analysis).
- 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…