Big Data Technologies

Spark Architecture, RDDs, DataFrames, Spark SQL and Machine Learning

PGCP-BDA

Apache Spark

Apache Spark is a distributed execution engine that builds task graphs for batch, SQL, machine-learning and streaming workloads across a cluster.

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

driver executor and cluster manager

The Spark driver plans work, the cluster manager allocates resources and executors run tasks and retain application data.

RDD

A resilient distributed dataset is an immutable partitioned collection described by lineage so lost partitions can be recomputed.

lineage and lazy evaluation

Spark records transformations as lineage and delays execution until an action requires a result.

transformation and action

A Spark transformation lazily defines a new dataset, while an action triggers execution and returns or writes a result.

narrow and wide transformation

A narrow transformation reads each output partition from few input partitions; a wide transformation requires a shuffle across partitions.

DataFrame and Catalyst

A Spark DataFrame is a distributed table with named typed columns; Catalyst analyzes and optimizes its logical plan before physical execution.

Spark SQL

Spark SQL analyzes SQL or DataFrame expressions into optimized logical and physical plans executed as distributed stages.

partition and cache

Partitions define Spark task parallelism and data placement, while caching retains reused partitions to avoid recomputation.

Spark ML pipeline

A Spark ML pipeline chains estimators and transformers so fitting produces a reproducible model pipeline that applies identical feature steps during.

Spark Execution Model

A Spark application contains a driver and executors. The driver builds a logical plan and schedules jobs. Executors run tasks and retain cached partitions. An action triggers execution of lazy transformations. Narrow transformations can process each output partition from one input partition while wide transformations require a shuffle across the cluster. Stages are separated at shuffle boundaries.

An RDD is an immutable partitioned collection with lineage used for recomputation. A DataFrame adds named columns and schema. Spark SQL can optimize DataFrame and SQL expressions through logical and physical planning. DataFrames normally outperform opaque row-by-row functions because the engine can prune columns, push filters and choose join strategies. Caching helps only when reused data justifies memory and recomputation cost.

Spark SQL and Machine Learning

Spark SQL reads tables and file formats through a common structured interface. Temporary views belong to a session or application while catalog tables preserve metadata. Partitioning and bucketing affect file layout. Broadcast joins avoid shuffling a small relation. Skewed keys and many tiny partitions produce slow tasks even when aggregate capacity is high.

Spark ML pipelines connect transformers and estimators. A transformer converts one DataFrame to another. An estimator learns a model through fit and produces a transformer. Feature preparation, training and evaluation must use separate data to avoid leakage. Parameters and fitted stages can be assembled into a repeatable pipeline. Distributed execution helps large feature sets but does not correct poor sampling, invalid labels or an unsuitable metric.

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.