Skip to main content

πŸ“ Apache Spark

Description​

< What is it? >​

Apache Spark is an open-source engine for processing data in parallel across a cluster. It supports SQL queries, DataFrame transformations, and streaming workloads. PySpark is its Python interface.

Spark reads data from storage, computes results, and writes them back. Storage can be a local filesystem, the Hadoop Distributed File System (HDFS), cloud object storage, or a connected database. See the Spark quick start.

  • Example: read sales records from many files, group them by region, and write regional totals for a reporting system.

Key points​

< Driver, executors, and partitions >​

ComponentRole
DriverRuns the application’s coordinating logic and schedules tasks
ExecutorsRun tasks and hold application data on worker machines
PartitionA portion of a dataset that a task can process
Cluster managerAllocates resources; options include Spark standalone, YARN, and Kubernetes
Application / driver β†’ Tasks β†’ Executors β†’ Read and write storage

A large dataset is divided into partitions so multiple tasks can run concurrently. The driver should coordinate work without collecting the entire dataset into its own memory. See the cluster overview.

< DataFrames and lazy execution >​

A DataFrame represents data as named columns with a schema. Transformations such as filtering and selecting columns build a plan; an action such as count(), collect(), or a write triggers computation. Spark SQL can optimize structured operations expressed through SQL or DataFrame APIs. See the Spark SQL guide.

A shuffle redistributes data between partitions, often for joins or grouped aggregations. Network transfer, sorting, and uneven partition sizes can dominate the cost of a job. Caching can help when intermediate data is reused, but it consumes resources and is not needed for every dataset.

< Batch and streaming >​

Batch processing works on a bounded dataset. Structured Streaming processes incrementally arriving data through Spark's structured APIs. Streaming applications also need checkpointing and appropriate handling of event time, late records, and output writes. See the Structured Streaming documentation.

Comparison​

< Spark and Hadoop MapReduce >​

AspectSparkHadoop MapReduce
ExecutionPlans multi-stage computations with reusable intermediate dataExecutes map, shuffle, and reduce stages; pipelines often chain jobs
InterfacesDataFrames, SQL, and lower-level APIsMapper and reducer programs
Typical workloadsData preparation, iterative analytics, SQL, and streamingLarge batch transformations and aggregations
Relationship to HadoopCan use HDFS for storage and YARN for resource managementA processing component of Hadoop

Spark can replace a MapReduce processing job while continuing to use the same Hadoop storage and cluster resources. Performance depends on the workload, data layout, and configuration.

Implementation​

< Aggregate sales with PySpark >​

This local example requires PySpark and a compatible Java installation for your Spark release. local[2] uses two local worker threads; it does not create a two-machine cluster. Follow the quick start for environment setup.

from pyspark.sql import SparkSession, functions as F

spark = (
SparkSession.builder
.master("local[2]")
.appName("regional-sales")
.getOrCreate()
)

sales = spark.createDataFrame(
[("East", 120), ("West", 90), ("East", 60)],
"region string, amount long",
)
totals = (
sales.groupBy("region")
.agg(F.sum("amount").alias("total"))
.orderBy("region")
)

for row in totals.collect():
print(row["region"], row["total"])

spark.stop()

Expected printed result, excluding Spark startup logs:

East 180
West 90

collect() is suitable here because the result contains two rows. For large results, write distributed output to storage instead of bringing it all into the driver.

Troubleshoot​

< Common problems >​

SymptomWhat to inspect
Driver runs out of memoryLarge collect() or toPandas() calls
One task takes much longer than othersSkewed keys and uneven partition sizes
A join or aggregation is slowShuffle volume, data scanned, and the query plan
A tiny job is slower than local PythonCluster startup and scheduling overhead relative to the work

Reference​