π 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 >β
| Component | Role |
|---|---|
| Driver | Runs the applicationβs coordinating logic and schedules tasks |
| Executors | Run tasks and hold application data on worker machines |
| Partition | A portion of a dataset that a task can process |
| Cluster manager | Allocates 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 >β
| Aspect | Spark | Hadoop MapReduce |
|---|---|---|
| Execution | Plans multi-stage computations with reusable intermediate data | Executes map, shuffle, and reduce stages; pipelines often chain jobs |
| Interfaces | DataFrames, SQL, and lower-level APIs | Mapper and reducer programs |
| Typical workloads | Data preparation, iterative analytics, SQL, and streaming | Large batch transformations and aggregations |
| Relationship to Hadoop | Can use HDFS for storage and YARN for resource management | A 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 >β
| Symptom | What to inspect |
|---|---|
| Driver runs out of memory | Large collect() or toPandas() calls |
| One task takes much longer than others | Skewed keys and uneven partition sizes |
| A join or aggregation is slow | Shuffle volume, data scanned, and the query plan |
| A tiny job is slower than local Python | Cluster startup and scheduling overhead relative to the work |
Related ideasβ
- Apache Hadoop provides HDFS and YARN, which Spark can use.
- MapReduce explains distributed mapping, shuffling, and aggregation.
- Snowflake provides a managed platform for storing and querying data.
- Distributed Systems & Data Processing compares their roles.