Big Data and Data Engineering
Hadoop Ecosystem; Apache Hive; Apache Spark
C-CAT
Hadoop Ecosystem
What is Hadoop?
Hadoop is an open-source framework for storing and processing Big Data in a distributed manner.
Core Components:
+--------------------------------------------------+
| HADOOP ECOSYSTEM |
| |
| HDFS → Distributed Storage |
| YARN → Resource Management |
| MapReduce → Distributed Processing |
+--------------------------------------------------+
YARN — Yet Another Resource Negotiator
- Manages cluster resources (CPU, memory) across jobs
- ResourceManager — cluster-level master; allocates resources
- NodeManager — node-level; executes tasks, manages containers
MapReduce
MapReduce is Hadoop's processing model — splits computation into two phases:
Input Data → MAP → (Intermediate key-value pairs) → SHUFFLE & SORT → REDUCE → Output
Word Count Example:
Input: "Hello World Hello"
MAP phase:
"Hello" → 1
"World" → 1
"Hello" → 1
SHUFFLE (group by key):
"Hello" → [1, 1]
"World" → [1]
REDUCE phase:
"Hello" → 2
"World" → 1
MapReduce Code Concept:
# Mapper
def map(text_line):
for word in text_line.split():
emit(word, 1)
# Reducer
def reduce(word, counts):
emit(word, sum(counts))
Hadoop Ecosystem Components
| Component | Function |
|---|---|
| HDFS | Distributed file storage |
| YARN | Resource management and scheduling |
| MapReduce | Batch processing framework |
| Hive | SQL-like interface for HDFS data |
| HBase | NoSQL database on HDFS |
| Pig | Scripting language for data transformation |
| Sqoop | Import/export between RDBMS and HDFS |
| Flume | Ingestion of streaming log data into HDFS |
| Oozie | Workflow scheduler for Hadoop jobs |
| ZooKeeper | Distributed coordination service |
Apache Hive
What is Hive?
Apache Hive is a data warehouse built on top of Hadoop that provides:
- HiveQL — SQL-like query language for querying HDFS data
- Converts HiveQL to MapReduce/Tez/Spark jobs
HiveQL Query → Hive → MapReduce job → HDFS result
-- Create external table pointing to HDFS data
CREATE EXTERNAL TABLE sales (
date STRING,
product STRING,
amount DECIMAL(10,2),
quantity INT
)
ROW FORMAT DELIMITED
FIELDS TERMINATED BY ','
LOCATION '/user/hive/warehouse/sales/';
-- Query the data
SELECT product, SUM(amount) as total_sales
FROM sales
WHERE date BETWEEN '2024-01-01' AND '2024-12-31'
GROUP BY product
ORDER BY total_sales DESC;
Hive vs RDBMS
| Feature | Hive | RDBMS |
|---|---|---|
| Schema | Schema on read | Schema on write |
| Data stored in | HDFS | Local storage |
| Latency | High (minutes) | Low (ms) |
| ACID | Limited | Full |
| Scalability | Massive | Limited |
| Updates | Limited | Full |
| Best for | Analytics (OLAP) | Transactions (OLTP) |
Apache Spark
What is Spark?
Apache Spark is a fast, general-purpose distributed computing framework that:
- Is 100x faster than Hadoop MapReduce (in-memory processing)
- Supports batch, streaming, SQL, ML and graph processing
- Written in Scala; APIs in Python (PySpark), Java, Scala, R, SQL
Spark Architecture
+--------------------------------------------------+
| SPARK CLUSTER |
| |
| Driver Program Worker Nodes |
| +----------+ +---------+ |
| | SparkContext| ──────> | Executor| → Tasks |
| | (your app) | +---------+ |
| +----------+ +---------+ |
| ──> | Executor| → Tasks |
| Cluster Manager +---------+ |
| (YARN, K8s, Mesos) ...N workers... |
+--------------------------------------------------+
RDD (Resilient Distributed Dataset)
The fundamental data structure in Spark:
- Resilient — fault-tolerant (recomputed from lineage if lost)
- Distributed — partitioned across cluster
- Dataset — collection of data elements
Spark Transformations (lazy): map, filter, flatMap, groupByKey, reduceByKey, join
Spark Actions (trigger execution): collect, count, save, reduce, take
PySpark Examples
from pyspark.sql import SparkSession
# Create Spark session
spark = SparkSession.builder \
.appName("Sales Analysis") \
.getOrCreate()
# Read data
df = spark.read.csv("hdfs:///data/sales.csv", header=True, inferSchema=True)
# DataFrame operations (lazy)
result = df \
.filter(df.amount > 1000) \
.groupBy("product") \
.agg({"amount": "sum"}) \
.orderBy("sum(amount)", ascending=False) \
.limit(10)
# Action: actually execute the computation
result.show()
# SQL style
df.createOrReplaceTempView("sales")
spark.sql("""
SELECT product, SUM(amount) as total_sales
FROM sales
WHERE amount > 1000
GROUP BY product
ORDER BY total_sales DESC
LIMIT 10
""").show()
# Spark Streaming
from pyspark.streaming import StreamingContext
ssc = StreamingContext(spark.sparkContext, 1) # 1 second batch interval
Continue learning
Related notes
Definition of AI; Need of AI
Artificial Intelligence
Introduction to Data Engineering; Big Data — The 5 V's; Types of Data
Big Data and Data Engineering
Introduction to C Programming; C Program Structure; Data Types and Variables
C Programming
What Is a Computer?; Machine Cycle: Fetch–Decode–Execute; CPU Organization
Computer Architecture
Put this topic into timed practice
Open mock tests when you want full-exam pacing, or keep drilling in practice mode.