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

ComponentFunction
HDFSDistributed file storage
YARNResource management and scheduling
MapReduceBatch processing framework
HiveSQL-like interface for HDFS data
HBaseNoSQL database on HDFS
PigScripting language for data transformation
SqoopImport/export between RDBMS and HDFS
FlumeIngestion of streaming log data into HDFS
OozieWorkflow scheduler for Hadoop jobs
ZooKeeperDistributed 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

FeatureHiveRDBMS
SchemaSchema on readSchema on write
Data stored inHDFSLocal storage
LatencyHigh (minutes)Low (ms)
ACIDLimitedFull
ScalabilityMassiveLimited
UpdatesLimitedFull
Best forAnalytics (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

Put this topic into timed practice

Open mock tests when you want full-exam pacing, or keep drilling in practice mode.