Backpressure
Behandeln Sie schnelle Produzenten sicher.
Backpressure ist eine kostenlose Scala for Backend Engineering & Functional Programming-Lektion auf CoddyKit. Dies ist Lektion 3 von 4. Du kannst die komplette Lektion unten kostenlos lesen – dann übst du sie direkt im Browser mit einem integrierten Code-Editor und einem KI-Tutor rund um die Uhr. Sie ist Teil des Scala for Backend Engineering & Functional Programming-Lernpfads, und dein Fortschritt wird über Web und CoddyKit-App synchronisiert. Der Scala for Backend Engineering & Functional Programming-Kurs umfasst insgesamt 4 Lektionen.
Teile dieser Lektion wurden noch nicht übersetzt und werden auf Englisch angezeigt.
What Is Backpressure?
Backpressure is a flow-control mechanism that prevents a fast producer from overwhelming a slow consumer. Instead of buffering unboundedly or dropping data, the consumer signals how much it can handle.
Akka Streams implements the Reactive Streams standard, where demand flows upstream and elements flow downstream.
Demand-Driven Flow
Each stage only emits when the next stage has signalled demand. A Sink requests N elements; that demand propagates upstream until a Source produces exactly what was asked for.
This pull-based protocol means producers never push more than consumers can process.
Why It Matters for Pipelines
Without backpressure, a fast Kafka consumer feeding a slow database would accumulate millions of in-flight records, exhausting memory and crashing the process.
Backpressure naturally throttles the upstream to the slowest stage, giving stable memory usage under load.
// Fast source, slow sink: backpressure slows the source
val g =
Source(1 to 1000000)
.map(_ * 2)
.to(slowDatabaseSink)Internal Buffering
Between asynchronous boundaries, Akka Streams keeps a small internal buffer (default 16 elements). It absorbs short bursts so stages need not lock-step on every element.
When the buffer fills, backpressure kicks in and the upstream stops producing until space frees up.
import akka.stream.Attributes
val buffered =
Flow[Int]
.map(identity)
.addAttributes(Attributes.inputBuffer(initial = 32, max = 32))Explicit buffer with Overflow Strategy
The buffer operator inserts an explicit buffer of a chosen size with an OverflowStrategy that decides what happens when it is full.
This lets you trade memory for the ability to decouple producer and consumer speed.
import akka.stream.OverflowStrategy
val withBuffer =
Source(1 to 1000)
.buffer(size = 100, OverflowStrategy.backpressure)Overflow Strategies
Strategies include backpressure (slow the upstream), dropHead/dropTail (discard oldest or newest), dropBuffer, dropNew, and fail (terminate with an error).
Dropping strategies suit live data like sensor readings where stale values can be discarded safely.
import akka.stream.OverflowStrategy
val latestWins =
liveTicks.buffer(1, OverflowStrategy.dropHead)
val strict =
liveTicks.buffer(50, OverflowStrategy.fail)Conflate to Summarize
When a consumer is slow, conflate merges pending elements into one using a combine function instead of buffering them all.
For example, collapse many numeric updates into their sum, so the consumer always sees an aggregate of what it missed.
val summarized =
fastMetrics
.conflate((acc, next) => acc + next)
// Slow downstream receives summed batchesExpand to Fill Demand
expand is the dual of conflate: when downstream demands faster than upstream produces, it synthesizes extra elements from the last seen value.
This is useful to keep emitting a most-recent reading at a steady rate.
val repeated =
sensor.expand(last => Iterator.continually(last))
// Downstream always gets the latest sensor valueAsync Boundaries
By default, fused stages run on a single actor with no buffering between them. Inserting async places a stage on its own actor, adding a buffer and enabling pipelined parallelism.
Async boundaries are where backpressure buffers actually live.
val pipelined =
Source(1 to 1000)
.map(slowStep).async
.map(anotherSlowStep).async
.to(Sink.ignore)Throttle as Explicit Rate Control
throttle imposes a deliberate maximum rate, generating backpressure upstream to honor it. This protects rate-limited external services even when the consumer could go faster.
A burst parameter allows short spikes above the steady rate.
import scala.concurrent.duration._
val limited =
requests
.throttle(
elements = 100, per = 1.second, maximumBurst = 20,
akka.stream.ThrottleMode.Shaping)Observing Backpressure
You can detect backpressure by watching for upstream slowdowns or by measuring buffer occupancy. The log operator and Akka's stream attributes help trace where a pipeline stalls.
A persistently full buffer indicates the slowest stage that is gating throughput.
val traced =
Source(1 to 100)
.log("after-source")
.map(_ * 2)
.log("after-map")
.to(Sink.ignore)Quick Check
Think about how Akka Streams keeps a fast producer from flooding a slow consumer.
Recap
Backpressure is the demand-driven backbone of Akka Streams: consumers signal demand upstream so producers cannot overwhelm them, keeping memory bounded.
You saw internal buffers, the explicit buffer operator with overflow strategies, summarizing with conflate, filling demand with expand, async boundaries, and deliberate rate control via throttle. Next you will run a complete pipeline.
Häufig gestellte Fragen
Ist die Lektion „Backpressure“ kostenlos?
Ja — der vollständige Text von „Backpressure“ ist hier im Web kostenlos zu lesen. Um sie interaktiv zu üben (integrierter Code-Editor und 24/7 KI-Tutor) und den Rest des Scala for Backend Engineering & Functional Programming-Kurses freizuschalten, upgrade auf CoddyKit PRO. Der Scala for Backend Engineering & Functional Programming-Kurs umfasst insgesamt 4 Lektionen.
Was lerne ich in „Backpressure“?
Behandeln Sie schnelle Produzenten sicher. Du übst Scala for Backend Engineering & Functional Programming mit praktischem Code, den du direkt im Browser ausführst, und ein 24/7 KI-Tutor beantwortet deine Fragen während du die Lektion bearbeitest.
Brauche ich Erfahrung, um Scala for Backend Engineering & Functional Programming zu starten?
Keine Vorkenntnisse erforderlich. Scala for Backend Engineering & Functional Programming auf CoddyKit ist für Anfänger bis fortgeschrittene Lernende strukturiert, sodass du hier starten oder von Anfang an beginnen und in deinem eigenen Tempo voranschreiten kannst. Dies ist Lektion 3 von 4.
Wie lange dauert die Lektion „Backpressure“?
Die meisten CoddyKit-Lektionen dauern etwa 5–10 Minuten. Jede ist kompakt und interaktiv, sodass du stetig Fortschritte machst und genau dort weitermachst, wo du aufgehört hast – im Web und in der App.
Kann ich in dieser Scala for Backend Engineering & Functional Programming-Lektion Code schreiben und ausführen?
Ja. Jede Scala for Backend Engineering & Functional Programming-Lektion enthält einen integrierten Code-Editor, sodass du echten Code direkt in deinem Browser schreibst und ausführst und sofort KI-Feedback erhältst — ohne lokale Einrichtung erforderlich.