Lecture Notes: Introduction to NOSQL Systems
24.1 Introduction to NOSQL Systems
24.1.1 Emergence of NOSQL Systems
Introduction
Traditional relational database systems (RDBMSs) are highly effective for:
Structured data
Transaction processing
Complex SQL queries
However, modern Internet-scale applications generate:
Massive volumes of data
Semi-structured or unstructured data
Rapidly growing datasets
This led to the emergence of:
NOSQL Systems
Meaning of NOSQL
NOSQL
Stands for:
“Not Only SQL”
NOSQL systems:
Do not necessarily replace SQL
Complement relational databases
Focus on scalability and flexibility
Why NOSQL Systems Emerged
Modern applications such as:
Gmail
Facebook
Google Maps
E-commerce platforms
Social networks
generate:
Huge amounts of heterogeneous data
Traditional relational systems faced challenges because:
Fixed schemas are restrictive
Complex joins reduce performance
Horizontal scalability is difficult
High concurrency creates bottlenecks
Example 1: Email Applications
Applications such as:
Gmail
Yahoo Mail
may contain:
Millions of users
Billions of emails
Challenges for Relational Systems
1. Too Many Services
Traditional SQL systems provide:
Complex query processing
Serializable transactions
Concurrency control
Many large-scale applications:
Do not require all these features
2. Restrictive Schema
Relational systems require:
Fixed predefined schemas
But modern applications often contain:
Semi-structured data
Dynamic attributes
Varying formats
Example 2: Facebook
Facebook stores:
User profiles
Friend relationships
Posts
Images
Videos
This data:
Is highly interconnected
Changes frequently
Is not fully structured
Hence:
Multiple data storage approaches are needed.
Development of NOSQL Systems
Large companies developed custom systems to handle:
Big data
Scalability
Distributed storage
Major NOSQL Systems
1. Google BigTable
Developed By
Google
Used In
Gmail
Google Maps
Web indexing
Category
Column-based / Wide-column store
Apache HBase
Description
Open-source system based on:
BigTable concepts
2. Amazon DynamoDB
Developed By
Amazon
Category
Key-value store
Characteristics
Fast access
Cloud-based
Highly scalable
3. Facebook Cassandra
Developed By
Facebook
Now Known As
Apache Cassandra
Characteristics
Combines:
Key-value concepts
Column-family concepts
4. MongoDB and CouchDB
Category
Document-based NOSQL systems
Features
JSON-like documents
Flexible schema
5. Graph Databases
Examples
Neo4J
GraphBase
Suitable For
Social networks
Relationship-heavy data
6. Hybrid NOSQL Systems
Example
OrientDB
Combines
Features from:
Graph databases
Document stores
Key-value systems
Other Systems Related to NOSQL
Object Databases
Based on:
Object-oriented data model
XML Databases
Store:
Native XML documents
24.1.2 Characteristics of NOSQL Systems
NOSQL systems differ from SQL systems in:
1. Distributed system characteristics
2. Data model and query language characteristics
A. Distributed System Characteristics
1. Scalability
Definition
Ability of a system to grow with increasing data and workload.
Types of Scalability
Horizontal Scalability
Adding:
More nodes
to the distributed system.
Node1 + Node2 + Node3 + ...
Preferred in NOSQL Systems
Vertical Scalability
Increasing:
CPU
RAM
Storage
of existing machines.
NOSQL Preference
NOSQL systems emphasize:
Horizontal Scalability
because:
It supports continuous expansion
Lower-cost commodity hardware can be used
Dynamic Scaling
NOSQL systems support:
Adding nodes while operational
Redistribution of data automatically
2. Availability, Replication, and Eventual Consistency
Availability
The system remains operational:
Even if some nodes fail
Replication
Definition
Maintaining copies of data on multiple nodes.
Advantages
✅ Improved availability
✅ Better fault tolerance
✅ Faster read performance
Challenge with Replication
Write operations become complex because:
All replicas must be updated
Eventual Consistency
Definition
All replicas will:
Eventually become consistent
but:
Immediate consistency is not guaranteed.
Serializable Consistency vs Eventual Consistency
Serializable Eventual
Consistency Consistency
Strong consistency Relaxed consistency
Immediate Delayed
synchronization synchronization
Slower writes Faster writes
Traditional SQL Many NOSQL
systems systems
3. Replication Models
Two major replication models exist:
1. Master-Slave Replication
2. Master-Master Replication
Master-Slave Replication
Characteristics
One master copy
Multiple slave copies
Write Operations
Must occur at:
Master node
Read Operations
May occur:
At master
Or at slave nodes
Advantages
✅ Simpler conflict management
✅ Easier consistency maintenance
Disadvantages
❌ Master bottleneck
❌ Single point of failure
Master-Master Replication
Characteristics
Reads and writes allowed at all replicas
Advantages
✅ High availability
✅ Better load balancing
Disadvantages
❌ Temporary inconsistencies
❌ Conflict resolution required
Conflict Resolution
Concurrent updates at different nodes require:
Reconciliation methods
4. Sharding of Files
Definition
Sharding is:
Horizontal partitioning of data
across multiple nodes.
Purpose of Sharding
✅ Load balancing
✅ Parallel access
✅ Improved scalability
✅ Faster processing
Example
Customer IDs:
1–1000 → Node A
1001–2000 → Node B
Sharding + Replication
Combining:
Sharding
Replication
improves:
Availability
Performance
Scalability
5. High-Performance Data Access
NOSQL systems optimize:
Fast object retrieval
using:
1. Hashing
2. Range Partitioning
Hashing
Method
Apply hash function:
h( K)
to object key K .
h( K)
Purpose
Determine storage location quickly.
Advantages of Hashing
✅ Very fast access
✅ Uniform distribution
Range Partitioning
Method
Objects stored according to:
Key ranges
K imin ≤ K ≤ K imax
K imin ≤ K ≤ K imax
Advantages of Range Partitioning
✅ Efficient range queries
✅ Ordered access
B. Data Model and Query Language Characteristics
1. Schema Flexibility
Traditional SQL
Requires:
Fixed schema
NOSQL Systems
Often:
Schema-less
Semi-structured
Semi-Structured Data
Common formats:
JSON
XML
JSON (JavaScript Object Notation)
Example:
{
"name": "Ashwin",
"department": "CSE",
"research": ["Blockchain", "AI"]
}
Advantages of Schema Flexibility
✅ Easy evolution of data
✅ Dynamic attributes
✅ Better adaptability
Limitation
Constraints must often be:
Enforced in application code
instead of database schema.
2. Less Powerful Query Languages
NOSQL systems emphasize:
Performance
Simplicity
rather than:
Complex querying
CRUD Operations
Most NOSQL systems support:
Operati
Meaning
on
Create Insert data
Retrieve
Read
data
Modify
Update
data
Remove
Delete
data
SCRUD Operations
Adds:
Search (Find)
operation.
Query APIs
NOSQL systems often provide:
Programming APIs
instead of full SQL support.
Lack of Joins
Many NOSQL systems:
Do not support joins directly
Joins must be:
Handled by applications
3. Versioning
Some NOSQL systems store:
Multiple versions of data
along with:
Timestamps
Advantages of Versioning
✅ Historical tracking
✅ Recovery support
✅ Conflict resolution
24.1.3 Categories of NOSQL Systems
NOSQL systems are classified into major categories.
1. Document-Based NOSQL Systems
Data Representation
Stored as:
Documents
usually in:
JSON format
Examples
MongoDB
CouchDB
Features
✅ Flexible schema
✅ Easy nested structures
✅ Fast document retrieval
2. Key-Value Stores
Structure
Key → Value
Examples
DynamoDB
Redis
Features
✅ Very fast access
✅ Simple model
✅ Highly scalable
3. Column-Based / Wide Column Stores
Structure
Data partitioned into:
Column families
Examples
BigTable
Cassandra
HBase
Features
✅ Efficient large-scale storage
✅ Compression advantages
✅ Supports versioning
4. Graph-Based NOSQL Systems
Data Model
Represented as:
Nodes
Edges
Examples
Neo4J
GraphBase
Features
✅ Efficient relationship traversal
✅ Ideal for social networks
✅ Fast graph queries
Additional Categories
5. Hybrid NOSQL Systems
Combine features from:
Multiple NOSQL categories
Example
OrientDB
6. Object Databases
Based on:
Object-oriented concepts
7. XML Databases
Store:
Native XML data
Search Engines as NOSQL Stores
Large search engines also:
Store massive data
Support rapid search
Hence can be viewed as:
Big-data NOSQL systems
Key Differences: SQL vs NOSQL
SQL Systems NOSQL Systems
Structured
Flexible schema
schema
Strong
Eventual consistency
consistency
Vertical scaling Horizontal scaling
Complex joins Simple access patterns
BASE principles often
ACID properties
used
Relational model Multiple data models
Key Concepts Summary
Concept Description
NOSQL Not Only SQL systems
Horizontal
Adding more nodes
Scalability
Eventual
Delayed consistency
Consistency
Replication Multiple data copies
Sharding Horizontal partitioning
Create, Read, Update,
CRUD
Delete
Semi-structured data
JSON
format
Stores JSON-like
Document Store
documents
Key-Value Store Key-based access
Graph Database Node-edge data model
Advantages of NOSQL Systems
✅ High scalability
✅ High availability
✅ Flexible schemas
✅ Fast access
✅ Handles big data efficiently
✅ Suitable for cloud computing
Limitations of NOSQL Systems
❌ Weak consistency in some systems
❌ Limited query capabilities
❌ Lack of standardization
❌ Complex application logic
❌ Joins often unsupported
Lecture Notes: The CAP Theorem
24.2 The CAP Theorem
Introduction
Distributed database systems often use:
Replication
Multiple nodes
Distributed processing
to improve:
Availability
Scalability
Performance
However, maintaining:
Consistency among replicas
becomes a major challenge.
The:
CAP Theorem
explains the trade-offs involved in distributed systems.
Background: ACID Properties in Distributed Databases
Traditional distributed databases attempt to guarantee:
ACID Properties
Property Meaning
Atomicity All-or-nothing execution
Consistenc
Database constraints preserved
y
Concurrent transactions behave
Isolation
serially
Durability Committed changes persist
Replication in Distributed Systems
To improve:
Availability
Fault tolerance
distributed systems store:
Multiple copies of data items
across:
Multiple nodes
Problem with Replication
Suppose data item:
X
has copies at different nodes.
Example Scenario
Transaction T 1updates copy of X at Node 1
Transaction T 2updates another copy at Node 2
Then:
Two inconsistent copies of X may exist.
Consequence
If:
Transaction T 3reads from Node 1
Transaction T 4reads from Node 2
they may obtain:
Different values of the same item
leading to:
Inconsistency
Strong Consistency
Traditional SQL systems enforce:
Serializable Consistency
This ensures:
All transactions appear serial
Replicated copies remain synchronized
Drawback of Strong Consistency
Enforcing serializability introduces:
High overhead
Complex synchronization
Reduced performance
especially in:
Large-scale distributed systems
NOSQL Perspective
NOSQL systems prioritize:
Scalability
Availability
Performance
Hence:
Strong consistency is often relaxed
The CAP Theorem
The CAP theorem explains the trade-offs in:
Distributed systems with replication
CAP Stands For
Lett
Property
er
C Consistency
A Availability
Partition
P
Tolerance
1. Consistency (C)
Definition
All nodes see:
The same data
At the same time
after an update.
Meaning in CAP
Consistency means:
Every read receives the latest write
or:
An error message
Example
If one replica updates:
X =100
then all nodes should immediately reflect:
X =100
Important Clarification
Consistency in:
CAP
is different from consistency in:
ACID
CAP Consistency
Refers to:
Consistency among replicated copies
ACID Consistency
Refers to:
Preservation of integrity constraints
Relationship Between Them
If consistency among replicas is treated as:
A database constraint
then:
CAP consistency and ACID consistency become related.
2. Availability (A)
Definition
Every request receives:
A successful response
or:
A failure message
Meaning
The system remains:
Operational
Accessible
even when some nodes fail.
Availability Ensures
✅ Fast response
✅ Continuous service
✅ No unnecessary downtime
3. Partition Tolerance (P)
Definition
The system continues functioning:
Even if network failures partition nodes
Network Partition
Occurs when:
Nodes cannot communicate with each other
Example
Partition 1: Node A ↔ Node B
Partition 2: Node C ↔ Node D
But:
Nodes in different partitions cannot communicate.
Partition Tolerance Means
The system:
Continues operating despite communication failures.
The CAP Theorem Statement
A distributed system:
Cannot simultaneously guarantee all three:
Consistency
Availability
Partition Tolerance
CAP Trade-Off
Only:
Two out of the three properties
can be fully guaranteed at a time.
CAP Triangle
Consistency
/\
/ \
/ \
/ \
Availability----Partition Tolerance
Choosing Two Properties
Distributed system designers must prioritize:
Which two properties are most important.
CP Systems (Consistency + Partition Tolerance)
Characteristics
Maintain consistency
Continue during partitions
May sacrifice availability
Behavior
If partition occurs:
Some operations may be rejected
to preserve:
Correctness
Example Systems
Traditional distributed SQL databases often prefer:
CP
AP Systems (Availability + Partition Tolerance)
Characteristics
Always available
Continue during partitions
May allow temporary inconsistency
Behavior
Nodes may:
Return outdated values temporarily
Example
Many NOSQL systems prefer:
AP
because:
High availability is critical
CA Systems (Consistency + Availability)
Characteristics
Strong consistency
High availability
but:
No partition tolerance
Limitation
In distributed environments:
Network partitions are unavoidable
Hence:
Pure CA systems are impractical for large-scale distributed systems.
Why NOSQL Systems Prefer AP
Modern Web applications require:
Continuous operation
Fast response
Large-scale scalability
Hence:
Availability and partition tolerance are prioritized.
Eventual Consistency
Many NOSQL systems use:
Eventual Consistency
Definition
Replicas may temporarily differ,
but:
Eventually all replicas become consistent.
Example
Suppose:
X =50
is updated to:
X =100
One node may show:
50
while another shows:
100
temporarily.
Eventually:
All replicas converge to:
100
Advantages of Eventual Consistency
✅ High availability
✅ Better performance
✅ Faster writes
✅ Improved scalability
Drawbacks of Eventual Consistency
❌ Temporary stale reads
❌ Conflicting updates possible
❌ More complex application logic
Strong vs Weak Consistency
Weak/Eventual
Strong Consistency
Consistency
Immediate
Delayed synchronization
synchronization
Slower performance Faster performance
Higher overhead Lower overhead
Traditional SQL Many NOSQL systems
CAP in SQL vs NOSQL
SQL Systems NOSQL Systems
Prefer
Prefer availability
consistency
Eventual
ACID guarantees
consistency
Strong Relaxed
transactions transactions
Less scalable Highly scalable
Impact of CAP Theorem
The CAP theorem influences:
System architecture
Replication strategy
Consistency models
Performance optimization
Real-World Examples
System Type Preferred CAP Properties
Banking Consistency + Partition
systems tolerance
Availability + Partition
Social media
tolerance
E-commerce Availability + Partition
carts tolerance
Airline
Consistency prioritized
reservations
Consistency Models in NOSQL
Different NOSQL systems use:
Different levels of consistency
Examples:
Eventual consistency
Tunable consistency
Causal consistency
Read-your-writes consistency
Importance of CAP Theorem
The CAP theorem helps designers:
Understand unavoidable trade-offs
Build scalable distributed systems
Select suitable architectures
Key Concepts Summary
Concept Description
Impossible to guarantee C, A, and P
CAP Theorem
simultaneously
Consistency Same data visible on all replicas
Concept Description
Availability System always responds
Partition Tolerance Operates despite network failures
Eventual Consistency Replicas become consistent over time
Serializable
Strongest consistency level
Consistency
Replication Multiple copies of data
Advantages of Relaxed Consistency
✅ Improved scalability
✅ Faster response time
✅ Better fault tolerance
✅ High availability
Challenges in Distributed Systems
❌ Synchronization overhead
❌ Network partitions
❌ Replication conflicts
❌ Temporary inconsistency
Lecture Notes: Document-Based NOSQL Systems and MongoDB
24.3 Document-Based NOSQL Systems and MongoDB
Introduction
Document-based NOSQL systems store data as:
Collections of Documents
These systems are also called:
Document-oriented databases
Document stores
Characteristics of Document-Based Systems
Documents:
Are self-describing
Do not require a fixed schema
Can contain nested structures
Can have varying attributes
Comparison with Relational Databases
Relational
Document-Based NOSQL
DBMS
Fixed schema Schema-less
Tables and rows Collections and documents
Structured data Semi-structured data
Embedded documents
Joins required
possible
Strict
Flexible denormalization
normalization
Similarity with XML and Objects
Documents resemble:
Complex objects
XML documents
But unlike XML or relational systems:
Schema specification is optional
Self-Describing Data
Each document contains:
Its own structure
Field names
Data types
This is known as:
Self-describing data
Flexible Document Structure
Documents in the same collection:
Need not have identical fields
Example:
{
"name": "Ashwin",
"department": "CSE"
}
Another document:
{
"name": "Priya",
"department": "ISE",
"research": "AI"
}
Advantages of Flexible Schema
✅ Easy evolution of applications
✅ Dynamic attributes
✅ Better adaptability
✅ Reduced schema maintenance
Popular Document-Based NOSQL Systems
Examples include:
MongoDB
CouchDB
MongoDB
Overview
MongoDB is:
A document-oriented NOSQL database
Open-source
Highly scalable
Distributed
24.3.1 MongoDB Data Model
BSON Format
MongoDB stores documents in:
BSON (Binary JSON)
BSON
Definition
Binary representation of JSON with:
Additional data types
More efficient storage
JSON Example
{
"project_name": "Smart Traffic System",
"department": "CSE"
}
Collections
Documents are grouped into:
Collections
Similar to:
Tables in relational systems
Creating a Collection
Syntax:
[Link]("project")
Example with Options
[Link](
"project",
{
capped: true,
size: 1310720,
max: 500
}
)
Capped Collections
Definition
Collections with:
Fixed storage size
Maximum document limit
Advantages
✅ Better performance
✅ Predictable storage
✅ Efficient insertion
Another Collection Example
[Link](
"worker",
{
capped: true,
size: 5242880,
max: 2000
}
)
ObjectId Field
Every document has:
_id field
which uniquely identifies the document.
Characteristics of _id
✅ Automatically indexed
✅ Unique within collection
✅ Similar to primary key
System-Generated ObjectId
MongoDB can automatically generate:
ObjectId
Components of ObjectId
System-generated ObjectId contains:
Compone
Size
nt
Timestamp 4
Compone
Size
nt
bytes
3
Node ID
bytes
2
Process ID
bytes
3
Counter
bytes
User-Defined ObjectId
Users may also define:
Their own _id values
provided they are:
Unique
Schema-less Design
MongoDB collections:
Do not require predefined schemas
Design Approaches
MongoDB supports:
1. Denormalized design
2. Normalized design
Denormalized Design
Related data embedded inside one document.
Example
{
"_id": "P1",
"project_name": "Database System",
"workers": [
{
"worker_id": "W1",
"name": "Ashwin"
},
{
"worker_id": "W2",
"name": "Priya"
}
]
}
Advantages of Denormalization
✅ Faster reads
✅ No joins needed
✅ Better performance
Disadvantages
❌ Data redundancy
❌ Difficult updates
Arrays in MongoDB
Square brackets:
[ ... ]
represent:
Arrays of values
Reference-Based Design
Instead of embedding workers:
Store references to worker documents
Example
{
"_id": "P1",
"workers": ["W1", "W2"]
}
Normalized Design
Separate collections for:
Projects
Employees
Works_on
Similar to:
Relational normalization
Design Choice Depends On
Access patterns
Query requirements
Performance needs
24.3.2 MongoDB CRUD Operations
CRUD stands for:
Lett Operati
er on
C Create
R Read
U Update
D Delete
1. Insert Operation
Used to:
Create documents
Syntax
db.collection_name.insert(document)
Example
[Link]({
"_id": "P1",
"project_name": "AI Research"
})
Multiple Document Insert
[Link]([
{...},
{...}
])
2. Remove Operation
Used to:
Delete documents
Syntax
db.collection_name.remove(condition)
Example
[Link]({
"_id": "P1"
})
3. Update Operation
Used to:
Modify documents
Example
[Link](
{ "_id": "P1" },
{ $set: { "project_name": "Blockchain Research" } }
)
$set Operator
Used to:
Modify specified fields
without replacing the entire document.
4. Find Operation
Used for:
Reading/querying documents
Syntax
db.collection_name.find(condition)
Example
[Link]({
"department": "CSE"
})
Query Conditions
MongoDB supports:
Boolean conditions
Comparison operators
Logical operators
24.3.3 MongoDB Distributed System Characteristics
MongoDB supports:
Replication
Sharding
Distributed transactions
Atomicity in MongoDB
Most updates are atomic:
At single-document level
Multi-Document Transactions
MongoDB also supports:
Transactions across multiple documents
using:
Two-Phase Commit Protocol
Replication in MongoDB
MongoDB uses:
Replica Sets
Replica Set
A group of nodes storing:
Copies of the same data
Components of Replica Set
1. Primary Node
2. Secondary Nodes
3. Arbiter (optional)
Primary Node
Handles:
All write operations
Secondary Nodes
Store:
Replicated copies
Arbiter
Participates in:
Elections
but:
Does not store data
Replica Set Example
Primary → Secondary1
→ Secondary2
→ Arbiter
Replication Process
1. Write applied to primary
2. Changes propagated to secondaries
Read Preferences
MongoDB allows:
Different read configurations
Default Read Preference
All reads occur at:
Primary node
Advantages
✅ Latest data guaranteed
Secondary Reads
Reads can also occur at:
Secondary nodes
Advantage
✅ Improved read performance
Limitation
❌ Secondary may not have latest updates
due to:
Replication delay
Failover in MongoDB
If primary fails:
Election selects new primary
Importance of Odd Number of Nodes
Replica sets should contain:
Odd number of voting members
to avoid:
Tie during election
Sharding in MongoDB
Definition
Horizontal partitioning of documents.
Purpose of Sharding
✅ Load balancing
✅ Horizontal scalability
✅ Faster processing
✅ Distributed storage
Shards
Each partition is called:
Shard
Shard Key
MongoDB partitions documents using:
Shard Key
Requirements of Shard Key
Must:
1. Exist in every document
2. Be indexed
Sharding Methods
MongoDB supports:
1. Range Partitioning
2. Hash Partitioning
Range Partitioning
Documents partitioned according to:
Key ranges
Example
1–1,000,000 → Shard 1
1,000,001–2,000,000 → Shard 2
Advantages
✅ Efficient range queries
Hash Partitioning
Uses:
h( K)
h( K)
where:
K = shard key
Purpose
Distributes documents:
Uniformly across shards
Advantages
✅ Better load balancing
✅ Randomized distribution
Range vs Hash Partitioning
Range
Hash Partitioning
Partitioning
Efficient range Efficient random
queries access
Ordered
Uniform distribution
distribution
Possible hotspot Better load
issue balancing
Query Router
MongoDB uses:
Query Router
to:
Determine which shard contains requested data
Query Routing Process
1. Query submitted
2. Router identifies target shards
3. Query forwarded to appropriate nodes
Combination of Sharding and Replication
Techniqu
Purpose
e
Scalability and load
Sharding
balancing
Replicatio Availability and fault
n tolerance
Additional MongoDB Features
MongoDB also supports:
Indexing
Security
Aggregation
Administration tools
Advantages of MongoDB
✅ Flexible schema
✅ Horizontal scalability
✅ High availability
✅ Efficient distributed storage
✅ Fast reads and writes
✅ Easy JSON integration
Limitations of MongoDB
❌ Weaker consistency compared to SQL
❌ Complex joins not efficient
❌ Potential redundancy in denormalized design
Key Concepts Summary
Concept Description
BSON Binary JSON format
Collection Group of documents
Self-describing data
Document
object
Unique document
ObjectId
identifier
Create, Read, Update,
CRUD
Delete
Group of replicated
Replica Set
nodes
Primary Node Handles writes
Concept Description
Secondary
Stores replicas
Node
Sharding Horizontal partitioning
Shard Key Partitioning attribute
Lecture Notes: 24.4 NOSQL Key-Value Stores
1. Introduction to Key-Value Stores
Key-value stores are a category of NOSQL databases designed for:
High performance
Scalability
High availability
Fast data retrieval
They store data as:
( Key → Value )
Where:
Key = Unique identifier
Value = Actual data object
The system retrieves data quickly using the key.
2. Basic Concept of Key-Value Stores
Structure
Compone
Description
nt
Key Unique identifier
Data associated with the
Value
key
Example:
Key Value
100 {Name: "Ashwin", Dept:
1 "CSE"}
100 {Name: "Ravi", Dept:
2 "ISE"}
3. Characteristics of Key-Value Stores
Important Features
1. Fast Access
Data retrieval uses the key directly.
Lookup operations are extremely fast.
2. Scalability
Easily expanded by adding more nodes.
Supports horizontal scaling.
3. Availability
Data is replicated across multiple nodes.
System continues operating even if some nodes fail.
4. Simple Data Model
Minimal structure.
Flexible storage.
5. Distributed Storage
Data distributed across cluster nodes.
4. Types of Data Stored
Different systems support different value formats.
Type Example
Unstructured Byte arrays
Semi- JSON
structured documents
Structured Records/Tuples
5. DynamoDB Overview
Introduction
Amazon DynamoDB
DynamoDB is a cloud-based key-value NOSQL database developed by
Amazon.
It is part of:
AWS (Amazon Web Services)
SDK-based cloud platforms
6. DynamoDB Data Model
Main Concepts
a) Tables
Collection of items
No fixed schema required
b) Items
Individual records/objects
c) Attributes
Fields inside items
Example Item
{
"EmpID": 101,
"Name": "Ashwin",
"Dept": "CSE"
}
7. Primary Keys in DynamoDB
A primary key is mandatory.
Two types exist:
7.1 Hash Type Primary Key
Single attribute
Used for hashing
Example:
EmpID
Characteristics
Fast lookup
Items not stored in order
7.2 Hash and Range Type Primary Key
Uses two attributes:
( A,B)
Where:
A → Hash key
B → Range key
Example:
Hash Range
Key Key
Timestam
ItemID
p
Advantages
Supports multiple versions
Enables ordered retrieval
8. Voldemort Key-Value Store
Introduction
Project Voldemort
Open-source NOSQL system
Inspired by Amazon DynamoDB
Used by LinkedIn
Focuses on:
Performance
Horizontal scalability
Replication
Sharding
9. Basic Operations in Voldemort
Assume store name = s
9.1 Put Operation
[Link](k,v)
Stores value v using key k.
9.2 Get Operation
v = [Link](k)
Retrieves value associated with key k.
9.3 Delete Operation
[Link](k)
Deletes item with key k.
10. Data Formats in Voldemort
Values may be stored as:
Byte arrays
JSON
Serialized objects
Applications provide:
Serializer classes
Conversion functions
11. Consistent Hashing
Definition
A distributed data placement technique used to:
Distribute data among nodes
Support scalability
Support replication
Working Principle
A hash function:
h(k )
is applied to key k.
The hash value determines the node storing the item.
12. Ring-Based Architecture
The hash space is visualized as a circular ring.
Nodes are placed on the ring.
Items are stored at the node immediately following the hash value in
clockwise direction.
Advantages
a) Horizontal Scalability
Easy addition of new nodes
b) Load Balancing
Data distributed evenly
c) Fault Tolerance
Failure handled gracefully
d) Minimal Data Movement
Only some data redistributed when adding/removing nodes
13. Replication in Voldemort
Multiple copies of data are stored on successive nodes.
Benefits:
High availability
Fault tolerance
14. Sharding
Definition
Horizontal partitioning of data across nodes.
Purpose
Improve performance
Balance workload
Benefits
Benefit Description
Faster access Smaller datasets per
Benefit Description
node
More nodes can be
Scalability
added
Better
Requests distributed
concurrency
15. Consistency and Versioning
Voldemort allows:
Concurrent writes
Multiple versions of same item
Each version has:
Vector clock
Read Repair
When reading:
Multiple versions may exist
System reconciles versions
Application may assist in conflict resolution
16. Other Key-Value Stores
16.1 Oracle NoSQL Database
Oracle NoSQL Database
Oracle’s key-value NOSQL system
Distributed architecture
High scalability
16.2 Redis
Redis
Features:
In-memory storage
Extremely fast access
Persistence support
Master-slave replication
Used for:
Caching
Real-time applications
16.3 Apache Cassandra
Apache Cassandra
Hybrid NOSQL system combining:
Key-value concepts
Column-family concepts
Features:
High availability
Scalability
Distributed architecture
Used by:
Facebook
Large-scale enterprise systems
17. Advantages of Key-Value Stores
Advantage Description
Very fast
Key-based lookup
retrieval
High Supports distributed
scalability systems
Advantage Description
Flexible
No rigid structure
schema
High
Replication support
availability
Good for big
Massive storage support
data
18. Limitations of Key-Value Stores
Limitation Description
Limited
Mostly key-based access
querying
Weak
Eventual consistency often used
consistency
Application must handle
No joins
relationships
Less structured Harder for complex analytics
19. Applications of Key-Value Stores
Common Use Cases
Social media systems
Shopping carts
Session management
User profiles
Real-time analytics
Recommendation systems
Caching systems
Cloud applications
20. Comparison: SQL vs Key-Value Stores
SQL Key-Value
Feature
DBMS Store
Schema Fixed Flexible
Query
SQL API operations
Language
Supporte Usually
Joins
d absent
Scalability Vertical Horizontal
Consistency Strong Eventual
Performance Moderate Very high
21. Key Terms
Term Meaning
Key Unique identifier
Value Associated data
Sharding Horizontal partitioning
Replication Multiple copies of data
Consistent Data distribution
Hashing technique
Version tracking
Vector Clock
method
Read Repair Conflict reconciliation
22. Summary
Key-value stores are simple and highly scalable NOSQL systems.
Data is stored as (key, value) pairs.
DynamoDB and Voldemort are important examples.
Consistent hashing enables:
o Sharding
o Replication
o Scalability
These systems prioritize:
o Availability
o Performance
o Horizontal scaling
Commonly used in modern cloud and big data applications.
Lecture Notes: 24.5 Column-Based or Wide Column NOSQL
Systems
1. Introduction to Column-Based NOSQL Systems
Column-based NOSQL systems are designed for:
Huge volumes of distributed data
High scalability
Fast data access
Efficient storage for big data applications
These systems are also called:
Wide Column Stores
Column Family Stores
Examples
Apache HBase
Apache Cassandra
Google BigTable
2. BigTable and HBase
Google BigTable
Developed by Google
Used in:
o Gmail
o Google Search
o Google Maps
Uses:
o Google File System (GFS)
Apache HBase
Open-source system similar to BigTable
Typically uses:
o Hadoop Distributed File System (HDFS)
Can also use:
o Amazon S3 storage
HBase Characteristics
Distributed
Persistent
Scalable
Sorted map of key-value pairs
3. Key Difference from Key-Value Stores
In traditional key-value systems:
Key → single identifier
In column-based systems:
Key is multidimensional
HBase Cell Key Components
A cell is identified by:
( Table , RowKey ,ColumnFamily , ColumnQualifier ,Timestamp )
4. HBase Data Model
HBase organizes data using:
1. Namespaces
2. Tables
3. Rows
4. Column Families
5. Column Qualifiers
6. Columns
7. Cells
8. Versions and Timestamps
5. Tables and Rows
Tables
Data stored in tables
Each table has:
o Table name
Rows
Data stored as rows
Each row has:
o Unique Row Key
Important Property
Row keys:
Must be lexicographically ordered
Stored in sorted order
6. Column Families
Definition
A table contains one or more column families.
Characteristics
Specified during table creation
Cannot be changed later
Group related columns together
Similarity
Column families resemble:
Vertical partitioning in distributed databases
7. Column Qualifiers
Definition
Column qualifiers are dynamic column names.
Important Feature
Unlike relational databases:
Column qualifiers are NOT predefined
They are created dynamically
8. Columns in HBase
A column is identified as:
ColumnFamily :ColumnQualifier
Example:
personal:name
personal:email
academic:cgpa
9. Self-Describing Data
HBase supports self-describing data because:
Column qualifiers are created dynamically
Different rows can contain different qualifiers
This gives:
Flexibility
Schema evolution
Sparse storage
10. Versions and Timestamps
HBase supports multiple versions of data.
Each version has:
Timestamp
Version history
Timestamp Characteristics
Stored as long integer
Represents milliseconds since:
January 1 , 1970 UTC
Benefits
Historical data access
Recovery
Auditing
Time-based analysis
11. Cells in HBase
Definition
A cell stores one data value.
Cell Address
A cell is identified by:
( Table , RowKey ,ColumnFamily , ColumnQualifier ,Timestamp )
If timestamp is omitted:
Latest version is returned
12. Namespaces
Definition
A namespace is:
A collection of tables
Similar To
Database in relational systems
13. HBase CRUD Operations
HBase mainly supports low-level CRUD operations.
CRUD =
Create
Read
Update
Delete
Create Operation
Creates a table with column families.
Example:
create 'student', 'personal', 'academic'
Put Operation
Used for:
Insert
Update
Adding new versions
Example:
put 'student', '101', 'personal:name', 'Ashwin'
Get Operation
Retrieves one row.
Example:
get 'student', '101'
Scan Operation
Retrieves all rows.
Example:
scan 'student'
14. Important Observation
HBase:
Does NOT support joins directly
Provides only low-level operations
Therefore:
Complex operations handled by application programs
15. HBase Distributed Storage Architecture
Regions
Each table is divided into:
Regions
Each region stores:
Range of row keys
Why Ordered Row Keys?
Because:
Regions are based on row key ranges
Hence:
Row keys must be lexicographically ordered
16. Stores
Within each region:
Each column family stored in separate store
Thus:
Region
├── Store for Column Family 1
├── Store for Column Family 2
17. Region Servers
Regions are assigned to:
Region Servers
These are storage nodes responsible for:
Storing data
Processing requests
18. Master Server
HBase has:
One Master Server
Responsibilities:
Monitor region servers
Split tables into regions
Assign regions to servers
19. Apache Zookeeper
HBase uses:
Apache ZooKeeper
Purpose
Coordination
Synchronization
Naming services
Replication management
Features
Distributed coordination system
Keeps critical metadata in memory
Improves performance
20. Hadoop Distributed File System (HDFS)
HBase stores files using:
Apache Hadoop HDFS
HDFS Provides
Distributed storage
Fault tolerance
Scalability
21. Relationship Among Components
Applications
↓
HBase
↓
--------------------------------
| ZooKeeper | HDFS |
--------------------------------
22. Advantages of Column-Based NOSQL Systems
High Scalability
Easily add nodes
High Performance
Fast reads/writes
Efficient Storage
Sparse data handled efficiently
Flexible Schema
Dynamic columns
Versioning Support
Multiple versions of data
Distributed Architecture
Fault tolerance
Availability
23. Limitations
Complex Querying
No SQL joins
Application Complexity
Logic implemented in application layer
Data Modeling Challenges
Requires careful row-key design
Eventual Consistency
Strong consistency may not always exist
24. Applications of HBase / Column Stores
Used in:
Big data analytics
Social media systems
Web indexing
Log processing
Time-series data
IoT applications
Real-time analytics
25. Comparison: Relational DBMS vs HBase
Relational
Feature HBase
DBMS
Schema Fixed Flexible
Tables/
Data Model Column Families
Relations
Query
SQL CRUD APIs
Language
Joins Supported Not Supported
Scalability Vertical Horizontal
Versioning Limited Built-in
Column-family
Storage Row-oriented
oriented
26. Key Terms
Term Meaning
Group of related
Column Family
columns
Column Dynamic column
Qualifier name
Cell Basic storage unit
Range partition of
Region
rows
Region Server Stores regions
Timestamp Version identifier
Namespace Collection of tables
Distributed file
HDFS
system
ZooKeeper Coordination service
27. Summary
Column-based NOSQL systems are optimized for big data.
HBase is an open-source implementation inspired by Google
BigTable.
Data is organized using:
o Tables
o Rows
o Column Families
o Column Qualifiers
HBase supports:
o Horizontal scalability
o Distributed storage
o Versioning
o High availability
HDFS and ZooKeeper provide storage and coordination support.
HBase is widely used in cloud and big data applications.
Lecture Notes: 24.6 NOSQL Graph Databases and Neo4j
1. Introduction to Graph Databases
Graph databases are a category of NOSQL systems where data is
represented as a graph.
A graph consists of:
Nodes (Vertices) → represent entities
Relationships (Edges) → represent connections between entities
Both nodes and relationships can store data.
2. Characteristics of Graph Databases
Main Features
Relationship-oriented data model
Efficient traversal of connected data
Flexible schema
Suitable for highly connected data
Fast graph traversal operations
3. Examples of Graph Databases
Popular graph databases include:
Neo4j
OrientDB
TigerGraph
4. Neo4j Overview
Neo4j
Open-source graph database
Implemented in Java
Uses graph-oriented storage model
Supports ACID transactions
Provides Cypher query language
5. Neo4j Data Model
Neo4j organizes data using:
1. Nodes
2. Relationships
3. Labels
4. Relationship Types
5. Properties
6. Paths
6. Nodes
Definition
Nodes represent entities or objects.
Examples:
Employee
Department
Project
Location
7. Labels
Definition
Labels classify nodes into groups.
Example labels:
EMPLOYEE
DEPARTMENT
PROJECT
LOCATION
Important Points
A node can have:
o Zero labels
o One label
o Multiple labels
8. Multiple Labels
A node may belong to several categories.
Example:
PERSON:EMPLOYEE:MANAGER
This resembles:
Superclass and subclass concepts in EER model
9. Properties
Definition
Properties store data associated with nodes or relationships.
Properties are represented as:
{name : value}
Example:
{Lname:'Smith', Fname:'John', Minit:'B'}
10. Relationships
Definition
Relationships connect nodes.
Characteristics
Directed
Have relationship types
Can also contain properties
11. Relationship Types
Examples:
WorksFor
Manager
LocatedIn
WorksOn
These classify relationships.
12. Directed Relationships
Neo4j relationships are directional.
Example:
(Employee)-[:WorksFor]->(Department)
Even though directed:
Traversal possible in both directions
13. Relationships with Properties
Relationships can also store data.
Example:
[:WorksOn {Hours:20}]
Here:
Hours is a relationship property
14. Paths
Definition
A path is a traversal through connected nodes and relationships.
Example:
(Employee)-[:WorksOn]->(Project)
Paths are heavily used in graph queries.
15. Comparison with ER/EER Model
Neo4j ER/EER
Concept Equivalent
Node Entity
Entity
Node Label
Type/Subclass
Relationship
Relationship
Instance
Relationship
Relationship Type
Type
Property Attribute
16. Differences from ER/EER Model
Neo4j Relationships are Directed
ER relationships are not directed.
Nodes May Have No Label
ER entities must belong to entity type.
Neo4j is Operational
Neo4j is:
Real database system
High-performance distributed DBMS
ER/EER:
Mainly design models
17. Optional Schema
Neo4j allows:
Schema-less databases
However:
Constraints
Indexes
can optionally be defined.
18. Constraints
Neo4j supports:
Uniqueness constraints
Property constraints
Example:
Employee ID must be unique
19. Indexing
Indexes improve retrieval performance.
Example indexes:
EmpId
Dno
Pno
20. Internal Node Identifiers
Neo4j automatically creates:
Unique internal IDs for nodes
These are system-generated.
21. Cypher Query Language
Neo4j uses:
Cypher query language
Cypher is:
Declarative
Pattern-based
Graph-oriented
22. Important Cypher Clauses
Clause Purpose
Create
CREATE
nodes/relationships
MATCH Find patterns
RETURN Display results
WHERE Apply conditions
Pass intermediate
WITH
results
Clause Purpose
ORDER
Sort results
BY
LIMIT Restrict output size
23. CREATE Command
Used to create nodes and relationships.
Example:
CREATE (e:EMPLOYEE {EmpId:1, Ename:'John'})
24. MATCH Clause
Used to specify graph patterns.
Example:
MATCH (e:EMPLOYEE)
RETURN e
25. RETURN Clause
Specifies query output.
Example:
RETURN [Link]
26. Example Query 1
Retrieve department locations.
MATCH (d:DEPARTMENT)-[:LocatedIn]->(loc:LOCATION)
WHERE [Link] = 5
RETURN loc
27. Example Query 2
Retrieve projects worked on by employee.
MATCH (e:EMPLOYEE)-[w:WorksOn]->(p:PROJECT)
WHERE [Link] = 2
RETURN p, [Link]
28. Example Query 3
Retrieve employees working on project 2.
MATCH (e:EMPLOYEE)-[w:WorksOn]->(p:PROJECT)
WHERE [Link] = 2
RETURN e, [Link]
29. ORDER BY Clause
Used for sorting.
Example:
ORDER BY [Link]
30. LIMIT Clause
Restricts number of returned rows.
Example:
LIMIT 10
31. WITH Clause
Used to:
Separate query stages
Perform aggregation
32. Aggregation Example
Find employees working on more than two projects.
MATCH (e:EMPLOYEE)-[:WorksOn]->(p:PROJECT)
WITH e, count(p) AS TotalProjects
WHERE TotalProjects > 2
RETURN e, TotalProjects
33. Graph Visualization
Neo4j supports:
Visual graph representation
This allows:
Graphical display of query results
34. Updating Nodes
Properties can be added or modified.
Example:
SET [Link] = 'Manager'
35. Neo4j Interfaces
Neo4j supports multiple interfaces:
Native Java API
Python Drivers
PHP Drivers
REST API
Cypher Interface
36. REST API
Neo4j supports:
Representational State Transfer (REST)
Allows:
Web-based communication
Integration with applications
37. Enterprise Edition vs Community Edition
Communi Enterpri
Feature
ty se
Cypher Yes Yes
ACID Support Yes Yes
Advance
Clustering Limited
d
Advance
Caching Basic
d
Advance
Replication Basic
d
Performance Enhance
Limited
Tools d
38. ACID Properties
Neo4j supports:
Atomicity
Consistency
Isolation
Durability
Thus:
Reliable transactions supported
39. Master-Slave Replication
Neo4j can run on distributed clusters.
Configuration
One Master Node
Multiple Slave Nodes
Features
Data replication
Synchronization
High availability
40. Caching
Neo4j supports:
Main memory caching
Benefits:
Faster query execution
Improved traversal speed
41. Logical Logs
Neo4j maintains:
Transaction logs
Used for:
Recovery
Fault tolerance
42. Advantages of Graph Databases
Efficient Relationship Traversal
Very fast for connected data.
Flexible Schema
Easy schema evolution.
Natural Representation
Suitable for networked data.
High Performance
Optimized graph operations.
43. Limitations
Complex Distributed Processing
Large graph distribution is challenging.
Not Ideal for Tabular Analytics
Relational systems better for heavy aggregation.
Query Complexity
Complex path queries can become expensive.
44. Applications of Graph Databases
Used in:
Social networks
Recommendation systems
Fraud detection
Knowledge graphs
Network management
Routing systems
Bioinformatics
Cybersecurity
45. Relational DBMS vs Graph Database
Relational Graph
Feature
DBMS Database
Data Model Tables Graphs
Relationshi
Joins Direct edges
ps
Schema Fixed Flexible
Expensive
Traversal Fast traversal
joins
Structured Connected
Best For
data data
46. Key Terms
Term Meaning
Node Entity/Object
Relationship Connection between nodes
Label Node classification
Relationship
Relationship classification
Type
Data associated with
Property
node/relationship
Path Traversal through graph
Cypher Query language
MATCH Pattern matching clause
CREATE Node/relationship creation
47. Summary
Graph databases store data as nodes and relationships.
Neo4j is a popular graph-oriented NOSQL system.
Cypher is used for graph querying.
Neo4j supports:
o Labels
o Properties
o Paths
o Directed relationships
Graph databases are highly efficient for connected data
applications.
Neo4j supports:
o ACID properties
o Replication
o Clustering
o Caching
o Visualization tools
Lecture Notes
25.1 What Is Big Data?
1. Introduction to Big Data
Definition of Big Data
Big Data refers to datasets whose size, complexity, and growth rate
exceed the capabilities of traditional database management systems
(DBMS) to:
Capture
Store
Manage
Process
Analyze data efficiently
Data Size Categories
Data Unit Size
Terabyte 10¹²
(TB) bytes
Petabyte 10¹⁵
(PB) bytes
Exabyte 10¹⁸
(EB) bytes
Big data may vary depending on:
Industry requirements
Data usage
Historical data storage
Type of applications
2. Characteristics of Big Data
The Gartner Group introduced the concept of the 3 Vs of Big Data in
2011:
1. Volume
2. Velocity
3. Variety
Later, researchers added:
4. Veracity
5. Value
3. Volume
Definition
Volume refers to the enormous amount of data generated, collected, and
stored.
Sources of High Data Volume
A. Sensor Data
Generated by:
Manufacturing plants
Smart devices
Environmental monitoring systems
B. Smart Devices
Examples:
Smart meters
Credit card readers
RFID devices
C. Social Media
Platforms such as:
Facebook
Twitter
generate billions of posts, tweets, images, and videos daily.
D. Multimedia Content
High-definition videos
Audio files
Images
YouTube uploads
E. Internet of Things (IoT)
Billions of interconnected devices continuously generate data.
4. Industrial Internet of Things (IIoT)
Definition
Industrial Internet of Things (IIoT) refers to interconnected industrial
devices that communicate and exchange data over the Internet.
Applications of IIoT
Traffic monitoring
Environmental monitoring
Smart healthcare
Inventory management
Energy management
Importance
IIoT improves:
Operational efficiency
Automation
Decision-making
Predictive maintenance
5. Examples of Big Data Generation
Application Type of Data
Area Generated
Gene sequencing Biological data
Real-time sensor
Traffic monitoring
data
Medical monitoring
Smart healthcare
data
Application Type of Data
Area Generated
Inventory tracking
RFID systems
data
User-generated
Social networks
content
Surveillance Video and image
systems data
6. Velocity
Definition
Velocity refers to:
The speed at which data is generated
The speed at which data is transmitted
The speed at which data is processed
Examples
A. Stock Market Transactions
Billions of transactions occur daily.
B. Social Media Streams
Real-time tweets
Facebook updates
Instagram uploads
C. Mobile Networks
Millions of call records are generated continuously.
D. Streaming Data
Continuous flow of:
Sensor data
Financial data
GPS data
7. Importance of Velocity
High-velocity data processing helps in:
Fraud detection
Real-time analytics
Trend identification
Emergency response systems
Online recommendation systems
8. Variety
Definition
Variety refers to the different forms and formats of data.
Big data contains:
Structured data
Semi-structured data
Unstructured data
9. Structured Data
Definition
Structured data follows a predefined data model.
Characteristics
Organized in rows and columns
Easy to store and query
Suitable for relational databases
Examples
Banking transactions
Student databases
Employee records
10. Semi-Structured Data
Definition
Semi-structured data does not follow a rigid structure but contains tags or
markers.
Examples
XML documents
JSON files
HTML pages
XML Applications
XML
Application
Type
BioML Biology
Geographic Information
GML
Systems
BeerXM
Brewing industry
L
11. Unstructured Data
Definition
Unstructured data lacks a predefined format or organization.
Examples
Emails
Blogs
Audio files
Videos
Images
PDFs
Social media posts
Web pages
Challenges
Difficult to process
Complex storage requirements
Requires advanced analytics tools
12. NoSQL Systems for Unstructured Data
Examples of NoSQL systems:
MongoDB
Neo4j
MongoDB
Stores:
Document-oriented data
JSON-like structures
Neo4j
Stores:
Graph-based data
Nodes and relationships
13. Veracity
Definition
Veracity refers to:
Data quality
Data accuracy
Data trustworthiness
Reliability of data sources
14. Features of Veracity
A. Credibility of Source
Determines whether data originates from a trusted source.
B. Suitability of Data
Checks whether data is appropriate for analysis.
15. Problems Related to Veracity
Big data may contain:
Incomplete data
Inaccurate data
Duplicate records
Noisy data
Uncertain information
Hence, data must undergo:
Quality testing
Validation
Cleansing
Credibility analysis
16. Importance of Big Data Technologies
Traditional database systems cannot efficiently manage:
Huge data volumes
Fast data generation
Complex data types
Therefore, new technologies were developed such as:
Hadoop
MapReduce
Distributed file systems
NoSQL databases
17. Hadoop and MapReduce Revolution
Hadoop
Apache Hadoop is a framework used for:
Distributed storage
Distributed processing
Big data analytics
MapReduce
MapReduce is a programming model used for parallel data processing.
These technologies form the foundation of modern big data systems.
18. Advantages of Big Data Analytics
Big data analytics helps organizations in:
Better decision-making
Predictive analysis
Customer behavior analysis
Fraud detection
Performance optimization
Personalized services
19. Applications of Big Data
Healthcare
Patient monitoring
Disease prediction
Banking
Fraud detection
Risk management
Social Media
Sentiment analysis
Trend detection
Smart Cities
Traffic regulation
Pollution monitoring
E-Commerce
Product recommendations
Customer analytics
20. Summary
Big Data refers to massive datasets beyond the capability of
traditional DBMS systems.
The major characteristics are:
o Volume
o Velocity
o Variety
o Veracity
IoT and social media contribute significantly to big data growth.
Structured, semi-structured, and unstructured data are all
components of big data.
Data quality and trustworthiness are important challenges.
Technologies such as Hadoop and MapReduce are designed to
process big data efficiently.
Lecture Notes
25.2 Introduction to MapReduce and Hadoop
1. Introduction
What is Hadoop?
Apache Hadoop is an open-source framework designed for:
Distributed storage
Distributed processing
Big data analytics
Hadoop is mainly used for processing huge volumes of structured and
unstructured data across clusters of commodity hardware.
2. Core Components of Hadoop
Hadoop mainly consists of two important components:
Component Purpose
HDFS (Hadoop Distributed File
Distributed storage
System)
Distributed data
MapReduce
processing
3. Historical Background of Hadoop
Origin of Hadoop
Hadoop originated from the need to build an open-source search engine
capable of processing massive web data.
Key Contributors
Doug Cutting
Mike Cafarella
4. Nutch Project
What is Nutch?
Apache Nutch was developed to:
Crawl web pages
Index web content
Build a scalable search engine
Nutch became an Apache Software Foundation project.
5. Influence of Google Technologies
The development of Hadoop was inspired by two important research
papers released by Google:
Yea
Technology
r
200
Google File System (GFS)
3
200 MapReduce Programming
4 Model
These papers introduced:
Distributed storage
Parallel data processing
Fault tolerance concepts
6. Development of Hadoop
Doug Cutting and Mike Cafarella:
Built a distributed file system
Developed a processing framework
Ported Nutch onto the new framework
This framework became known as Hadoop.
Naming of Hadoop
The name “Hadoop” came from:
The stuffed elephant toy belonging to Doug Cutting’s son.
7. Hadoop at Yahoo
In 2006:
Doug Cutting joined Yahoo
Yahoo adopted Hadoop to improve:
o Search engine processing
o Scalability
o Data analytics
Yahoo later:
Expanded Hadoop clusters
Improved Hadoop technology
Created large-scale Hadoop infrastructures
8. Growth of Hadoop Ecosystem
Major Milestones
Yea
Event
r
200 Hadoop development at
6 Yahoo
200
Launch of Cloudera
8
201
Hortonworks spin-off
1
Commercial Hadoop Companies
Cloudera
Hortonworks
9. Hadoop Adoption
Hadoop experienced rapid growth because:
It is open source
It supports massive scalability
It provides fault tolerance
It runs on commodity hardware
Large organizations adopted Hadoop for:
Big data analytics
Search engines
Machine learning
Data warehousing
10. Introduction to MapReduce
Definition
MapReduce is a programming model used for:
Parallel processing
Distributed computation
Large-scale data analysis
It was introduced by:
Jeffrey Dean
Sanjay Ghemawat
11. Features of MapReduce
MapReduce provides:
Automatic parallelization
Fault tolerance
Data distribution
Load balancing
Simplified programming
12. Motivation Behind MapReduce
Google needed to process:
Billions of web pages
Search logs
Web graphs
Search request statistics
Challenges included:
Data distribution
Fault handling
Parallel execution
Scalability
MapReduce solved these challenges efficiently.
13. Functional Programming Background
MapReduce concepts originated from:
Functional programming languages
Especially:
o LISP
The programming paradigm uses two major functions:
1. Map
2. Reduce
14. Key-Value Pair Model
MapReduce processes data in the form of:
(Key, Value)
Examples:
("student", "Ashwin")
("subject", "Big Data")
15. General Form of MapReduce Functions
Map Function
map[ K 1, V 1]→ List [K 2 ,V 2]
Purpose
Takes input key-value pairs
Produces intermediate key-value pairs
Reduce Function
reduce (K 2, List [V 2])→ List [ K 3 , V 3]
Purpose
Takes grouped intermediate values
Produces final output
16. Workflow of MapReduce
The execution of MapReduce consists of:
1. Input Splitting
2. Mapping
3. Shuffle and Sort
4. Reducing
5. Output Generation
17. Word Count Example
Objective
Count the frequency of each word in a document.
18. Map Phase
Input
Big Data Hadoop Hadoop
Mapper Output
(Big,1)
(Data,1)
(Hadoop,1)
(Hadoop,1)
The Map function:
Splits text into words
Emits:
(word,1)
19. Reduce Phase
The Reduce function:
Groups identical keys
Adds frequencies
Output
(Big,1)
(Data,1)
(Hadoop,2)
20. Data Types in Hadoop
Common Hadoop data types:
Data Type Description
Text String data
LongWritab Integer-like numeric
le type
21. Context Object in MapReduce
The context object is used to:
Send configuration information
Read data from HDFS
Output key-value pairs
Report task status
22. Applications of MapReduce
A. Distributed Grep
Searches for patterns in files.
Characteristics
Map-only task
No Reduce phase required
B. Reverse Web-Link Graph
Purpose
Creates:
(target URL, source URL)
Used in:
Search engines
Web analytics
C. Inverted Index
Purpose
Creates an index:
(word → list of document IDs)
Used in:
Search engines
Document retrieval systems
23. Components of a MapReduce Job
A MapReduce job contains:
Component Description
Map Task Mapping logic
Reduce Task Reducing logic
Input Path HDFS input location
Component Description
Structure of input
Input Format
data
Output storage
Output Path
location
Output
Structure of output
Format
Parallelism Number of reducers
24. JobTracker
Definition
JobTracker manages:
Job scheduling
Task execution
Monitoring of jobs
Responsibilities:
Assign tasks
Handle failures
Coordinate cluster processing
25. Hadoop Releases
Hadoop 1.x
Features
Original MapReduce framework
Improved HDFS
Added security features
Better support for HBase
26. Hadoop 2.x
Major improvements introduced in Hadoop 2.x include:
A. YARN
Apache Hadoop YARN stands for:
Yet Another Resource Negotiator
Features
Resource management
Job scheduling
Better scalability
B. New MapReduce Runtime
Runs on top of YARN.
Benefits:
Improved cluster utilization
Better performance
C. Improved HDFS
Features:
Federation support
Increased availability
Better scalability
27. Advantages of Hadoop and MapReduce
Open-source technology
Handles huge datasets
Fault tolerant
Scalable architecture
Parallel processing
Cost-effective
28. Limitations
Batch-oriented processing
Complex configuration
High latency
Requires skilled administration
29. Summary
Hadoop is an open-source big data processing framework.
Hadoop originated from the Nutch search engine project.
Google’s GFS and MapReduce papers inspired Hadoop development.
MapReduce processes large datasets using parallel computation.
Hadoop stores data using HDFS.
Hadoop evolved from version 1.x to 2.x with YARN support.
Hadoop became one of the most important big data technologies.
Lecture Notes
25.3 Hadoop Distributed File System (HDFS)
1. Introduction to HDFS
What is HDFS?
Hadoop Distributed File System (HDFS) is the storage component of the
Hadoop ecosystem.
It is designed to:
Store massive datasets
Run on clusters of commodity hardware
Support distributed data processing
Provide high-throughput access to data
2. Features of HDFS
HDFS provides:
Distributed storage
Fault tolerance
Scalability
High availability
Parallel data access
High throughput
3. Design Goals of HDFS
HDFS was developed with the following assumptions and goals:
4. Hardware Failure Handling
Assumption
Commodity hardware failures are common in large clusters.
Solution
HDFS automatically:
Detects failures
Recovers data
Replicates missing blocks
This provides:
Reliability
Continuous operation
5. Batch Processing Support
HDFS is optimized for:
Batch processing
Sequential file access
Large-scale analytics
Characteristics
High throughput is prioritized
Low latency is not the main goal
6. Support for Large Datasets
HDFS supports:
Gigabyte-scale files
Terabyte-scale files
Petabyte-scale storage
Large files are divided into blocks and distributed across cluster nodes.
7. Simple Coherency Model
HDFS follows:
Single writer
Multiple readers
Characteristics
Files cannot be modified directly
Data can only be appended
This simplifies:
Synchronization
Data consistency
8. HDFS Architecture
HDFS follows a:
Master-Slave Architecture
9. Components of HDFS Architecture
Component Function
NameNode Master node managing metadata
DataNode Slave node storing data blocks
Client Accesses files in HDFS
Secondary Performs checkpointing and backup
NameNode tasks
10. NameNode
Definition
The NameNode is the master server of HDFS.
Responsibilities include:
Managing file system namespace
Maintaining metadata
Tracking block locations
Managing file permissions
11. Metadata Managed by NameNode
The NameNode maintains:
File names
Directory structure
Ownership
Access permissions
File timestamps
Block mappings
12. Inodes in HDFS
HDFS uses:
Inodes (Index Nodes)
to store metadata information about:
Files
Directories
13. DataNodes
Definition
DataNodes are slave nodes responsible for storing actual file blocks.
Responsibilities:
Block storage
Read/write operations
Block replication
Block deletion
14. Working of DataNodes
Each DataNode:
Stores blocks in the local file system
Maintains block metadata
Sends status reports to NameNode
15. Secondary NameNode
Purpose
Secondary NameNode performs:
Checkpointing
Backup support
Functions
Merges checkpoint and journal files
Reduces NameNode restart time
16. File Read Operation in HDFS
Steps
Step 1
Client contacts NameNode.
Step 2
NameNode returns block locations.
Step 3
Client connects directly to DataNodes.
Step 4
Data blocks are read.
17. Advantages of Metadata and Data Separation
HDFS separates:
Metadata operations
Data transfer operations
Benefits
Faster metadata access
Better performance
Reduced network bottlenecks
18. Replication in HDFS
Purpose
Replication improves:
Reliability
Availability
Fault tolerance
19. Default Replication Factor
Typically:
3 Copies
of each block are maintained.
20. Benefits of Replication
Data recovery during failures
Improved read performance
Load balancing
High availability
21. Network Traffic Optimization
HDFS minimizes network traffic by:
Reading nearest replicas
Writing local replicas first
Avoiding unnecessary rack transfers
22. Rack Awareness in HDFS
Definition
Rack awareness means understanding network topology in data centers.
Distance Calculation
Distanc
Location
e
Same Node 0
Same Rack 2
Different 4
Distanc
Location
e
Rack
23. NameNode Storage Mechanism
The NameNode stores:
Namespace image
Journal logs
24. Journal File
Definition
A write-ahead commit log that records file system changes.
Purpose:
Recovery support
Fault tolerance
25. Checkpointing
Purpose
Checkpointing:
Saves stable system images
Reduces restart time
Improves recovery speed
26. Heartbeat Mechanism
Definition
DataNodes periodically send status information to NameNode.
This process is called:
Heartbeat
27. Block Report
Definition
A report sent by DataNodes containing:
Block ID
Block size
Generation stamp
28. Commands Sent by NameNode
The NameNode may instruct DataNodes to:
Replicate blocks
Delete replicas
Send reports
Shutdown or re-register
29. File I/O Operations in HDFS
HDFS supports:
File read
File write
File append
30. Write Pipeline in HDFS
Data is written using:
64 KB packets
through a write pipeline.
31. HDFS Checksum Mechanism
Purpose
Checksums detect:
Data corruption
Transmission errors
Process
Checksums generated for blocks
Verified during reads
32. Handling Corrupt Blocks
When corruption is detected:
1. Client informs NameNode
2. NameNode initiates replication
3. Corrupt block is removed
33. Replica Selection During Reads
HDFS attempts to:
Read nearest replica first
Reduce network usage
Improve performance
34. Block Placement Policy
Goal
Balance:
Reliability
Availability
Write cost
Read bandwidth
35. Replica Placement Strategy
HDFS places replicas:
One on local node
Others on different racks
Benefits:
Fault tolerance
Better scheduling flexibility
36. Replica Management
The NameNode tracks:
Replica count
Replica location
37. Over-Replication Handling
If too many replicas exist:
Extra replicas are removed
Disk utilization is balanced
38. HDFS Scalability
HDFS supports:
Thousands of nodes
Millions of files
Petabytes of storage
39. Yahoo HDFS Cluster Statistics
Yahoo achieved:
Feature Actual
Storage
14 PB
Capacity
Nodes 4,000
60
Files
Million
Clients 15,000
40. NameNode Memory Requirements
Observation
The NameNode requires large RAM because:
Metadata is stored in memory
Approximation
1 GB RAM ≈ 1 PB Storage
(under certain assumptions)
41. NameNode Bottleneck
Problem
Large numbers of:
Clients
Requests
MapReduce tasks
create bottlenecks at the NameNode.
42. Hadoop Ecosystem
The Hadoop ecosystem includes additional tools built on top of Hadoop.
43. Apache Pig
Apache Pig provides:
Pig Latin scripting language
Simplified data processing
Feature
Scripts are translated into:
Directed Acyclic Graph (DAG)
of MapReduce jobs.
44. Apache Hive
Apache Hive provides:
SQL-like querying
Data warehousing support
Features
SQL-92 support
Advanced analytics
SerDe mechanism
45. Oozie
Apache Oozie is used for:
Workflow scheduling
Job coordination
Supports:
MapReduce jobs
Hive queries
Pig scripts
46. Sqoop
Apache Sqoop is used to:
Transfer data between:
o Relational databases
o HDFS
47. HBase
Apache HBase is:
A NoSQL database
Built on HDFS
Applications
Time-series analysis
Streaming data
Data warehousing
Multi-dimensional lookups
48. Advantages of HDFS
Fault tolerance
Distributed storage
Scalability
Cost-effectiveness
High throughput
Automatic recovery
49. Limitations of HDFS
Not suitable for low-latency access
Single NameNode bottleneck
Limited random write capability
Complex management
50. Summary
HDFS is the distributed file system of Hadoop.
It uses a master-slave architecture.
NameNode manages metadata.
DataNodes store actual data blocks.
Replication improves reliability and availability.
HDFS supports large-scale distributed storage.
Hadoop ecosystem tools provide additional analytics and workflow
support.
Lecture Notes: 25.4 MapReduce – Additional Details
Introduction to MapReduce Runtime
The MapReduce framework provides a distributed runtime environment for
processing massive datasets over clusters of commodity hardware. It
follows a master-slave architecture and works closely with the Hadoop
Distributed File System (HDFS).
Main objectives of MapReduce runtime:
Parallel processing of large datasets
Fault tolerance
Scalability
Efficient task scheduling
Data locality optimization
25.4.1 MapReduce Runtime
1. Architecture of MapReduce Runtime
A. JobTracker (Master Node)
The JobTracker is the master component in Hadoop v1.
Responsibilities of JobTracker
1. Accepts and manages jobs submitted by clients
2. Splits jobs into tasks
3. Schedules tasks on cluster nodes
4. Monitors task execution
5. Handles failures and rescheduling
6. Tracks job completion
Functions
Job submission
Job initialization
Resource allocation
Task scheduling
Progress monitoring
B. TaskTracker (Worker Node)
Each worker machine runs a TaskTracker daemon.
Responsibilities of TaskTracker
Executes Map tasks
Executes Reduce tasks
Sends heartbeat messages to JobTracker
Reports status updates
Handles task failures
Performs Shuffle service
Important Features
One TaskTracker per worker node
Tasks run in separate JVM processes
Monitors task health continuously
Overall Flow of a MapReduce Job
The execution of a MapReduce job proceeds in several stages.
Step 1: Job Submission
Client submits:
Map and Reduce code (JAR file)
Configuration details
Input paths
Output paths
Job is sent to the JobTracker.
Step 2: Job Initialization
JobTracker:
Places job into Job Queue
Creates Map tasks from InputSplits
Creates Reduce tasks
Step 3: Task Assignment
Scheduler assigns tasks to TaskTrackers.
Important Concept: Data Locality
Tasks are scheduled on nodes where data already exists to reduce
network traffic.
Step 4: Task Execution
TaskTracker:
Launches JVM
Executes task
Monitors progress
Sends heartbeat updates
Step 5: Job Completion
After all tasks finish:
Intermediate files are cleaned
Output stored in HDFS
Job marked complete
Fault Tolerance in MapReduce
Fault tolerance is one of the strongest features of Hadoop.
There are three major failures:
1. Task Failure
Occurs due to:
Runtime exceptions
JVM crash
Task timeout
Recovery
TaskTracker informs JobTracker
JobTracker reschedules task
2. TaskTracker Failure
Occurs when:
Worker node crashes
Network disconnection occurs
Recovery
All incomplete tasks rescheduled
Completed map outputs recreated
3. JobTracker Failure
In Hadoop v1:
JobTracker is a Single Point of Failure (SPOF)
Manual restart required
All jobs must be resubmitted
This limitation was solved later using YARN in Hadoop v2.
Semantics During Failures
If Map and Reduce functions are deterministic:
Output remains correct
Duplicate completions ignored
Only first successful completion accepted
Shuffle Procedure in MapReduce
The Shuffle is the heart of MapReduce.
It ensures:
All records with same key reach same reducer
Shuffle has three phases:
1. Map Phase
Activities
Mapper outputs intermediate key-value pairs
Data buffered in memory
Data partitioned by reducer
Data sorted by keys
Components Used
Partitioner
Decides which reducer receives a key.
Comparator
Maintains sorting order.
Combiner
Performs local aggregation to reduce network traffic.
2. Copy Phase
Reducers fetch partitioned data from all mappers.
Features
Parallel copying
Managed by TaskTracker
Data transferred over network
3. Reduce Phase
Reducer:
Merges mapper outputs
Groups values by keys
Executes Reduce function
Produces final output
Job Scheduling in Hadoop
Scheduling determines how tasks are allocated to cluster resources.
1. FIFO Scheduler
Characteristics
Jobs processed in submission order
Simple implementation
Problems
Long jobs delay short jobs
Poor cluster utilization
2. Fair Scheduler
Goal:
Equal resource sharing
Features
Jobs grouped into pools
Fast response for small jobs
Better utilization
Advantages
Interactive queries execute faster
Balanced resource allocation
3. Capacity Scheduler
Designed for large enterprises.
Features
Multiple queues
Resource guarantees
Access control lists (ACLs)
Prevents monopolization
Benefits
Multi-tenant support
Department-wise resource allocation
25.4.2 Achieving Joins in MapReduce
Join is the most important relational algebra operation.
Suppose:
R(A,B)
S(B,C)
Join condition:
R . B=S . B
MapReduce supports several join strategies.
1. Sort-Merge Join
Process
Map phase partitions data
Shuffle sorts by join key
Reducer merges matching rows
Advantages
Works for very large datasets
Disadvantages
High shuffle cost
2. Map-Side Hash Join
Used when one table is small.
Steps
Small table loaded into memory
Large table streamed through mapper
Hash join performed
Advantages
Very fast
No shuffle phase
3. Partition Join
Both tables partitioned on join key.
Advantages
Only corresponding partitions joined
Reduces network traffic
4. Bucket Join
Combination of:
Map-side join
Partition join
One relation bucketed by join key.
5. N-Way Map-Side Join
Supports joining multiple tables in one MR job.
Example:
Fact table
Multiple dimension tables
Common in data warehouses.
6. Simple N-Way Join
Multiple smaller tables cached in reducer memory.
Advantages
Efficient for star schema joins
Handling Data Skew
Problem:
Some keys have extremely large number of records.
Solution:
Special partitioning strategies
Unique reducers for skewed keys
25.4.3 Apache Pig
Introduction
Apache Pig is a high-level platform for creating MapReduce programs.
Developed at:
Yahoo Research
Purpose:
Simplify Hadoop programming
Pig Latin Language
Pig Latin is:
Procedural scripting language
Easier than writing raw MapReduce code
Example Query
SQL Query:
SELECT category, AVG(pagerank)
FROM urls
WHERE pagerank > 0.2
GROUP BY category
HAVING COUNT(*) > 10^6;
Equivalent Pig Latin:
good_urls = FILTER urls BY pagerank > 0.2;
groups = GROUP good_urls BY category;
big_groups = FILTER groups BY COUNT(good_urls) > 10^6;
output = FOREACH big_groups GENERATE category,
AVG(good_urls.pagerank);
Features of Pig
1. Dataflow Language
Programs written as sequence of transformations.
2. Supports Complex Data
Supports:
JSON
XML
Nested structures
3. User Defined Functions (UDFs)
Allows custom processing logic.
Pig Data Types
1. Atom
Single value.
Example:
Integer
String
2. Tuple
Collection of fields.
3. Bag
Collection of tuples.
4. Map
Key-value collection.
Applications of Pig
Log analysis
Clickstream analysis
Search engine analytics
Geographic analysis
25.4.4 Apache Hive
Introduction
Apache Hive provides SQL-like querying on Hadoop.
Developed by:
Facebook
Purpose:
Data warehouse processing on Hadoop
Hive Architecture Components
Main components:
1. HiveQL
2. SQL Compiler
3. Metadata Repository
4. JDBC/ODBC Interface
5. HDFS Integration
HiveQL
HiveQL resembles SQL.
Supports:
SELECT
JOIN
GROUP BY
WHERE
Aggregations
Hive and HDFS Integration
Partitioning
Tables partitioned by:
Day
Month
User ID
Benefits:
Faster query execution
Reduced scanning
Bucketing
Data divided into buckets using hashing.
Improves:
Joins
Sampling
Query performance
SerDe (Serialization/Deserialization)
Defines how data is interpreted.
Supported formats:
CSV
JSON
Avro
ORC
Parquet
Optimizations in Hive
Logical Optimizations
Column pruning
Predicate pushdown
Physical Optimizations
Map-side joins
Join reordering
Partition pruning
Advanced SQL Features
Hive supports:
Subqueries
Window functions
Rollups
Grouping sets
Common Table Expressions (CTE)
Hive Runtime Improvements
Hive now supports:
Tez runtime
Spark integration
Advantages:
Faster execution
Reduced disk writes
Better DAG execution
25.4.5 Advantages of Hadoop and MapReduce
1. Massive Scalability
Thousands of machines can work together.
Example:
100 TB scanned using 1000 machines in minutes
2. High Throughput
Large block sizes reduce disk seek overhead.
3. Fault Tolerance
Automatic recovery from failures.
4. Handles Semi-Structured Data
Supports:
Logs
XML
JSON
Key-value data
5. Cost Effective
Uses commodity hardware.
6. Parallel Processing
Large datasets divided into smaller chunks.
7. Linear Scalability
Adding nodes increases performance proportionally.
Comparison: Traditional DBMS vs Hadoop
Traditional
Feature Hadoop
DBMS
Schema Fixed Flexible
Scalability Limited Very High
Fault
Moderate Excellent
Tolerance
Structured +
Data Type Structured
Unstructured
Cost Expensive Low Cost
Processing Centralized Distributed
Key Terms
Term Meaning
Master process managing
JobTracker
jobs
Worker process executing
TaskTracker
tasks
Transfers grouped data to
Shuffle
reducers
Local aggregation before
Combiner
shuffle
Partitioner Assigns keys to reducers
Scripting language for
Pig Latin
Hadoop
HiveQL SQL-like query language
SerDe Serialization/Deserialization
Bucket Hash-based partition
Partition
Scan only relevant partitions
Pruning
Summary
MapReduce enables distributed processing of big data.
JobTracker and TaskTracker manage execution.
Shuffle is central to data grouping.
Hadoop supports fault tolerance and scalability.
Joins can be efficiently implemented using multiple strategies.
Pig simplifies procedural data analysis.
Hive provides SQL-like querying on Hadoop.
Hadoop ecosystem enables large-scale analytics and data
warehousing.
Lecture Notes: Hadoop v2 (YARN)
25.5 Hadoop v2 alias YARN
1. Introduction to Hadoop v2 / YARN
Hadoop v2 is also known as:
MRv2 (MapReduce Version 2)
YARN (Yet Another Resource Negotiator)
YARN was developed to overcome the limitations of Hadoop v1.
Earlier Hadoop architecture tightly coupled:
Resource Management
Job Scheduling
MapReduce Execution
This created scalability and performance problems in large enterprise
clusters.
YARN separates:
1. Cluster Resource Management
2. Application Management
This separation improved:
Scalability
Flexibility
Resource utilization
Fault tolerance
2. Limitations of Hadoop v1
2.1 JobTracker Bottleneck
In Hadoop v1:
JobTracker handled:
o Scheduling
o Resource allocation
o Monitoring
o Failure handling
As cluster size increased:
Memory issues occurred
Heartbeat delays increased
Scheduling became slow
Problem:
JobTracker became:
A bottleneck
Single Point of Failure (SPOF)
2.2 Poor Resource Utilization
Resources were divided into fixed:
Map slots
Reduce slots
Issues:
Reduce slots remained idle during Map phase
Map slots remained idle during Reduce phase
Result:
Cluster underutilization
2.3 Limited to MapReduce
Hadoop v1 mainly supported:
MapReduce applications
But enterprises needed support for:
Machine learning
Graph processing
Interactive analytics
Streaming systems
MapReduce was not suitable for:
Iterative algorithms
Real-time processing
2.4 Upgrade Difficulties
Frequent Hadoop releases caused:
Cluster downtime
Job resubmission
Maintenance complexity
3. Multitenancy Issues in Hadoop v1
Definition
Multitenancy means:
Multiple users/applications sharing the same cluster resources
simultaneously.
Hadoop on Demand (HOD)
Yahoo introduced HOD to create separate clusters for users.
Advantages
Different Hadoop versions
Isolated failures
Custom configurations
Problems
Poor data locality
Static allocation
High latency
Idle resources
Hence, HOD was abandoned.
4. Goals of YARN
YARN was designed with the following goals:
Objectives
1. High scalability
2. Better cluster utilization
3. No single point of failure
4. Support multiple application frameworks
5. Backward compatibility with Hadoop v1
5. YARN Architecture
Core Components
YARN architecture contains:
1. ResourceManager (RM)
2. NodeManager (NM)
3. ApplicationMaster (AM)
4. Containers
6. ResourceManager (RM)
The ResourceManager is the master daemon.
Responsibilities
Resource allocation
Cluster management
Scheduling containers
Monitoring node availability
Important Features
Highly scalable
Only manages resources
Does NOT execute tasks
Interfaces Exposed by RM
1. Client API
Used to start ApplicationMasters.
2. AM-RM Protocol
Used by ApplicationMasters to request resources.
3. NM-RM Protocol
Used by NodeManagers to report node health.
7. Resource Requests
ApplicationMasters request resources using:
ResourceRequest
It specifies:
Number of containers
CPU requirement
Memory requirement
Locality preference
Priority
8. NodeManager (NM)
A NodeManager runs on every worker node.
Responsibilities
Manage containers
Monitor resource usage
Launch applications
Report node health
Services Provided by NM
1. Container Launch Service
Starts processes inside containers.
2. Log Aggregation
Uploads logs to HDFS.
3. Auxiliary Services
Example:
MR Shuffle service
9. ApplicationMaster (AM)
The ApplicationMaster manages one application.
Responsibilities
Request resources
Launch tasks
Monitor execution
Handle failures
Communicate with RM
Example
For a MapReduce job:
AM launches:
o Map tasks
o Reduce tasks
Thus, many JobTracker responsibilities moved to AM.
10. Containers in YARN
A Container is:
A resource allocation unit in YARN.
It contains:
CPU
Memory
Network resources
Applications run inside containers.
11. Workflow of YARN
Step-by-Step Flow
Step 1: Client submits application
Client sends application to RM.
Step 2: RM allocates container
RM launches AM in a container.
Step 3: AM requests resources
AM sends ResourceRequests.
Step 4: RM grants containers
Containers allocated based on scheduler policy.
Step 5: NM launches tasks
NodeManager starts task containers.
Step 6: Tasks execute
Tasks process data.
Step 7: Completion
AM reports completion to RM.
12. Advantages of YARN Architecture
12.1 Better Scalability
RM only manages resources.
Hence:
Scales to thousands of nodes
12.2 Dynamic Resource Allocation
Resources allocated dynamically.
Improves:
Throughput
Utilization
12.3 Multiple Framework Support
YARN supports:
MapReduce
Spark
Tez
Giraph
HBase
12.4 No Global Cluster Downtime
Different application versions can coexist.
12.5 Improved Fault Isolation
Failure of one AM affects only one application.
13. Fault Tolerance in YARN
ResourceManager Failure
RM still remains a critical component.
Recovery
State recovered from persistent storage
Applications restarted
ApplicationMaster Failure
Only affects one application.
AM can:
Restart failed tasks
Recover state
Container Failure
Handled by:
ApplicationMaster
NodeManager
14. YARN Scheduling Policies
Schedulers are pluggable.
Types
1. Fair Scheduler
Equal resource sharing
Better response time for small jobs
2. Capacity Scheduler
Resource guarantees for organizations
Supports multitenancy
15. Other Frameworks on YARN
15.1 Apache Tez
Definition
A DAG-based execution engine on YARN.
Features
Directed Acyclic Graph (DAG) execution
Avoids unnecessary disk writes
Faster than MapReduce
Tez Advantages
Example
MapReduce:
Join → Write to HDFS → Read Again → Group By
Tez:
Join output directly streamed to Group By.
Benefits
Lower latency
Better performance
15.2 Apache Giraph
Definition
Graph processing framework based on Google Pregel.
Features
Bulk Synchronous Parallel (BSP) model
Iterative graph processing
Used for PageRank
Applications
Used at Facebook for:
Social graph analysis
Friend recommendations
15.3 Hoya (HBase on YARN)
Purpose
Run elastic HBase clusters on YARN.
Features
Dynamic scaling
Multiple HBase versions
On-demand HBase clusters
16. Comparison: Hadoop v1 vs Hadoop v2
Hadoop Hadoop v2
Feature
v1 (YARN)
Resource JobTracke
ResourceManager
Management r
JobTracke
Job Management ApplicationMaster
r
Scalability Limited Very High
Fault Tolerance Weak Improved
Static Dynamic
Resource Allocation
Slots Containers
Multiple
Framework Support Only MR
Frameworks
SPOF Issue Severe Reduced
Cluster Utilization Poor Better
17. Key Terms
Term Meaning
Yet Another Resource
YARN
Negotiator
RM ResourceManager
Term Meaning
NM NodeManager
AM ApplicationMaster
Containe
Unit of resource allocation
r
DAG Directed Acyclic Graph
HOD Hadoop on Demand
BSP Bulk Synchronous Parallel
18. Advantages of Hadoop v2
Major Benefits
High scalability
Better utilization
Dynamic scheduling
Multiple frameworks
Improved fault tolerance
Efficient resource sharing
Better support for enterprise workloads
19. Disadvantages of YARN
RM still critical component
Complex architecture
Requires careful tuning
Resource management overhead
20. Important Exam Questions
Short Questions
1. Define YARN.
2. What is ResourceManager?
3. What is ApplicationMaster?
4. What is a container in YARN?
5. What is multitenancy?
6. What are the limitations of Hadoop v1?
7. Define NodeManager.
8. What is Tez?
9. What is Giraph?
10. What is Hoya?
Long Questions
1. Explain YARN architecture with neat diagram.
2. Discuss the limitations of Hadoop v1.
3. Explain ResourceManager and NodeManager.
4. Describe ApplicationMaster in YARN.
5. Explain fault tolerance in YARN.
6. Compare Hadoop v1 and Hadoop v2.
7. Explain Apache Tez and its advantages.
8. Discuss frameworks running on YARN.
21. Summary
YARN transformed Hadoop from:
A MapReduce-only system
into:
A general-purpose distributed data operating platform.
Key improvements include:
Dynamic resource allocation
Better scalability
Support for multiple processing frameworks
Higher cluster utilization
Improved fault tolerance
YARN enabled the Hadoop ecosystem to support:
Big Data analytics
Machine learning
Graph processing
SQL analytics
Lecture Notes: General Discussion on Big Data Technologies
(Section 25.6)
1. Introduction to Big Data Ecosystem
Big Data technologies evolved rapidly between 2004–2014 with major
developments in:
Hadoop v1
YARN (Hadoop v2 / MRv2)
Cloud Computing
Distributed Analytics Platforms
Many organizations and research institutions contributed to this growth:
Hortonworks
Cloudera
MapR Technologies
University of California, Berkeley AMPLab
The section discusses:
Hadoop vs Parallel RDBMS
Big Data and Cloud integration
Data locality challenges
YARN as a data services platform
Current challenges in Big Data
Future direction with Spark and in-memory analytics
25.6.1 Hadoop/MapReduce vs Parallel RDBMS
Comparison Study
Researchers such as:
Daniel Abadi
David DeWitt
Michael Stonebraker
compared:
Parallel Relational DBMS
Hadoop MapReduce
using a 100-node cluster.
Key Observation
Parallel DBMS systems were:
3–6 times faster than Hadoop v1
Reasons for Better RDBMS Performance
1. Indexing Support
B+ Trees accelerate:
o Selection
o Filtering
o Searching
2. Advanced Storage Models
Column-oriented storage improves analytical queries.
3. Compression Techniques
Operations directly on compressed data reduce I/O overhead.
4. Query Optimization
Sophisticated parallel query optimizers improve execution speed.
Hadoop Improvements Over Time
Modern Hadoop systems introduced:
Advanced File Formats
ORC (Optimized Row Columnar)
Parquet
Advantages:
Columnar storage
Compression
Predicate pushdown
Faster aggregation queries
SQL Support Improvements
Hadoop Earlier
Weak SQL support
Current Hadoop Ecosystem
HiveQL supports:
o Nested queries
o Aggregations
o Group By
o Warehouse functions
Additional Platforms
Apache Hive
Apache Spark
Apache Mahout
Apache Giraph
Handling Semi-Structured Data
Hadoop handles:
CSV
JSON
XML
more easily because:
Schema is applied during read time.
RDBMS systems also evolved:
PostgreSQL added JSON and key-value support.
XML support became common.
Fault Tolerance
Hadoop Advantage
MapReduce provides:
Automatic recovery
Replication
Task re-execution
Trade-off
Fault tolerance requires:
Intermediate file materialization
Additional disk I/O
Summary Comparison
Parallel
Feature Hadoop/YARN
RDBMS
Open source, low
Cost Expensive
cost
SQL Support Improving rapidly Mature
Semi-structured
Strong Moderate
Data
Fault Tolerance Excellent Moderate
Traditionally
Performance Improving
faster
25.6.2 Big Data in Cloud Computing
Relationship Between Big Data and Cloud
Big Data heavily depends on cloud computing because cloud provides:
1. Scalability
Two types:
Scale Out → Add more nodes
Scale Up → Add more resources to a node
2. Flexible Resource Management
Cloud infrastructure:
Handles failures gracefully
Supports virtualized computing
Enables elastic expansion
3. Cost Benefits
Cloud enables:
Operational expenditure instead of capital expenditure
Pay-as-you-use model
Ideal for startups and analytics projects.
Common Big Data Cloud Workflow
Typical architecture:
1. Data stored in cloud storage
2. Processing using:
o MapReduce
o Hive
o Pig
Examples:
Amazon S3
Cloud databases
Netflix Case Study
Netflix uses:
Amazon S3
Elastic MapReduce (EMR)
instead of HDFS.
Why Netflix Uses S3
Advantages
High Durability
99.999999999% durability
High Availability
99.99% availability
Elastic Storage
Easily scales from terabytes to petabytes
Cluster Flexibility
Clusters can:
Fail
Restart
Relocate geographically
without data loss.
25.6.3 Data Locality and Resource Optimization
Data Locality Problem
Performance suffers when:
Storage and computation are far apart.
Example:
Data in Amazon S3
Computation in Amazon EC2
causes:
High network latency
Slower execution
MapReduce Network Challenge
In Reduce phase:
Each reducer reads outputs from all mappers.
Large network traffic occurs.
Proposed Solutions
1. Locality-Aware Frameworks
Goal:
Place computation close to data.
Benefits:
Reduced network distance
Faster processing
2. Caching Techniques
PACMan
In-memory caching
MixApart
Disk-based caching
Benefits:
Faster repeated access
Reduced network load
Cloud Resource Optimization
Goal:
Optimize globally across all jobs.
Google BigQuery
Google BigQuery
Features:
SQL-like queries
Handles billions of rows
Automatic resource optimization
25.6.4 YARN as a Data Service Platform
Hadoop v1 Limitation
Earlier:
Hadoop focused mainly on MapReduce.
Now:
YARN supports multiple services.
Hadoop as a Data Lake
HDFS stores:
Historical data
Logs
Clickstreams
Sensor data
Social media data
Processing occurs directly where data resides.
Advantages of YARN Platform
Multiple Services Can Coexist
Examples:
Streaming
Machine Learning
Graph Processing
OLAP
SQL Engines
Applications on YARN
1. Apache Storm
Apache Storm
Used for:
Real-time stream processing
2. SAS on YARN
SAS Institute
Provides:
Statistical analytics on YARN
3. IBM Big SQL
IBM
Shared-nothing SQL engine running on YARN.
4. Actian Analytics
Actian
Provides SQL processing directly in Hadoop.
25.6.5 Challenges Faced by Big Data Technologies
1. Heterogeneity of Information
Data differs in:
Formats
Structures
Semantics
Sources
Challenges:
Data integration
Metadata management
Provenance tracking
2. Privacy and Confidentiality
Sensitive data sources:
Smartphones
GPS logs
Clickstreams
Transaction records
Risks:
Identity exposure
Privacy violations
Confidentiality breaches
Example regulation:
HIPAA
3. Visualization Challenges
Big Data systems generate huge outputs.
Need:
Human-friendly interfaces
Interactive visualization
Outlier detection
Collaborative analysis
4. Incomplete and Inconsistent Data
Problems:
Missing values
Uncertain information
Conflicting records
Especially severe in:
Crowdsourced systems
Sensor-based systems
25.6.6 Moving Forward
Need for Unified Big Data Solutions
Modern analytics combines:
ETL
Machine Learning
Graph Processing
Reporting
Challenge:
Different engines use different programming models.
In-Memory Computing
Modern hardware provides:
Large RAM
Flash storage
This enables:
In-memory analytics
Faster response times
SAP HANA
SAP HANA
Features:
In-memory processing
Columnar storage
Scale-out architecture
Apache Spark
Apache Spark
Developed from:
Berkeley AMPLab
Spark Components
Component Purpose
Spark Core Data processing
Spark SQL SQL analytics
Component Purpose
Graph
GraphX
processing
Machine
MLlib
learning
Spark Stream
Streaming processing
Resilient Distributed Datasets (RDDs)
Central abstraction in Spark.
Characteristics:
Distributed collections
Fault tolerant
Re-creatable from lineage
Can be cached:
o In memory
o On disk
Advantages of Spark
Unified programming model
Fast in-memory computation
Supports diverse analytics workloads
Integrates with Hadoop/YARN
Key Takeaways
Hadoop Evolution
Hadoop evolved from batch processing to a complete data platform.
YARN Importance
Decoupled resource management from applications.
Enabled multiple analytics frameworks.
Cloud Integration
Big Data and cloud computing are tightly connected.
Spark Future
In-memory computing and unified analytics are the future of Big
Data.
Ongoing Challenges
Privacy
Data heterogeneity
Visualization
Resource optimization
Real-time analytics