Lesson 21 / 25

Scala for Data: Apache Spark

Process large datasets with Spark's Scala API.

Big data with Spark

Apache Spark is a distributed engine for large-scale data processing, written mostly in Scala, with APIs for Scala, Java, Python and SQL. A SparkSession is the entry point. DataFrames (Dataset[Row]) represent tables with schemas and are optimised by the Catalyst optimiser; Datasets (Dataset[Order]) add compile-time types using case classes. Transformations such as filter, select, groupBy, agg and join are lazy, building a plan that runs only when an action (show, count, write, collect) is called. Spark reads and writes Parquet, CSV, JSON, Delta Lake and databases, and runs on clusters (Kubernetes, YARN, standalone, or managed platforms such as Databricks, EMR and Dataproc). Important practices: avoid collect on large data, prefer built-in functions to UDFs, partition data sensibly and watch for skewed joins. Version note: Spark has historically been built for Scala 2.12 and 2.13, so Spark jobs are usually written in Scala 2.13 syntax, even when the rest of a team uses Scala 3; check the Scala version your Spark distribution supports.

Revenue by city with the DataFrame API

Lazy transformations, an aggregation and a Parquet output.

import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions._

object RevenueJob {
  def main(args: Array[String]): Unit = {
    val spark = SparkSession.builder().appName("revenue-by-city").getOrCreate()

    val orders = spark.read
      .option("header", "true")
      .option("inferSchema", "true")
      .csv("s3a://shop-data/orders/2026-09/*.csv")

    val revenue = orders
      .filter(col("status") === "paid")
      .groupBy(col("city"))
      .agg(
        count(lit(1)).as("orders"),
        round(sum(col("total_paise")) / 100, 2).as("revenue_rupees")
      )
      .orderBy(desc("revenue_rupees"))         // still lazy: nothing has run yet

    revenue.show(10)                            // action: triggers the job
    revenue.write.mode("overwrite").parquet("s3a://shop-data/reports/revenue-by-city/")

    spark.stop()
  }
}
// written in Scala 2.13-compatible syntax, as Spark jobs typically are

Avoid collect on big data

collect() pulls the whole dataset into the driver's memory and can crash it. Use show, take(n) or write the result to storage, and keep the heavy lifting on the executors.

Quick check: When do Spark transformations such as filter and groupBy actually execute?

  • Immediately when called
  • At compile time
  • Only when an action such as show, count or write is called
  • Never; Spark only stores plans
Answer

Only when an action such as show, count or write is called — Spark builds a lazy plan and executes it when an action needs the result.