Contents

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.