Backpressure
Gestisca in sicurezza i produttori veloci.
Backpressure è una lezione Scala for Backend Engineering & Functional Programming gratuita su CoddyKit. Questa è la lezione 3 di 4. Puoi leggere la lezione completa qui gratuitamente — poi esercitati direttamente nel browser con un editor di codice integrato e un tutor IA disponibile 24/7. Fa parte del percorso di apprendimento Scala for Backend Engineering & Functional Programming, e i tuoi progressi si sincronizzano tra il web e l'app CoddyKit. Il corso Scala for Backend Engineering & Functional Programming include 4 lezioni in totale.
Che cos'è la backpressure?
La backpressure è un meccanismo di controllo del flusso che impedisce a un produttore veloce di sovraccaricare un consumatore lento. Invece di usare un buffering illimitato o di scartare i dati, il consumatore segnala quanti dati è in grado di gestire.
Akka Streams implementa lo standard Reactive Streams, in cui la richiesta fluisce a monte e gli elementi fluiscono a valle.
Flusso guidato dalla richiesta
Ogni stage emette dati solo quando lo stage successivo ha segnalato una richiesta. Un Sink richiede N elementi; la richiesta si propaga a monte finché un Source produce esattamente quanto richiesto.
Questo protocollo basato sul pull fa sì che i produttori non inviino mai più dati di quanti i consumatori possano elaborare.
Perché è importante nelle pipeline
Senza backpressure, un consumatore Kafka veloce che alimenta un database lento accumulerebbe milioni di record in elaborazione, esaurendo la memoria e causando il crash del processo.
La backpressure limita naturalmente la velocità dello stream a monte in base allo stage più lento, mantenendo stabile l'uso della memoria sotto carico.
// Fast source, slow sink: backpressure slows the source
val g =
Source(1 to 1000000)
.map(_ * 2)
.to(slowDatabaseSink)Buffering interno
Tra i confini asincroni, Akka Streams mantiene un piccolo buffer interno (16 elementi per impostazione predefinita). Questo assorbe i picchi brevi, evitando che gli stage debbano procedere in perfetta sincronia per ogni elemento.
Quando il buffer si riempie, entra in azione la backpressure e lo stream a monte smette di produrre finché non si libera spazio.
import akka.stream.Attributes
val buffered =
Flow[Int]
.map(identity)
.addAttributes(Attributes.inputBuffer(initial = 32, max = 32))Buffer esplicito con Overflow Strategy
L'operatore buffer inserisce un buffer esplicito di dimensione scelta, con una OverflowStrategy che decide cosa fare quando il buffer è pieno.
In questo modo è possibile scambiare memoria con la possibilità di disaccoppiare la velocità del produttore da quella del consumatore.
import akka.stream.OverflowStrategy
val withBuffer =
Source(1 to 1000)
.buffer(size = 100, OverflowStrategy.backpressure)Strategie di overflow
Le strategie includono backpressure (rallenta lo stream a monte), dropHead/dropTail (scarta rispettivamente l'elemento più vecchio o quello più recente), dropBuffer, dropNew e fail (termina con un errore).
Le strategie di scarto sono adatte ai dati in tempo reale, come le letture dei sensori, per i quali è possibile eliminare senza rischi i valori obsoleti.
import akka.stream.OverflowStrategy
val latestWins =
liveTicks.buffer(1, OverflowStrategy.dropHead)
val strict =
liveTicks.buffer(50, OverflowStrategy.fail)Conflate per riassumere
Quando un consumatore è lento, conflate unisce gli elementi in attesa in un unico elemento usando una funzione di combinazione, invece di inserirli tutti nel buffer.
Ad esempio, può condensare molti aggiornamenti numerici nella loro somma, così il consumatore vede sempre un aggregato di ciò che non ha elaborato.
val summarized =
fastMetrics
.conflate((acc, next) => acc + next)
// Slow downstream receives summed batchesExpand per soddisfare la richiesta
expand è il duale di conflate: quando lo stream a valle richiede dati più velocemente di quanto lo stream a monte li produca, genera elementi aggiuntivi a partire dall'ultimo valore osservato.
È utile per continuare a emettere la lettura più recente a una frequenza costante.
val repeated =
sensor.expand(last => Iterator.continually(last))
// Downstream always gets the latest sensor valueConfini asincroni
Per impostazione predefinita, gli stage fusi vengono eseguiti su un unico actor, senza buffering tra loro. Inserire async colloca uno stage su un actor dedicato, aggiungendo un buffer e abilitando il parallelismo in pipeline.
È ai confini asincroni che risiedono effettivamente i buffer della backpressure.
val pipelined =
Source(1 to 1000)
.map(slowStep).async
.map(anotherSlowStep).async
.to(Sink.ignore)Throttle come controllo esplicito della frequenza
throttle impone una frequenza massima deliberata, generando backpressure a monte per rispettarla. In questo modo protegge i servizi esterni soggetti a limiti di frequenza, anche quando il consumatore potrebbe procedere più velocemente.
Un parametro di burst consente brevi picchi superiori alla frequenza normale.
import scala.concurrent.duration._
val limited =
requests
.throttle(
elements = 100, per = 1.second, maximumBurst = 20,
akka.stream.ThrottleMode.Shaping)Monitoraggio della backpressure
È possibile rilevare la backpressure osservando i rallentamenti a monte o misurando l'occupazione del buffer. L'operatore log e gli attributi degli stream di Akka aiutano a individuare il punto in cui una pipeline si arresta.
Un buffer costantemente pieno indica lo stage più lento, che limita il throughput.
val traced =
Source(1 to 100)
.log("after-source")
.map(_ * 2)
.log("after-map")
.to(Sink.ignore)Verifica rapida
Rifletta su come Akka Streams impedisce a un produttore veloce di sommergere un consumatore lento.
Riepilogo
La backpressure è l'elemento portante di Akka Streams, basato sulla richiesta: i consumatori segnalano la richiesta a monte, così i produttori non possono sovraccaricarli e la memoria rimane sotto controllo.
Ha visto i buffer interni, l'operatore buffer esplicito con le strategie di overflow, il riepilogo con conflate, il soddisfacimento della richiesta con expand, i confini asincroni e il controllo deliberato della frequenza tramite throttle. Ora eseguirà una pipeline completa.
Domande Frequenti
La lezione «Backpressure» è gratuita?
Sì — il testo completo di «Backpressure» è gratuito qui sul web. Per esercitarvi in modo interattivo (un editor di codice integrato e un tutor IA 24/7) e sbloccare il resto del corso Scala for Backend Engineering & Functional Programming, passa a CoddyKit PRO. Il corso Scala for Backend Engineering & Functional Programming include 4 lezioni in totale.
Cosa imparerò in «Backpressure»?
Gestisca in sicurezza i produttori veloci. Eserciti Scala for Backend Engineering & Functional Programming con codice pratico che esegui direttamente nel browser, e un tutor IA 24/7 risponde alle tue domande mentre lavori sulla lezione.
Ho bisogno di esperienza per iniziare Scala for Backend Engineering & Functional Programming?
Non è richiesta alcuna esperienza precedente. Scala for Backend Engineering & Functional Programming su CoddyKit è strutturato per principianti e studenti avanzati, quindi puoi iniziare da qui o dall'inizio e procedere al tuo ritmo. Questa è la lezione 3 di 4.
Quanto tempo richiede la lezione «Backpressure»?
La maggior parte delle lezioni CoddyKit richiede circa 5–10 minuti. Ogni lezione è breve e interattiva, quindi fai progressi costanti e riprendi esattamente da dove hai lasciato su web e app.
Posso scrivere ed eseguire codice in questa lezione Scala for Backend Engineering & Functional Programming?
Sì. Ogni lezione Scala for Backend Engineering & Functional Programming include un editor di codice integrato, quindi scrivi ed esegui codice reale direttamente nel tuo browser e ricevi feedback istantaneo dall'IA — nessuna configurazione locale necessaria.