Free Handbook · Every example compiled & verified

Futures & Concurrency

Run work concurrently with Future and an ExecutionContext, compose it with map, flatMap and for, handle failures, and share state safely with atomics.

0 / 142 lessons🔥 0 day streak
ShareXLinkedIn

Module 10 · what you'll be able to do

  • Start work on another thread with Future and explain what the ExecutionContext is for
  • Combine Futures with map, flatMap and for, and start independent ones before the for so they run in parallel
  • Handle a failed Future with recover, recoverWith and fallbackTo, and read the exception it carries
  • Turn a list of Futures into one with Future.sequence and Future.traverse, and bridge callback APIs with Promise
  • Keep blocking calls off the global pool and update shared counters safely with java.util.concurrent atomics
01

Future and ExecutionContext

A Future[A] is a value of type A that will exist later. Writing Future { work } hands work to a thread pool and returns immediately, so the current thread can keep going. The thread pool is chosen by an ExecutionContext — a given (implicit) value every Future method asks for. For examples and small programs, import scala.concurrent.ExecutionContext.Implicits.global supplies the standard one: a pool with roughly one thread per CPU core.

scalaMain.scala
import scala.concurrent.{Future, Await}
import scala.concurrent.ExecutionContext.Implicits.global
import scala.concurrent.duration.*

def slowSquare(n: Int): Int =
  Thread.sleep(200)
  n * n

@main def run(): Unit =
  val f: Future[Int] = Future(slowSquare(7))
  println("future started, main keeps going")
  println(s"done yet? ${f.isCompleted}")

  val result = Await.result(f, 2.seconds)
  println(s"result: $result")
  println(s"done yet? ${f.isCompleted}")
  println(f.value)
Outputcompiled & run with real Scala
future started, main keeps going
done yet? false
result: 49
done yet? true
Some(Success(49))

Future(slowSquare(7)) returned at once; the square was computed on a pool thread while main printed. Await.result blocks until the value arrives or the timeout expires. f.value is an Option[Try[Int]]: None while running, then Some(Success(...)) or Some(Failure(...)).

Your turn

Start two Futures, slowSquare(3) and slowSquare(4), then await both and print their sum.

Await only at the edge
Await.result parks a whole thread doing nothing. In a real service you never call it in the middle of your code: you return the Future and let the web framework or main wait once, at the very edge. Every example on this page ends with one Await because that is the simplest way to get a deterministic printed result.
Error you will hit

Cannot find an implicit ExecutionContext

scala
import scala.concurrent.{Future, Await}
import scala.concurrent.duration.*

@main def run(): Unit =
  val f = Future(21 * 2)
  println(Await.result(f, 2.seconds))
-- [E172] Type Error: Main.scala:5:24
5 |  val f = Future(21 * 2)
  |                        ^
  |Cannot find an implicit ExecutionContext. You might add
  |an (implicit ec: ExecutionContext) parameter to your method.
  |
  |The ExecutionContext is used to configure how and on which
  |thread pools asynchronous tasks (such as Futures) will run,
  |so the specific ExecutionContext that is selected is important.
  |
  |If your application does not define an ExecutionContext elsewhere,
  |consider using Scala's global ExecutionContext by defining
  |the following:
  |
  |implicit val ec: scala.concurrent.ExecutionContext = scala.concurrent.ExecutionContext.global
  |
  |The following import might fix the problem:
  |
  |  import scala.concurrent.ExecutionContext.Implicits.global
  |
1 error found
Compilation failed
Why the compiler said that

Future.apply has a using parameter of type ExecutionContext — it needs to know which thread pool should run the block. Nothing in scope provides one, so the compiler stops. map, flatMap, recover and onComplete ask for one too.

The fix

In a small program, import the global pool. In a library or service, take it as a parameter ((using ec: ExecutionContext)) so the caller decides which pool runs your work.

scala
import scala.concurrent.{Future, Await}
import scala.concurrent.ExecutionContext.Implicits.global
import scala.concurrent.duration.*

@main def run(): Unit =
  val f = Future(21 * 2)
  println(Await.result(f, 2.seconds))
Error you will hit

TimeoutException from Await.result

scala
import scala.concurrent.{Future, Await}
import scala.concurrent.ExecutionContext.Implicits.global
import scala.concurrent.duration.*

@main def run(): Unit =
  val slow = Future { Thread.sleep(2000); "done" }
  println(Await.result(slow, 500.millis))
Exception in thread "main" java.util.concurrent.TimeoutException: Future timed out after [500 milliseconds]
	at scala.concurrent.Future$.timeoutError(Future.scala:612)
	at scala.concurrent.impl.Promise$DefaultPromise.tryAwait0(Promise.scala:272)
	at scala.concurrent.impl.Promise$DefaultPromise.result(Promise.scala:285)
	at scala.concurrent.Await$.result$$anonfun$1(package.scala:216)
	at scala.concurrent.BlockContext$DefaultBlockContext$.blockOn(BlockContext.scala:67)
	at scala.concurrent.Await$.result(package.scala:133)
	at Main$package$.run(Main.scala:7)
	at run.main(Main.scala:5)
Why the compiler said that

The work takes two seconds but the caller only agreed to wait half a second. Await.result gives up and throws TimeoutException on the waiting thread. The Future itself is not cancelled — it keeps running in the background; Scala Futures have no cancel button.

The fix

Pick a timeout that matches what the work really needs, or decide what the caller should see when the work is too slow. For a real deadline on a request, handle the timeout as a value instead of letting it crash main.

scala
import scala.concurrent.{Future, Await}
import scala.concurrent.ExecutionContext.Implicits.global
import scala.concurrent.duration.*
import scala.util.Try

@main def run(): Unit =
  val slow = Future { Thread.sleep(2000); "done" }
  val answer = Try(Await.result(slow, 500.millis)).getOrElse("still working, try later")
  println(answer)
02

Composing Futures with map, flatMap and for

You rarely wait for a Future just to use its value. Instead you describe what should happen when it arrives. map transforms the eventual value; flatMap chains a step that itself returns a Future; a for-comprehension is the readable way to write several flatMaps (exactly as with Option in Module 08). Nothing blocks: each step is scheduled to run after the previous one completes.

scalaMain.scala
import scala.concurrent.{Future, Await}
import scala.concurrent.ExecutionContext.Implicits.global
import scala.concurrent.duration.*

case class User(id: Int, name: String)
case class Order(userId: Int, total: Int)

def fetchUser(id: Int): Future[User] = Future(User(id, "Asha"))
def fetchOrders(u: User): Future[List[Order]] =
  Future(List(Order(u.id, 250), Order(u.id, 120)))

@main def run(): Unit =
  val nameF: Future[String] = fetchUser(1).map(_.name.toUpperCase)
  val ordersF: Future[List[Order]] = fetchUser(1).flatMap(fetchOrders)

  val summaryF: Future[String] =
    for
      user   <- fetchUser(1)
      orders <- fetchOrders(user)
    yield s"${user.name} has ${orders.size} orders worth ${orders.map(_.total).sum}"

  println(Await.result(nameF, 2.seconds))
  println(Await.result(ordersF, 2.seconds))
  println(Await.result(summaryF, 2.seconds))
Outputcompiled & run with real Scala
ASHA
List(Order(1,250), Order(1,120))
Asha has 2 orders worth 370

Each of nameF, ordersF and summaryF is itself a Future — you are building a pipeline, not running it step by step. The for is sequential by nature here: the orders call needs the user first.

Your turn

Add def fetchDiscount(u: User): Future[Int] = Future(10) and a third generator in the for, then subtract the discount from the total.

Error you will hit

map where you needed flatMap: Future[Future[Int]]

scala
import scala.concurrent.Future
import scala.concurrent.ExecutionContext.Implicits.global

def fetchId(name: String): Future[Int] = Future(name.length)
def fetchScore(id: Int): Future[Int] = Future(id * 10)

@main def run(): Unit =
  val score: Future[Int] = fetchId("asha").map(id => fetchScore(id))
  println(score)
-- [E007] Type Mismatch Error: Main.scala:8:63
8 |  val score: Future[Int] = fetchId("asha").map(id => fetchScore(id))
  |                                                     ^^^^^^^^^^^^^^
  |                                    Found:    scala.concurrent.Future[Int]
  |                                    Required: Int
1 error found
Compilation failed
Why the compiler said that

map expects a function that returns a plain value. Here the function returns another Future, so the result would be a Future[Future[Int]] — a box inside a box — and it does not fit the declared Future[Int]. Without the type annotation it would compile and you would be left holding the nested type.

The fix

When the next step returns a Future, use flatMap (or a for), which chains and flattens.

scala
import scala.concurrent.{Future, Await}
import scala.concurrent.ExecutionContext.Implicits.global
import scala.concurrent.duration.*

def fetchId(name: String): Future[Int] = Future(name.length)
def fetchScore(id: Int): Future[Int] = Future(id * 10)

@main def run(): Unit =
  val score: Future[Int] = fetchId("asha").flatMap(id => fetchScore(id))
  println(Await.result(score, 2.seconds))
03

Parallel vs sequential: start Futures before the for

This is the Future mistake that ships to production most often. A Future starts running the moment it is created. Inside a for, the second generator is only evaluated after the first Future completes — so if you create both Futures inside the for, they run one after the other. When the steps are independent, create the Futures first, into vals, and only combine them in the for.

scalaMain.scala
import scala.concurrent.{Future, Await}
import scala.concurrent.ExecutionContext.Implicits.global
import scala.concurrent.duration.*

def call(name: String): Future[String] = Future { Thread.sleep(500); name }

def timeMs[A](body: => A): (A, Long) =
  val start = System.nanoTime()
  val a = body
  (a, (System.nanoTime() - start) / 1000000)

@main def run(): Unit =
  val (seq, seqMs) = timeMs {
    val f = for
      a <- call("prices")
      b <- call("stock")
    yield s"$a+$b"
    Await.result(f, 5.seconds)
  }
  println(s"sequential: $seq, at least 1000 ms: ${seqMs >= 1000}")

  val (par, parMs) = timeMs {
    val pricesF = call("prices")
    val stockF  = call("stock")
    val f = for
      a <- pricesF
      b <- stockF
    yield s"$a+$b"
    Await.result(f, 5.seconds)
  }
  println(s"parallel:   $par, under 1000 ms: ${parMs < 1000}")
Outputcompiled & run with real Scala
sequential: prices+stock, at least 1000 ms: true
parallel:   prices+stock, under 1000 ms: true

Same result, half the time. The only difference is where the Futures are created. Timings vary from machine to machine, so the program prints a comparison rather than the raw milliseconds.

Your turn

Add a third 500 ms call, call("reviews"), to both versions. How long does each take now?

VisualizeWhy the first version is sequentialStep 1 / 4
val f = for
a <- call("prices")
b <- call("stock")
yield s"$a+$b"
// the compiler rewrites it to:
call("prices").flatMap(a =>
call("stock").map(b => s"$a+$b"))
Line 7

call("prices") runs immediately and starts a 500 ms task. flatMap registers a callback to run when it finishes.

Variables now
runningprices
All 4 steps as a table
StepLineWhat happenedVariables now
17call("prices") runs immediately and starts a 500 ms task. flatMap registers a callback to run when it finishes.running = prices
28This line is inside the callback. call("stock") is not evaluated until prices has completed, 500 ms later.running = stock a = "prices"
38After another 500 ms, map builds the string. Total: about 1000 ms.a = "prices" b = "stock"
42In the parallel version both calls ran when the vals were defined, before the for. The for then only waits for two tasks that were already running side by side: about 500 ms.

Two independent Futures can also be combined with zip (a pair) or zipWith (a function), which reads as "both of these, together" and avoids the trap completely: pricesF.zipWith(stockF)(_ + "+" + _).

04

Failed Futures: recover, recoverWith, fallbackTo

If the code inside a Future throws, the Future completes as a failure holding that exception — exactly like Try. A failure then skips every later map and flatMap, the way a Left skips steps in an Either. recover turns selected exceptions into a value, recoverWith turns them into another Future (a retry, a cache lookup), and fallbackTo uses a second Future if the first fails.

scalaMain.scala
import scala.concurrent.{Future, Await}
import scala.concurrent.ExecutionContext.Implicits.global
import scala.concurrent.duration.*

def fetchRate(currency: String): Future[Double] = Future {
  currency match
    case "EUR" => 0.92
    case "INR" => 83.1
    case other => throw new NoSuchElementException(s"no rate for $other")
}

def cachedRate(currency: String): Future[Double] = Future(1.0)

@main def run(): Unit =
  val bad = fetchRate("XYZ").map(_ * 100)
  Await.ready(bad, 2.seconds)
  println(bad.value)

  val safe = fetchRate("XYZ").recover { case _: NoSuchElementException => 0.0 }
  println(Await.result(safe, 2.seconds))

  val cached = fetchRate("XYZ").recoverWith { case _ => cachedRate("XYZ") }
  println(Await.result(cached, 2.seconds))

  val either = fetchRate("GBP").fallbackTo(fetchRate("EUR"))
  println(Await.result(either, 2.seconds))

  val asEither = fetchRate("XYZ").map(Right(_)).recover { case e => Left(e.getMessage) }
  println(Await.result(asEither, 2.seconds))
Outputcompiled & run with real Scala
Some(Failure(java.util.NoSuchElementException: no rate for XYZ))
0.0
1.0
0.92
Left(no rate for XYZ)

Await.ready waits without throwing, so you can inspect the failure in value. The last line is a common pattern: turn a failed Future into a successful Future[Either[String, Double]] so the error travels as data.

Your turn

Use recover with two cases: NoSuchElementException gives 1.0, and anything else gives -1.0.

Callbacks and a program that exits too early
f.onComplete { ... } and f.foreach(println) register code to run when the Future finishes, on a pool thread. Nothing waits for them: if main returns first, the JVM exits and the callback never runs, because the global pool uses daemon threads. Use callbacks for side effects such as logging and metrics, and keep composing with map/flatMap for results.
Error you will hit

A failed Future rethrown by Await.result

scala
import scala.concurrent.{Future, Await}
import scala.concurrent.ExecutionContext.Implicits.global
import scala.concurrent.duration.*

def fetchPrice(sku: String): Future[Int] = Future {
  if sku == "B-2" then throw new IllegalStateException(s"price service has no $sku")
  100
}

@main def run(): Unit =
  val total = fetchPrice("B-2").map(_ * 2)
  println(Await.result(total, 2.seconds))
Exception in thread "main" java.lang.IllegalStateException: price service has no B-2
	at Main$package$.fetchPrice$$anonfun$1(Main.scala:6)
	at scala.concurrent.Future$.apply$$anonfun$1(Future.scala:722)
	at scala.concurrent.impl.Promise$Transformation.run(Promise.scala:522)
	at java.base/java.util.concurrent.ForkJoinTask$RunnableExecuteAction.compute(Unknown Source)
	at java.base/java.util.concurrent.ForkJoinTask$RunnableExecuteAction.compute(Unknown Source)
	at java.base/java.util.concurrent.ForkJoinTask$InterruptibleTask.exec(Unknown Source)
	at java.base/java.util.concurrent.ForkJoinTask.doExec(Unknown Source)
	at java.base/java.util.concurrent.ForkJoinPool$WorkQueue.topLevelExec(Unknown Source)
	at java.base/java.util.concurrent.ForkJoinPool.runWorker(Unknown Source)
	at java.base/java.util.concurrent.ForkJoinWorkerThread.run(Unknown Source)
Why the compiler said that

The exception was thrown on a pool thread (line 6), stored in the Future, and passed untouched through map. Await.result rethrows it on the main thread. Notice the stack trace: it shows where the exception was created — a ForkJoin worker running your lambda — and never mentions line 12 where main was waiting. Stack traces from async code show the worker, not the caller.

The fix

Decide what a failure means for the caller and say so with recover (or turn it into an Either) before the Future reaches the edge. When debugging, put the input in the exception message, as this one does, because the stack trace will not tell you who asked.

scala
import scala.concurrent.{Future, Await}
import scala.concurrent.ExecutionContext.Implicits.global
import scala.concurrent.duration.*

def fetchPrice(sku: String): Future[Int] = Future {
  if sku == "B-2" then throw new IllegalStateException(s"price service has no $sku")
  100
}

@main def run(): Unit =
  val total = fetchPrice("B-2").map(_ * 2).recover { case _: IllegalStateException => 0 }
  println(Await.result(total, 2.seconds))
05

Many Futures at once: sequence and traverse

A list of ids often becomes a list of calls. Future.sequence turns a List[Future[A]] into a single Future[List[A]]. Future.traverse(items)(f) does the map and the sequence in one step. All the Futures run concurrently, and the result keeps the input order, no matter which call finished first. If any one fails, the combined Future fails.

scalaMain.scala
import scala.concurrent.{Future, Await}
import scala.concurrent.ExecutionContext.Implicits.global
import scala.concurrent.duration.*
import scala.util.{Try, Success}

def lookup(id: Int): Future[String] = Future {
  Thread.sleep(400 - id * 100)   // later ids finish first
  s"item-$id"
}

@main def run(): Unit =
  val one: List[Future[String]] = List(1, 2, 3).map(lookup)
  println(Await.result(Future.sequence(one), 2.seconds))

  println(Await.result(Future.traverse(List(1, 2, 3))(lookup), 2.seconds))

  val parsed = Future.traverse(List("1", "2", "x"))(s => Future(s.toInt))
  Await.ready(parsed, 2.seconds)
  println(parsed.value)

  def attempt(s: String): Future[Try[Int]] = Future(s.toInt).transform(t => Success(t))
  val each = Future.traverse(List("1", "2", "x"))(attempt)
  println(Await.result(each, 2.seconds).map(_.toOption))
Outputcompiled & run with real Scala
List(item-1, item-2, item-3)
List(item-1, item-2, item-3)
Some(Failure(java.lang.NumberFormatException: For input string: "x"))
List(Some(1), Some(2), None)

Item 3 finished first, yet the list is in input order. transform(t => Success(t)) turns every outcome — success or failure — into a successful Future[Try[Int]], so one bad item no longer fails the whole batch.

Your turn

Print how many of the attempts failed, using count(_.isFailure).

In real jobs: limit how many run at once
Future.traverse(tenThousandIds)(callApi) fires ten thousand requests at once and gets you rate-limited or knocks over the service you call. Real code batches the input — ids.grouped(50), awaiting each batch in turn with foldLeft — or uses a streaming library (Pekko Streams, fs2) with an explicit parallelism setting.
06

Promise: completing a Future yourself

Every Future has a writer side called a Promise. You create a Promise[A], hand out promise.future, and later call success(value) or failure(exception) exactly once. The typical use is wrapping a callback-based Java API — one that takes onSuccess and onError functions — so the rest of your code gets an ordinary Future.

scalaMain.scala
import scala.concurrent.{Future, Promise, Await}
import scala.concurrent.ExecutionContext.Implicits.global
import scala.concurrent.duration.*

// A callback-style API, like many Java client libraries
def legacyGet(key: String, onDone: String => Unit, onError: Throwable => Unit): Unit =
  val t = new Thread(() =>
    if key.isEmpty then onError(new IllegalArgumentException("empty key"))
    else onDone(s"value-for-$key")
  )
  t.start()

def get(key: String): Future[String] =
  val p = Promise[String]()
  legacyGet(key, v => p.success(v), e => p.failure(e))
  p.future

@main def run(): Unit =
  println(Await.result(get("config"), 2.seconds))
  val upper = get("user").map(_.toUpperCase)
  println(Await.result(upper, 2.seconds))

  val bad = get("")
  Await.ready(bad, 2.seconds)
  println(bad.value)
Outputcompiled & run with real Scala
value-for-config
VALUE-FOR-USER
Some(Failure(java.lang.IllegalArgumentException: empty key))

After the wrapper, the callback API composes like any other Future: map, for, recover. Future.successful(x) and Future.failed(e) are the shortcuts for a Future that is already complete.

Error you will hit

Completing a Promise twice

scala
import scala.concurrent.Promise

@main def run(): Unit =
  val p = Promise[Int]()
  p.success(1)
  p.success(2)
  println(p.future.value)
Exception in thread "main" java.lang.IllegalStateException: Promise already completed.
	at scala.concurrent.Promise.complete(Promise.scala:60)
	at scala.concurrent.Promise.complete$(Promise.scala:39)
	at scala.concurrent.impl.Promise$DefaultPromise.complete(Promise.scala:116)
	at scala.concurrent.Promise.success(Promise.scala:97)
	at scala.concurrent.Promise.success$(Promise.scala:39)
	at scala.concurrent.impl.Promise$DefaultPromise.success(Promise.scala:116)
	at Main$package$.run(Main.scala:6)
	at run.main(Main.scala:3)
Why the compiler said that

A Future can only ever have one outcome. The second success tries to change it, which is always a bug. It happens in real code when a callback API calls both onDone and onError, or when a timeout and a result race to complete the same Promise.

The fix

When more than one party may legitimately try to complete the Promise, use trySuccess / tryFailure: they return false instead of throwing when someone got there first.

scala
import scala.concurrent.Promise

@main def run(): Unit =
  val p = Promise[Int]()
  println(p.trySuccess(1))
  println(p.trySuccess(2))
  println(p.future.value)
07

Blocking calls and thread pools

The global ExecutionContext has about as many threads as your machine has cores. That is ideal for CPU work, and terrible for code that waits — a JDBC query, a file read, Thread.sleep. Eight Futures each blocked on a slow database call can occupy every thread, and nothing else in the program makes progress. There are two standard remedies: wrap the waiting call in blocking { ... }, which lets the global pool add temporary threads, or run blocking work on a separate pool built from java.util.concurrent.Executors.

scalaMain.scala
import scala.concurrent.{Future, Await, ExecutionContext, blocking}
import scala.concurrent.duration.*
import java.util.concurrent.Executors

@main def run(): Unit =
  val pool = Executors.newFixedThreadPool(4)
  given ExecutionContext = ExecutionContext.fromExecutorService(pool)

  def slowQuery(id: Int): Future[(Int, String)] = Future {
    blocking(Thread.sleep(100))          // stands in for a JDBC call
    (id * id, Thread.currentThread.getName)
  }

  val all = Future.sequence((1 to 8).map(slowQuery))
  val results = Await.result(all, 5.seconds)

  println(results.map(_._1))
  println(s"threads used: at most 4? ${results.map(_._2).distinct.size <= 4}")
  pool.shutdown()
Outputcompiled & run with real Scala
Vector(1, 4, 9, 16, 25, 36, 49, 64)
threads used: at most 4? true

given ExecutionContext = ... makes this pool the one every Future in scope uses. Eight tasks shared four threads. The thread names differ from run to run, so the program prints a fact about them rather than the names themselves.

Your turn

Change the pool size to 8 and the timeout to 150 millis. Does it still finish in time?

Shut your pools down
Threads from Executors.newFixedThreadPool are not daemon threads: if you forget pool.shutdown(), main ends but the JVM keeps running forever, waiting for work that will never come. The global pool uses daemon threads, which is why earlier examples exit without any cleanup.
Kind of workWhere to run it
CPU-bound: parsing, maths, transforming collectionsThe global pool (one thread per core is right)
A few occasional blocking callsThe global pool, wrapped in blocking { ... }
Steady blocking I/O: JDBC, legacy HTTP clients, filesA dedicated fixed pool sized to what the resource allows (e.g. the DB connection count)
Many cheap waits on JDK 21+Executors.newVirtualThreadPerTaskExecutor() as the ExecutionContext
08

Shared state: atomics and synchronized

Immutable values are safe to share between threads — nobody can change them. Trouble starts with a shared var or mutable collection. count += 1 is really three steps (read, add, write), and two threads doing it at once can both read 41 and both write 42, losing an update. This is a race condition, and its result changes from run to run:

scalaMain.scala
// WRONG: four tasks increment a plain var 10,000 times each
var count = 0
val tasks = (1 to 4).map(_ => Future { for _ <- 1 to 10_000 do count += 1 })
Await.result(Future.sequence(tasks), 5.seconds)
println(count)   // 40000 expected; prints 23817, 31002, 40000... it varies

Not a verified example on purpose: its output is different on every run, which is exactly the bug.

The fix is to make each update indivisible. java.util.concurrent.atomic gives you AtomicInteger, AtomicLong and AtomicReference, whose updates happen as one step. synchronized lets only one thread at a time run a block. For maps, use java.util.concurrent.ConcurrentHashMap or Scala's collection.concurrent.TrieMap.

scalaMain.scala
import scala.concurrent.{Future, Await}
import scala.concurrent.ExecutionContext.Implicits.global
import scala.concurrent.duration.*
import java.util.concurrent.atomic.{AtomicInteger, AtomicReference}

@main def run(): Unit =
  val hits = AtomicInteger(0)
  val seen = AtomicReference(Set.empty[Int])

  var total = 0
  val lock = new Object

  val tasks = (1 to 4).map { worker =>
    Future {
      for i <- 1 to 10_000 do
        hits.incrementAndGet()
        lock.synchronized { total += 1 }
      seen.updateAndGet(s => s + worker)
    }
  }
  Await.result(Future.sequence(tasks), 5.seconds)

  println(hits.get)
  println(total)
  println(seen.get.toList.sorted)
Outputcompiled & run with real Scala
40000
40000
List(1, 2, 3, 4)

Both counters are exact. updateAndGet applies a function atomically, retrying if another thread changed the value in between — a clean way to keep an immutable Set or Map in one shared reference. The set is sorted before printing because a Set has no fixed order.

Your turn

Add an AtomicLong that sums i across all workers and print it. What should it be?

Prefer not sharing at all
The cleanest concurrent code has no shared mutable state. Let each Future return its partial result and combine them afterwards: Future.traverse(chunks)(c => Future(c.sum)).map(_.sum). Atomics are for the cases you cannot avoid, such as counters and metrics.
09

Beyond Future: Cats Effect, ZIO and Pekko

Future is in the standard library and it is what Spark, Play and most Java interop code use. It has two well-known limits: it starts running eagerly as soon as it is created (you just saw the consequence), and it cannot be cancelled. Larger Scala services often use a library instead. These snippets need their libraries, so they are shown as static code.

scala
//> using dep org.typelevel::cats-effect:3.6.3
import cats.effect.{IO, IOApp}
import cats.syntax.all.*
import scala.concurrent.duration.*

object Main extends IOApp.Simple:
  def call(name: String): IO[String] = IO.sleep(500.millis).as(name)

  val run: IO[Unit] =
    (call("prices"), call("stock"))
      .parMapN((a, b) => s"$a+$b")
      .timeout(2.seconds)
      .flatMap(IO.println)

Cats Effect: an IO is a description of work; nothing runs until the runtime runs it. parMapN runs both in parallel, and timeout really cancels.

scala
//> using dep dev.zio::zio:2.1.21
import zio.*

object Main extends ZIOAppDefault:
  def call(name: String): UIO[String] = ZIO.succeed(name).delay(500.millis)

  val run =
    call("prices").zipPar(call("stock"))
      .map((a, b) => s"$a+$b")
      .flatMap(Console.printLine(_))

ZIO: the same idea with an error type in the signature — ZIO[Environment, Error, Result].

scala
//> using dep org.apache.pekko::pekko-actor-typed:1.2.1
import org.apache.pekko.actor.typed.{ActorSystem, Behavior}
import org.apache.pekko.actor.typed.scaladsl.Behaviors

enum Command:
  case Increment
  case Report

def counter(n: Int): Behavior[Command] = Behaviors.receiveMessage {
  case Command.Increment => counter(n + 1)
  case Command.Report    => println(s"count = $n"); Behaviors.same
}

@main def run(): Unit =
  val system = ActorSystem(counter(0), "counter")
  (1 to 3).foreach(_ => system ! Command.Increment)
  system ! Command.Report
  system.terminate()

Apache Pekko (the open-source fork of Akka): actors process one message at a time, so their private state needs no locks.

Stay with Future when

  • You work in Spark, Play, or code that calls Java async APIs
  • The concurrency is simple: a few parallel calls, then combine
  • You want no extra dependency

Reach for an effect library when

  • You need cancellation, timeouts that stop work, retries and resource safety
  • The service is large and concurrency is its core job
  • The team already uses Typelevel (Cats Effect, fs2, http4s) or ZIO
Future
A value that will be available later, computed on a thread pool; completes once, with a success or a failure.
ExecutionContext
The thread pool that runs a Future's work and callbacks; passed as a given (implicit) parameter.
Await.result
Blocks the current thread until a Future completes or a timeout passes; use only at the edge of a program.
recover / recoverWith
Turn selected failures of a Future into a value, or into another Future.
Future.sequence
Turns a collection of Futures into one Future of a collection, in input order.
Future.traverse
Maps each element to a Future and sequences the results in one step.
Promise
The writable side of a Future: complete it once with success or failure.
blocking
Marks code that waits, so the global pool can add a temporary thread instead of starving.
Race condition
A bug where the result depends on the timing of threads, such as two threads losing each other's updates.
AtomicInteger
A JVM integer whose updates, such as incrementAndGet, happen as one indivisible step.
Quick check

Two independent calls take 1 second each. Which code finishes in about 1 second?

Quick check

What does Future.traverse(List(1, 2, 3))(f) return if f(2) fails?

Frequently asked questions

What is an ExecutionContext in Scala?
It is the thread pool that runs the body of a Future and every callback attached with map, flatMap or recover. Methods take it as a given (implicit) parameter. ExecutionContext.Implicits.global is a sensible default for CPU work; blocking I/O belongs in blocking { } or on a dedicated pool.
Why do my Scala Futures run sequentially inside a for-comprehension?
A for over Futures becomes nested flatMap calls, and the second Future is only created inside the first one's callback. Create independent Futures into vals before the for (or combine them with zip) so they start at the same time.
Should I use Await.result in production Scala code?
Almost never. It blocks a thread and, inside a pool, can deadlock or starve other work. Return the Future and let the framework complete the request, or await exactly once at the top of main in a command-line program.

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.