Trasformare gli stream
Mappi e filtri i dati in movimento.
Trasformare gli stream è una lezione Scala for Backend Engineering & Functional Programming gratuita su CoddyKit. Questa è la lezione 2 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.
Gli operatori come trasformazioni
Akka Streams fornisce un ricco insieme di operatori per Source e Flow, simili alle API delle collezioni di Scala, ma eseguiti in modo asincrono e nel rispetto della backpressure.
Ogni operatore restituisce un nuovo blueprint, quindi le trasformazioni vengono composte in modo dichiarativo prima ancora che lo stream venga eseguito.
map e filter
map applica una funzione sincrona a ogni elemento; filter elimina gli elementi che non soddisfano un predicato. Sono gli strumenti fondamentali per le trasformazioni elemento per elemento.
Entrambi mantengono l'ordine e propagano il completamento e gli errori a valle.
val flow =
Flow[Int]
.filter(_ % 2 == 0)
.map(n => n * n)mapConcat per trasformazioni uno-a-molti
Quando un input deve produrre diversi output, usi mapConcat. Riceve una funzione che restituisce un iterable e appiattisce i risultati nello stream.
Restituire una collezione vuota equivale di fatto a eliminare l'elemento.
val explode: Flow[String, String, akka.NotUsed] =
Flow[String].mapConcat(line => line.split(",").toList)
val words = Source(List("a,b", "c,d,e"))
.via(explode)grouped e sliding
grouped(n) raggruppa gli elementi consecutivi in una Seq contenente fino a n elementi, utile per le scritture batch nei database. sliding(n) emette finestre sovrapposte.
Il raggruppamento riduce il sovraccarico per elemento nelle pipeline con molte operazioni di I/O.
val batches: Source[Seq[Int], akka.NotUsed] =
Source(1 to 1000).grouped(100)
val windows =
Source(1 to 10).sliding(3, step = 1)scan e fold
scan emette l'accumulatore aggiornato dopo ogni elemento, creando uno stream dello stato in evoluzione. fold emette soltanto il valore accumulato finale, una volta completato lo stream a monte.
Usi scan per i contatori in tempo reale e fold per gli aggregati finali.
val running =
Source(1 to 5).scan(0)(_ + _) // 0,1,3,6,10,15
val total =
Source(1 to 5).fold(0)(_ + _) // 15mapAsync per le operazioni asincrone
mapAsync(parallelism) chiama una funzione che restituisce un Future ed emette i risultati in ordine, eseguendo contemporaneamente fino a parallelism Future.
Lo usi per chiamate asincrone come le ricerche nei database o le richieste HTTP quando l'ordine è importante.
import scala.concurrent.Future
val enriched =
Flow[UserId]
.mapAsync(parallelism = 4)(id => lookup(id))
def lookup(id: UserId): Future[User] = ???mapAsyncUnordered
mapAsyncUnordered si comporta come mapAsync, ma emette ogni risultato appena viene completato, ignorando l'ordine degli input.
Può migliorare il throughput quando l'ordine non è importante per lo stage a valle, perché un Future lento non blocca più quelli più veloci.
val fast =
Flow[UserId]
.mapAsyncUnordered(parallelism = 8)(id => lookup(id))Trasformazioni con stato con statefulMapConcat
Per le trasformazioni elemento per elemento che richiedono uno stato locale mutabile, statefulMapConcat crea un nuovo stato a ogni materializzazione e restituisce un iterable di output.
È il modo sicuro per mantenere contatori o buffer senza condividere lo stato tra diverse esecuzioni dello stream.
val withIndex: Flow[String, (Int, String), akka.NotUsed] =
Flow[String].statefulMapConcat { () =>
var i = 0
elem => { i += 1; List((i, elem)) }
}Operatori basati sul tempo
Gli stream possono trasformare i dati in base al tempo oltre che al contenuto. throttle limita la frequenza di emissione, groupedWithin raggruppa in base alla dimensione o al tempo trascorso e takeWithin limita la durata.
Questi operatori sono essenziali per limitare la frequenza delle API esterne.
import scala.concurrent.duration._
val limited =
Source(1 to 1000)
.throttle(10, 1.second)
.groupedWithin(100, 500.millis)Gestione degli errori nelle trasformazioni
Per impostazione predefinita, un'eccezione generata all'interno di un operatore causa il fallimento dell'intero stream. Una supervision strategy può invece resume (eliminando l'elemento non valido) o restart lo stage.
Applichi la strategia al Flow con withAttributes.
import akka.stream.{ActorAttributes, Supervision}
val safe =
Flow[String].map(_.toInt)
.withAttributes(
ActorAttributes.supervisionStrategy(_ => Supervision.Resume))Composizione dei Flow
È possibile comporre piccoli Flow in Flow più grandi con via, ottenendo un unico Flow riutilizzabile. In questo modo ogni trasformazione rimane focalizzata e può essere testata autonomamente.
Il Flow composto ha il tipo di input del primo Flow e il tipo di output dell'ultimo.
val parse = Flow[String].map(_.toInt)
val square = Flow[Int].map(n => n * n)
val parseAndSquare: Flow[String, Int, akka.NotUsed] =
parse.via(square)Verifica rapida
Consideri le trasformazioni asincrone e le relative garanzie sull'ordine.
Riepilogo
Ha esplorato gli operatori di trasformazione: map/filter elemento per elemento, mapConcat uno-a-molti, grouped per il raggruppamento, scan e fold per l'accumulazione e le operazioni asincrone con mapAsync.
Ha inoltre visto le trasformazioni con stato, gli operatori basati sul tempo come throttle, le supervision strategy per gli errori e la composizione dei Flow con via. Ora vedrà come la backpressure mantiene sicuri questi stage.
Domande Frequenti
La lezione «Trasformare gli stream» è gratuita?
Sì — il testo completo di «Trasformare gli stream» è 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 «Trasformare gli stream»?
Mappi e filtri i dati in movimento. 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 2 di 4.
Quanto tempo richiede la lezione «Trasformare gli stream»?
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.
Tutte le lezioni di questo corso
- Source, Flow e Sink
- Trasformare gli stream
- Backpressure
- Eseguire una pipeline