Data Engineering › Batch & Distributed Processing
Apache Spark
The most widely used engine for distributed batch and streaming processing.
Also known as: Spark, PySpark, Spark SQL, Spark Structured Streaming, Apache Spark
Apache Spark is the most widely used engine for distributed data processing: it splits a big dataset into partitions, spreads the work across a cluster of machines and combines the results. It handles batch jobs, SQL, streaming and machine learning in one framework, and you can use it from Python (PySpark), Scala, Java, R and SQL.
from pyspark.sql import SparkSession, functions as F
spark = SparkSession.builder.appName("daily-revenue").getOrCreate()
orders = spark.read.parquet("s3://lake/orders/")
daily = (orders
.filter(F.col("status") == "paid")
.groupBy("order_date")
.agg(F.sum("total_cents").alias("revenue")))
daily.write.mode("overwrite").parquet("s3://lake/marts/daily_revenue/")
How it works, in short
- A driver program builds a plan, and executors on worker machines run tasks on partitions of the data in parallel (driver and executors, partitions).
- Operations are lazy: transformations build a plan, and an action triggers execution, so Spark can optimize the whole thing first (transformations vs actions).
- The DataFrame / SQL API is the main interface (DataFrames), optimized by Spark’s query planner.
- Steps that regroup data by key (joins, group-bys, sorts) cause a shuffle: data moves across the network, which is usually the most expensive part.
- It can spill to disk when memory is short (spill to disk).
Where it fits
- Large-scale batch transformations over data lakes (batch processing).
- Streaming with Structured Streaming (stream processing).
- ML pipelines on big data, and a common engine behind lakehouse platforms.
Practical lessons
- You may not need it. For data that fits on one machine, a single-node engine or plain SQL in your warehouse is simpler and often faster. Distribution has overhead.
- Performance problems usually come from shuffles, skew and small files: a few hot keys that overload one task (data skew), too many or too few partitions, or reading thousands of tiny files (small files).
- Prefer built-in functions to Python UDFs, which are slow (UDFs).
- Avoid
collect()on big data, which pulls everything to the driver and crashes it. - Filter and select columns early, and use columnar formats (Parquet).
- Learn to read the Spark UI (stages, tasks, shuffle sizes). It tells you where time goes.
- Memory tuning and cluster sizing are real work. Managed platforms handle some of it.
It’s a big ecosystem. Start with DataFrames and SQL, and the rest makes more sense.