Practical Materialized Values in Akka Streams: Coordinating Side Effects Without Blocking
Learn how to use Akka Streams' materialized values to coordinate side effects like database writes while keeping backpressure intact, with a concrete Scala example and trade‑off analysis.
09 Feb 2026, 05:08 UTC

The problem: side effects in a non‑blocking pipeline
Akka Streams excels at moving data through a series of transformations while preserving backpressure. But real‑world pipelines often need to perform side effects—writing to a database, publishing metrics, or sending a notification—without breaking the non‑blocking contract. If you shove a blocking call into a map stage you stall the whole stream; if you fire‑and‑forget you lose visibility into success or failure.
The framework solves this with materialized values: each stage can produce a value (often a Future or CompletionStage) that becomes available once the stream is materialized (i.e., started). By wiring these values together you can coordinate side effects while the stream continues to flow.
What a materialized value actually is
When you connect a Source, a Flow, and a Sink you get a RunnableGraph. Calling run() on that graph materializes it—allocates actors, opens connections, and returns a materialized value defined by the sink (or by a custom graph). The type of that value is declared in the sink’s signature, e.g. Sink.foreachAsync yields a Future[Done], while Sink.fold yields a Future[A] with the folded result.
Because the materialized value is produced after the stream starts, you can attach callbacks (onComplete, map) that run when the side effect finishes, all without blocking the stream’s processing threads.
Worked example: counting processed records while writing to PostgreSQL
Below is a minimal Scala 2.13 / Akka 2.8 snippet that reads JSON lines from a file, parses them, inserts each record into a table, and finally returns the total number of rows inserted. The example uses slick for DB access, but any async driver works the same way.
import akka.actor.ActorSystem
import akka.stream.scaladsl._
import akka.stream.Materializer
import scala.concurrent.{ExecutionContext, Future}
import slick.jdbc.PostgresProfile.api._
case class Event(id: Long, payload: String)
object EventStream extends App {
implicit val system: ActorSystem = ActorSystem("event-stream")
implicit val mat: Materializer = Materializer(system)
implicit val ec: ExecutionContext = system.dispatcher
val db = Database.forConfig("postgres") // defined in application.conf
// 1. Source: read lines from a file (one JSON per line)
val source: Source[String, Future[IOResult]] =
FileIO.fromPath(java.nio.file.Paths.get("events.jsonl"))
.via(Framing.delimiter(ByteString("\n"), maximumFrameLength = 256, allowTruncation = true))
.map(_.utf8String)
// 2. Flow: parse JSON → case class (using your favourite JSON lib)
val parseFlow: Flow[String, Event, NotUsed] =
Flow[String].map { line =>
// placeholder for real parsing, e.g. circe.decode[Event](line).getOrElse(throw ...)
Event(line.split(",")(0).toLong, line.split(",")(1))
}
// 3. Sink: async DB insert, returns Future[Int] = count of inserted rows
val insertSink: Sink[Event, Future[Int]] =
Sink.foldAsync[Int, Event](0) { (count, evt) =>
val action = sqlu"INSERT INTO events (id, payload) VALUES (${evt.id}, ${evt.payload})"
db.run(action).map(_ => count + 1)(ec)
}
// 4. Assemble and run – the materialized value is Future[Int]
val resultFuture: Future[Int] = source.via(parseFlow).runWith(insertSink)
// 5. Handle completion
resultFuture.onComplete {
case scala.util.Success(cnt) =>
println(s"Inserted $cnt records")
system.terminate()
case scala.util.Failure(ex) =>
println(s"Stream failed: ${ex.getMessage}")
system.terminate()
}
}
Where to run this: compile with sbt (requires akka-stream, akka-actor-typed, slick, postgresql dependencies). Execute the main class. You need a running PostgreSQL instance and a table events(id bigint, payload text). The application.conf must contain a postgres slick profile with URL, user, password.
Expected checks: after the program finishes, query SELECT count(*) FROM events; – it should match the printed count. If the file contains malformed lines, the parseFlow will throw and the stream will fail, propagating the exception into resultFuture.
Risks: foldAsync processes elements sequentially; for high throughput you may replace it with mapAsyncUnordered feeding a Sink.ignore and a separate counter, but then you lose the single materialized count.
Trade‑off: materialized values vs. side‑effecting sinks
- Materialized value (as above) – gives you a typed handle (e.g.
Future[Int]) that composes with other futures. Ideal when you need to know *when* the side effect finishes or aggregate results. - Side‑effecting sink (
Sink.foreachAsync) – simpler code, returnsFuture[Done]. You lose the aggregated result unless you add a separateSink.folddownstream, which means two passes over the data or a custom graph. - Latency impact –
foldAsyncforces sequential DB calls. If you need parallel inserts, usemapAsync(parallelism)into aSink.ignoreand keep a separate atomic counter; the trade‑off is more complex materialization logic.
Actionable closing: verify the pattern in your codebase
- Identify a stream that currently performs blocking I/O inside a
maporforeach. - Replace that stage with a sink that returns a materialized
Future(e.g.Sink.foldAsyncor a customGraphStagewith a materialized value). - Run the stream under a realistic load (use
akka.stream.testkit.scaladsl.TestSourcefor unit tests) and assert the materialized future completes with the expected value. - Monitor backpressure metrics (
akka.stream.Materializerexposessubscriptioncallbacks) to ensure the new sink does not introduce unbounded buffers.
By making the side effect’s completion a first‑class materialized value you keep the stream reactive, gain composable error handling, and avoid the classic “blocking in a map” anti‑pattern.
0 replies
A thoughtful contribution can make all the difference. Be the first to share one.