Source, Flow e Sink
I componenti fondamentali dello streaming.
Source, Flow e Sink è una lezione Scala for Backend Engineering & Functional Programming gratuita su CoddyKit. Questa è la lezione 1 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.
I tre elementi fondamentali
Akka Streams modella una pipeline di dati come un grafo di fasi di elaborazione. Le tre fasi lineari fondamentali sono Source (produce elementi), Flow (li trasforma) e Sink (li consuma).
Una Source ha un'uscita, un Sink ha un ingresso e un Flow ha esattamente un ingresso e un'uscita. Collegandoli si descrive che cosa deve accadere, non quando.
Definire una Source
Una Source[Out, Mat] emette elementi di tipo Out ed espone un valore materializzato di tipo Mat. Le Source più semplici provengono da collezioni in memoria o da intervalli.
Finché lo stream non viene eseguito, una Source è solo un blueprint immutabile che può essere riutilizzato liberamente.
import akka.stream.scaladsl.Source
val numbers: Source[Int, akka.NotUsed] =
Source(1 to 100)
val single: Source[String, akka.NotUsed] =
Source.single("hello")Definire un Sink
Un Sink[In, Mat] consuma elementi di tipo In. Il valore materializzato spesso contiene il risultato del consumo, ad esempio un Future che si completa quando lo stream termina.
Sink.foreach esegue un effetto collaterale per ogni elemento; Sink.fold accumula un singolo risultato.
import akka.stream.scaladsl.Sink
import scala.concurrent.Future
val printSink: Sink[Int, Future[akka.Done]] =
Sink.foreach(println)
val sumSink: Sink[Int, Future[Int]] =
Sink.fold(0)(_ + _)Definire un Flow
Un Flow[In, Out, Mat] si trova tra una Source e un Sink e trasforma ogni elemento. I Flow sono riutilizzabili autonomamente e possono essere composti prima di essere collegati a un endpoint.
In questo caso un Flow raddoppia gli interi e li converte in stringhe.
import akka.stream.scaladsl.Flow
val doubleToString: Flow[Int, String, akka.NotUsed] =
Flow[Int]
.map(_ * 2)
.map(n => s"value=$n")Collegare una Source a un Sink
L'operatore via collega un Flow a una Source, mentre to collega un Sink. Collegando direttamente una Source a un Sink con to si produce un RunnableGraph: un blueprint completo ed eseguibile.
Gli elementi non vengono ancora trasferiti; si sta solo descrivendo la topologia.
import akka.stream.scaladsl.{Source, Sink, RunnableGraph}
val graph: RunnableGraph[akka.NotUsed] =
Source(1 to 10).to(Sink.foreach(println))via: inserire un Flow
Usi via per inserire un Flow nella pipeline. Un Source.via(flow) produce un nuovo Source il cui tipo di output corrisponde a quello dell'output del Flow.
È possibile concatenare le chiamate a via per costruire lunghe pipeline di trasformazione a partire da piccoli Flow, facili da testare.
val pipeline =
Source(1 to 10)
.via(Flow[Int].filter(_ % 2 == 0))
.via(Flow[Int].map(_ * 10))
.to(Sink.foreach(println))Sicurezza dei tipi tra gli stage
Il compilatore verifica che il tipo di output di ogni stage corrisponda al tipo di input dello stage successivo. Un Source[Int] non può collegarsi a un Sink[String] senza un Flow intermedio che converta il tipo.
Questa verifica statica rileva gli errori di collegamento della pipeline prima dell'esecuzione.
// Source[Int] -> Flow[Int, String] -> Sink[String]
val ok =
Source(1 to 3)
.via(Flow[Int].map(_.toString))
.to(Sink.foreach[String](println))Valori materializzati
Ogni blueprint contiene un valore materializzato: un riferimento prodotto quando lo stream viene eseguito. Un Sink.fold materializza un Future del risultato. Per impostazione predefinita, la combinazione degli stage conserva il valore materializzato più a sinistra (NotUsed per i source semplici).
Usi toMat e Keep per scegliere il valore del lato da conservare.
import akka.stream.scaladsl.Keep
import scala.concurrent.Future
val g: RunnableGraph[Future[Int]] =
Source(1 to 100)
.toMat(Sink.fold(0)(_ + _))(Keep.right)Componenti riutilizzabili
Poiché Source, Flow e Sink sono valori immutabili, è possibile definirli una volta e riutilizzarli in molte pipeline. In questo modo si favorisce la creazione di una libreria di piccoli stage di elaborazione con nomi significativi.
Un Flow definito per l'analisi dei dati può essere inserito sia in una pipeline basata su file sia in una pipeline HTTP.
val parse: Flow[String, Int, akka.NotUsed] =
Flow[String].map(_.trim.toInt)
val fromFile = lines.via(parse)
val fromHttp = requestBody.via(parse)Costruttori comuni di Source
Akka Streams include molte factory per Source: Source.single, Source.repeat, Source.tick per l'emissione temporizzata, Source.future a partire da un Future e Source.empty.
Scegliere il costruttore corretto rende esplicita l'intenzione del produttore di dati.
import scala.concurrent.duration._
val ticks = Source.tick(0.seconds, 1.second, "tick")
val onceF = Source.future(scala.concurrent.Future.successful(42))
val forever = Source.repeat("x")Costruttori comuni di Sink
Analogamente, tra i sink sono disponibili Sink.head (il primo elemento come Future), Sink.seq (raccoglie tutto in una Seq), Sink.ignore (esaurisce lo stream e scarta gli elementi) e Sink.last.
Nelle pipeline che restituiscono un risultato al codice, Sink.seq e Sink.fold sono i più comuni.
import scala.concurrent.Future
val collect: Sink[Int, Future[Seq[Int]]] = Sink.seq
val firstOne: Sink[Int, Future[Int]] = Sink.head
val drain: Sink[Int, Future[akka.Done]] = Sink.ignoreVerifica rapida
Verifichi la Sua comprensione dei tipi degli stage lineari.
Riepilogo
Ha appreso i tre elementi costitutivi lineari: Source produce, Flow trasforma e Sink consuma. Sono blueprint immutabili e riutilizzabili, collegati con via e to.
Collegare un Source a un Sink produce un RunnableGraph che contiene un valore materializzato, ma non trasferisce dati finché non viene eseguito. Ora trasformerà gli stream con operatori più avanzati.
Domande Frequenti
La lezione «Source, Flow e Sink» è gratuita?
Sì — il testo completo di «Source, Flow e Sink» è 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 «Source, Flow e Sink»?
I componenti fondamentali dello streaming. 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 1 di 4.
Quanto tempo richiede la lezione «Source, Flow e Sink»?
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