Free Handbook · Every example compiled & verified

Scala + Tools

Scala as it is used at work: sbt builds and MUnit tests, Apache Spark, Kafka, PostgreSQL, Redis, AWS, Docker and CI, each as a small, focused lab.

0 / 142 lessons🔥 0 day streak
ShareXLinkedIn

Module 12 · what you'll be able to do

  • Read and write a build.sbt with dependencies, a Scala version and a test framework, and run it with sbt
  • Write a Spark job with the Dataset API and explain how it maps to the collection methods you already know
  • Produce to and consume from Kafka, and query PostgreSQL with parameters from Scala
  • Cache reads in Redis and call AWS services from Scala through the Java SDK
  • Package a Scala service into a Docker image and run its tests on every push with GitHub Actions
01

The Scala toolchain at a glance

Every program so far has been a single Main.scala run with scala run. A real Scala job is a project built with sbt (or Mill, or scala-cli for small tools), dozens of libraries from Maven Central, a test suite, and usually one of two worlds: data engineering (Spark, Kafka, Databricks, Airflow) or backend services (http4s, Pekko, ZIO, Play, with PostgreSQL). Because Scala runs on the JVM, every Java library is available too, which is how Scala code talks to AWS, JDBC and Redis. These labs show the smallest real version of each. Most need libraries or servers, so they are static snippets; where the idea runs on the standard library there is a verified example as well.

ToolJobYou meet it when
sbt (or Mill, scala-cli)Compile, test, manage dependencies, packageDay one of any Scala job
MUnit / ScalaTestAutomated testsEvery pull request
Apache SparkDistributed batch and streaming data processingMost Scala data-engineering roles
Apache KafkaEvent streams between servicesStreaming pipelines and event-driven backends
PostgreSQL (JDBC, Doobie, Skunk)Store and query dataAny service that keeps data
RedisCache, rate limits, sessionsAs soon as a read gets slow
AWS (S3, EMR, Glue)Storage and managed compute for data jobsRunning Spark in the cloud
Docker + GitHub ActionsPackage and ship; test every pushDeployment and CI
02

Scala + sbt: the build file

sbt is Scala's standard build tool. build.sbt declares the Scala version and the libraries; %% appends the Scala binary version to the artifact name (munit_3), because Scala libraries are published once per Scala version, while a single % is for plain Java libraries. Source goes in src/main/scala, tests in src/test/scala. Start sbt once and keep it running: ~test re-runs the tests on every save.

scalabuild.sbt
ThisBuild / scalaVersion := "3.3.6"   // the Scala 3 LTS line
ThisBuild / organization := "com.example"

lazy val root = (project in file("."))
  .settings(
    name := "orders",
    libraryDependencies ++= Seq(
      "org.postgresql" % "postgresql" % "42.7.7",        // Java library: one %
      "org.scalameta" %% "munit" % "1.1.1" % Test,       // Scala library: %% adds _3
    ),
    scalacOptions ++= Seq("-deprecation", "-feature", "-Werror"),
  )

// project/build.properties:  sbt.version=1.11.6
//
// sbt compile     compile main sources
// sbt test        run every test
// sbt ~test       re-run tests on every file save
// sbt run         run the main method
// sbt assembly    fat jar (needs the sbt-assembly plugin)
Spark projects pin the Scala version
Spark 3.x is built for Scala 2.12 and 2.13, and Spark 4 for Scala 2.13, so a Spark job's build.sbt uses scalaVersion := "2.13.16" and marks Spark as % "provided" (the cluster supplies it). Everything in this handbook compiles on Scala 3; the Scala 2.13 syntax differences you will see in Spark code (braces, implicit instead of given) are covered in Module 00.
03

Scala + MUnit: automated tests

MUnit is the lightweight test library most new Scala 3 projects start with; ScalaTest is the long-standing alternative you will find in older and Spark codebases. A test is a named block with assertions. Pure functions, which Scala encourages, are the easiest code in the world to test: same input, same output, nothing to mock.

scalasrc/test/scala/PricingSuite.scala
class PricingSuite extends munit.FunSuite:

  test("adds tax to the net price"):
    assertEqualsDouble(Pricing.gross(100.0, taxRate = 0.18), 118.0, 0.001)

  test("a bulk discount applies from 10 items"):
    assertEquals(Pricing.discount(qty = 9), 0.0)
    assertEquals(Pricing.discount(qty = 10), 0.05)

  test("rejects a negative price"):
    val e = intercept[IllegalArgumentException](Pricing.gross(-1.0, 0.18))
    assertEquals(e.getMessage, "requirement failed: net must not be negative")

// src/main/scala/Pricing.scala
object Pricing:
  def gross(net: Double, taxRate: Double): Double =
    require(net >= 0, "net must not be negative")
    net * (1 + taxRate)

  def discount(qty: Int): Double = if qty >= 10 then 0.05 else 0.0

The same idea runs without any library: a tiny check helper is enough to verify a pure function while you learn, and it is exactly what an interviewer wants to see you write in a coding round.

scalaMain.scala
object Pricing:
  def gross(net: Double, taxRate: Double): Double =
    require(net >= 0, "net must not be negative")
    net * (1 + taxRate)

  def discount(qty: Int): Double = if qty >= 10 then 0.05 else 0.0

def check(name: String, ok: Boolean): Unit =
  println(s"${if ok then "PASS" else "FAIL"} $name")

@main def run(): Unit =
  check("adds tax", math.abs(Pricing.gross(100.0, 0.18) - 118.0) < 0.001)
  check("no discount at 9", Pricing.discount(9) == 0.0)
  check("discount at 10", Pricing.discount(10) == 0.05)
  val rejected =
    try { Pricing.gross(-1.0, 0.18); false }
    catch case e: IllegalArgumentException => e.getMessage == "requirement failed: net must not be negative"
  check("rejects negative", rejected)
Outputcompiled & run with real Scala
PASS adds tax
PASS no discount at 9
PASS discount at 10
PASS rejects negative
Your turn

Break discount on purpose (change >= to >) and run it again. Exactly one check fails, and its name tells you where.

04

Scala + Apache Spark: a Dataset job

Spark is written in Scala, and Scala is its most complete API. A Dataset[T] is a distributed collection of typed rows: the operations have the same names as the collection methods from Module 06 (map, filter, groupByKey), but they build a plan that Spark optimises and runs across a cluster only when an action (show, count, write) is called. Case classes from Module 05 become the row schema.

Scala + Apache Spark

Revenue per country from order files, with the Dataset and DataFrame APIs

Dependency: "org.apache.spark" %% "spark-sql" % "3.5.6" % "provided" with scalaVersion := "2.13.16". import spark.implicits.* brings the encoders that turn case classes into rows and enables $"column". The typed filter(_.amount > 0) is checked by the compiler; the untyped groupBy and agg are what the optimiser handles best, so real jobs mix both. Run locally with sbt run (master local[*]), or on a cluster with spark-submit. The full engine, joins and tuning are in the Spark with Scala handbook.

scala
import org.apache.spark.sql.{SparkSession, functions as F}

case class Order(id: Long, country: String, amount: Double)

object Revenue {
  def main(args: Array[String]): Unit = {
    val spark = SparkSession.builder()
      .appName("revenue-by-country")
      .master("local[*]")          // drop this line when using spark-submit
      .getOrCreate()
    import spark.implicits._

    val orders = spark.read
      .option("header", "true")
      .option("inferSchema", "true")
      .csv("data/orders.csv")
      .as[Order]                  // typed Dataset[Order]

    val revenue = orders
      .filter(_.amount > 0)       // typed lambda
      .groupBy($"country")
      .agg(F.round(F.sum($"amount"), 2).as("revenue"), F.count("*").as("orders"))
      .orderBy(F.desc("revenue"))

    revenue.show()                // action: the job runs here
    revenue.write.mode("overwrite").parquet("out/revenue")
    spark.stop()
  }
}
Spark with Scala handbook →

The logic of that job is ordinary Scala. Written against a local List, the same filter, group, sum and sort runs without Spark, which is also how you unit-test a transformation before running it on a cluster:

scalaMain.scala
case class Order(id: Long, country: String, amount: Double)

def revenueByCountry(orders: Seq[Order]): Seq[(String, Double, Int)] =
  orders
    .filter(_.amount > 0)
    .groupBy(_.country)
    .map((country, os) => (country, BigDecimal(os.map(_.amount).sum).setScale(2, BigDecimal.RoundingMode.HALF_UP).toDouble, os.size))
    .toSeq
    .sortBy((country, revenue, _) => (-revenue, country))

@main def run(): Unit =
  val orders = List(
    Order(1, "IN", 120.50), Order(2, "US", 99.99), Order(3, "IN", 80.00),
    Order(4, "DE", 45.10), Order(5, "US", -10.00), Order(6, "DE", 60.00),
  )
  for (country, revenue, n) <- revenueByCountry(orders) do
    println(f"$country%-3s $revenue%8.2f  $n order(s)")
Outputcompiled & run with real Scala
IN    200.50  2 order(s)
DE    105.10  2 order(s)
US     99.99  1 order(s)

The refund (a negative amount) is filtered out before grouping, exactly as in the Spark job. Keeping the transformation a pure function of Seq[Order] makes it testable in milliseconds.

05

Scala + Apache Kafka: producing and consuming events

Kafka is a durable, partitioned log of events. Producers append records to a topic; consumers in a consumer group share the topic's partitions and track how far they have read with offsets. From Scala you use the official Java client directly (shown below), or a streaming library built on it: fs2-kafka, ZIO Kafka, Pekko Connectors Kafka, or Spark Structured Streaming's Kafka source.

Scala + Apache Kafka

A producer and a consumer loop with the Java client

Dependency: "org.apache.kafka" % "kafka-clients" % "4.0.0". The key decides the partition, so every event for the same order lands in the same partition and stays in order. enable.auto.commit=false plus commitSync() after processing gives at-least-once delivery: a crash replays the last batch, so the handler must be idempotent. Using.resource closes the client even if the loop throws.

scala
import org.apache.kafka.clients.consumer.KafkaConsumer
import org.apache.kafka.clients.producer.{KafkaProducer, ProducerRecord}
import java.time.Duration
import java.util.Properties
import scala.jdk.CollectionConverters.*
import scala.util.Using

val common = Map("bootstrap.servers" -> "localhost:9092")

def props(m: Map[String, String]): Properties =
  val p = Properties()
  m.foreach((k, v) => p.put(k, v))
  p

@main def produce(): Unit =
  val cfg = common ++ Map(
    "key.serializer"   -> "org.apache.kafka.common.serialization.StringSerializer",
    "value.serializer" -> "org.apache.kafka.common.serialization.StringSerializer",
    "acks"             -> "all",
  )
  Using.resource(KafkaProducer[String, String](props(cfg))) { producer =>
    producer.send(ProducerRecord("orders", "order-42", """{"id":42,"status":"paid"}""")).get()
  }

@main def consume(): Unit =
  val cfg = common ++ Map(
    "group.id"           -> "billing",
    "enable.auto.commit" -> "false",
    "auto.offset.reset"  -> "earliest",
    "key.deserializer"   -> "org.apache.kafka.common.serialization.StringDeserializer",
    "value.deserializer" -> "org.apache.kafka.common.serialization.StringDeserializer",
  )
  Using.resource(KafkaConsumer[String, String](props(cfg))) { consumer =>
    consumer.subscribe(List("orders").asJava)
    while true do
      val records = consumer.poll(Duration.ofMillis(500)).asScala
      records.foreach(r => println(s"${r.key} @ partition ${r.partition} offset ${r.offset}: ${r.value}"))
      if records.nonEmpty then consumer.commitSync()
  }
06

Scala + PostgreSQL: parameterised queries

On the JVM every database is reached through JDBC. Plain JDBC works from Scala, but most teams use a Scala library on top: Doobie (JDBC wrapped in Cats Effect, queries as values), Skunk (a native, non-blocking Postgres driver), Slick or Quill (queries written as collection-like code), or ScalikeJDBC. All of them bind parameters instead of concatenating strings.

Scala + PostgreSQL

A Doobie query mapped to a case class

Dependencies: "org.tpolecat" %% "doobie-core" % "1.0.0-RC10" and "org.tpolecat" %% "doobie-postgres" % "1.0.0-RC10". The sql"..." interpolator turns $domain into a bound ? parameter, so it is injection-safe even though it looks like string interpolation. .query[Customer] maps each row to the case class by column position. Nothing touches the database until the IO runs. For SQL itself, joins, indexes and query plans, see SQL Mastery.

scala
import cats.effect.{IO, IOApp}
import doobie.*
import doobie.implicits.*

case class Customer(id: Long, email: String)

object CustomersApp extends IOApp.Simple:
  val xa = Transactor.fromDriverManager[IO](
    driver = "org.postgresql.Driver",
    url = "jdbc:postgresql://localhost:5432/shop",
    user = "shop",
    password = sys.env.getOrElse("DB_PASSWORD", ""),
    logHandler = None,
  )

  def byDomain(domain: String): ConnectionIO[List[Customer]] =
    sql"SELECT id, email FROM customers WHERE email LIKE ${"%@" + domain} ORDER BY id"
      .query[Customer]
      .to[List]

  def run: IO[Unit] =
    byDomain("example.com").transact(xa).flatMap(cs => IO.println(cs))

// List(Customer(1,[email protected]), Customer(4,[email protected]))
07

Scala + Redis: cache-aside

Cache-aside is the most common caching pattern: look in the cache first; on a miss, load from the slow source, store the result with an expiry, and return it. From Scala you can use a Java client such as Jedis or Lettuce, or redis4cats in Cats Effect codebases.

Scala + Redis

Caching a slow lookup with Jedis and a TTL

Dependency: "redis.clients" % "jedis" % "6.0.0". A JedisPool is shared and thread-safe; each call borrows a connection with Using.resource and returns it. setex stores the value with a 300-second expiry, so stale data fixes itself. The cache is an optimisation: if Redis is empty, the code still returns the right answer.

scala
import redis.clients.jedis.JedisPool
import scala.util.Using

class PriceCache(pool: JedisPool, loadFromDb: String => BigDecimal):
  def price(sku: String): BigDecimal =
    Using.resource(pool.getResource) { redis =>
      Option(redis.get(s"price:$sku")) match
        case Some(cached) => BigDecimal(cached)          // hit
        case None =>                                      // miss: load and store
          val fresh = loadFromDb(sku)
          redis.setex(s"price:$sku", 300, fresh.toString)
          fresh
    }

// val cache = PriceCache(JedisPool("localhost", 6379), sku => slowQuery(sku))
// cache.price("SKU-1")  // first call: database; next 5 minutes: Redis

The pattern itself is plain Scala. Here it is with an in-memory map standing in for Redis, counting how often the slow source is actually hit:

scalaMain.scala
import scala.collection.mutable

class CacheAside[K, V](load: K => V):
  private val store = mutable.Map.empty[K, V]
  var loads = 0

  def get(key: K): V =
    store.getOrElseUpdate(key, { loads += 1; load(key) })

@main def run(): Unit =
  val prices = CacheAside[String, BigDecimal](sku => BigDecimal(sku.length * 10))
  println(prices.get("SKU-1"))
  println(prices.get("SKU-1"))
  println(prices.get("SKU-22"))
  println(s"slow loads: ${prices.loads}")
Outputcompiled & run with real Scala
50
50
60
slow loads: 2

getOrElseUpdate takes its second argument by name (Module 03), so the block only runs on a miss. Three lookups, two loads.

08

Scala + AWS: reading and writing S3

Scala data platforms live in the cloud: raw files land in S3, Spark jobs run on EMR, Glue or Databricks, and results go back to S3 as Parquet or Delta/Iceberg tables. Spark reads s3a:// paths directly; for everything else, Scala calls the AWS SDK for Java v2.

Scala + AWS

Upload and list objects with the AWS SDK for Java v2

Dependency: "software.amazon.awssdk" % "s3" % "2.31.50". Credentials come from the default chain (environment variables, ~/.aws/credentials, or the IAM role of the machine), never from code. listObjectsV2Paginator handles pages of 1,000 keys for you; asScala turns the Java iterable into a Scala one.

scala
import software.amazon.awssdk.core.sync.RequestBody
import software.amazon.awssdk.regions.Region
import software.amazon.awssdk.services.s3.S3Client
import software.amazon.awssdk.services.s3.model.{ListObjectsV2Request, PutObjectRequest}
import scala.jdk.CollectionConverters.*
import scala.util.Using

@main def s3Demo(): Unit =
  Using.resource(S3Client.builder().region(Region.AP_SOUTH_1).build()) { s3 =>
    val bucket = "acme-data-lake"

    s3.putObject(
      PutObjectRequest.builder().bucket(bucket).key("raw/orders/2024-06-01.csv").build(),
      RequestBody.fromString("id,country,amount\n1,IN,120.50\n"),
    )

    val listing = s3.listObjectsV2Paginator(
      ListObjectsV2Request.builder().bucket(bucket).prefix("raw/orders/").build()
    )
    for obj <- listing.contents().asScala do
      println(s"${obj.key} ${obj.size} bytes")
  }
09

Scala + Docker: packaging a service

A Scala service ships as a JVM application in a container. The two common routes are a fat jar built by sbt-assembly and copied into a small JRE image (below), or the sbt-native-packager plugin, whose Docker/publishLocal task writes the Dockerfile for you. Either way, build in one stage and run in a slim one, so the final image carries a JRE and your jar, not sbt and the compiler.

Scala + Docker

A multi-stage Dockerfile for an sbt project

The first stage has sbt and a JDK; it downloads dependencies in a separate layer (copying only the build files first, so the layer is cached until they change), then builds the fat jar. The second stage is a plain JRE. -XX:MaxRAMPercentage makes the JVM size its heap from the container's memory limit rather than the host's. Needs addSbtPlugin("com.eed3si9n" % "sbt-assembly" % "2.3.1") in project/plugins.sbt.

dockerfile
# build stage
FROM sbtscala/scala-sbt:eclipse-temurin-21.0.7_6_1.11.6_3.3.6 AS build
WORKDIR /app
COPY build.sbt ./
COPY project ./project
RUN sbt update                      # dependencies cached in their own layer
COPY src ./src
RUN sbt assembly                    # target/scala-3.3.6/orders-assembly-0.1.0.jar

# run stage
FROM eclipse-temurin:21-jre
WORKDIR /app
COPY --from=build /app/target/scala-3.3.6/*-assembly-*.jar app.jar
USER 1000
EXPOSE 8080
ENTRYPOINT ["java", "-XX:MaxRAMPercentage=75", "-jar", "app.jar"]

# docker build -t orders:1.0 .
# docker run -p 8080:8080 -e DB_PASSWORD=... orders:1.0
10

Scala + GitHub Actions: CI on every push

Continuous integration compiles and tests every push and pull request, so a broken build never reaches the main branch. For Scala the two things that make CI fast are caching the dependency directories (Coursier and Ivy) and running sbt once with several commands instead of starting it several times.

Scala + GitHub

A workflow that checks formatting, compiles with warnings as errors, and tests

actions/setup-java with cache: sbt restores the dependency cache; sbt/setup-sbt installs the launcher. scalafmtCheckAll fails the build if code is not formatted (it needs the sbt-scalafmt plugin and a .scalafmt.conf), and with -Werror in scalacOptions a new compiler warning fails too. Add the docker build step on the main branch to publish the image from the previous lab.

yaml
# .github/workflows/ci.yml
name: ci
on:
  push:
    branches: [main]
  pull_request:

jobs:
  test:
    runs-on: ubuntu-latest
    steps:
      - uses: actions/checkout@v4
      - uses: actions/setup-java@v4
        with:
          distribution: temurin
          java-version: 21
          cache: sbt
      - uses: sbt/setup-sbt@v1
      - name: Format, compile, test
        run: sbt scalafmtCheckAll compile test
sbt
The standard Scala build tool: compiles, tests, resolves dependencies and packages, driven by build.sbt.
%% vs %
In build.sbt, %% appends the Scala binary version to a Scala library's name (munit_3); % is for plain Java libraries.
Dataset[T]
Spark's typed, distributed collection. Transformations build a plan; actions such as show or write run it.
Consumer group
Kafka consumers that share a topic's partitions, each partition read by one member, with progress stored as offsets.
Cache-aside
Check the cache, load from the source on a miss, store with an expiry, return.
Fat jar
One jar containing your code and every dependency, runnable with java -jar; built by sbt-assembly.
Quick check

In build.sbt, why is it "org.scalameta" %% "munit" % "1.1.1" but "org.postgresql" % "postgresql" % "42.7.7"?

Frequently asked questions

What build tool do Scala projects use?
sbt is the standard and the one most jobs use; Mill is a simpler alternative gaining ground, and scala-cli is ideal for scripts and small tools. Spark projects are often built with sbt or Maven. All of them pull libraries from Maven Central.
Do I need Scala to use Apache Spark?
No, Spark also has Python, SQL, Java and R APIs, and PySpark is very popular. But Spark itself is written in Scala, the Scala API is the most complete and type-safe one, and many data-engineering roles, especially on high-volume or streaming pipelines, ask for Scala.
Can Scala use Java libraries?
Yes. Scala compiles to JVM bytecode, so any Java library works directly: JDBC drivers, the Kafka client, the AWS SDK, Jedis. scala.jdk.CollectionConverters converts between Java and Scala collections with asScala and asJava.

Finish the Scala handbook, then get hired

Sit the exam for your certificate, run your resume through the ATS checker, and see the jobs that ask for exactly this.

Check my resume
Found this course useful? Share it.
ShareXLinkedIn

Comments

0

Join the conversation. Sign in to leave a comment — we'd love to hear your thoughts.