Source, Flow und Sink
Die Grundbausteine von Streams.
Source, Flow und Sink ist eine kostenlose Scala for Backend Engineering & Functional Programming-Lektion auf CoddyKit. Dies ist Lektion 1 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.
Die drei Bausteine
Akka Streams modelliert eine Datenpipeline als Graph aus Verarbeitungsstufen. Die drei grundlegenden linearen Stufen sind Source (erzeugt Elemente), Flow (transformiert sie) und Sink (verbraucht sie).
Eine Source hat einen Ausgang, ein Sink einen Eingang und ein Flow genau einen Eingang und einen Ausgang. Durch ihre Verbindung wird beschrieben, was geschehen soll, nicht wann.
Eine Source definieren
Eine Source[Out, Mat] gibt Elemente des Typs Out aus und stellt einen materialisierten Wert des Typs Mat bereit. Die einfachsten Sources entstehen aus In-Memory-Collections oder Wertebereichen.
Bis der Stream ausgeführt wird, ist eine Source lediglich ein unveränderlicher Bauplan, den Sie frei wiederverwenden können.
import akka.stream.scaladsl.Source
val numbers: Source[Int, akka.NotUsed] =
Source(1 to 100)
val single: Source[String, akka.NotUsed] =
Source.single("hello")Ein Sink definieren
Ein Sink[In, Mat] verarbeitet Elemente des Typs In]. Der materialisierte Wert enthält häufig das Ergebnis der Verarbeitung, etwa ein Future, das abgeschlossen wird, wenn der Stream endet.
Sink.foreach führt für jedes Element einen Seiteneffekt aus; Sink.fold akkumuliert ein einzelnes Ergebnis.
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)(_ + _)Einen Flow definieren
Ein Flow[In, Out, Mat] befindet sich zwischen einer Source und einem Sink und transformiert jedes Element. Flows können eigenständig wiederverwendet und vor ihrer Verbindung mit einem Endpunkt kombiniert werden.
Hier verdoppelt ein Flow Ganzzahlen und wandelt sie in Zeichenketten um.
import akka.stream.scaladsl.Flow
val doubleToString: Flow[Int, String, akka.NotUsed] =
Flow[Int]
.map(_ * 2)
.map(n => s"value=$n")Source und Sink verbinden
Der Operator via verbindet einen Flow mit einer Source, und to verbindet einen Sink. Wenn Sie eine Source direkt mit einem Sink über to verbinden, entsteht ein RunnableGraph: ein geschlossener, ausführbarer Bauplan.
Noch werden keine Elemente übertragen; hier wird lediglich die Topologie beschrieben.
import akka.stream.scaladsl.{Source, Sink, RunnableGraph}
val graph: RunnableGraph[akka.NotUsed] =
Source(1 to 10).to(Sink.foreach(println))via: Einen Flow einfügen
Verwenden Sie via, um einen Flow in die Pipeline einzufügen. Ein Source.via(flow) liefert eine neue Source, deren Ausgabetyp dem Ausgabetyp des Flows entspricht.
Durch das Verketten von via-Aufrufen können Sie aus kleinen, testbaren Flow-Bausteinen lange Transformationspipelines erstellen.
val pipeline =
Source(1 to 10)
.via(Flow[Int].filter(_ % 2 == 0))
.via(Flow[Int].map(_ * 10))
.to(Sink.foreach(println))Typsicherheit über mehrere Stufen hinweg
Der Compiler stellt sicher, dass der Ausgabetyp jeder Stufe mit dem Eingabetyp der nächsten übereinstimmt. Eine Source[Int] kann nicht ohne einen zwischengeschalteten Flow, der den Typ umwandelt, mit einem Sink[String] verbunden werden.
Diese statische Prüfung erkennt Fehler bei der Pipeline-Verdrahtung bereits vor der Laufzeit.
// Source[Int] -> Flow[Int, String] -> Sink[String]
val ok =
Source(1 to 3)
.via(Flow[Int].map(_.toString))
.to(Sink.foreach[String](println))Materialisierte Werte
Jeder Bauplan enthält einen materialisierten Wert: ein Handle, das beim Ausführen des Streams erzeugt wird. Ein Sink.fold materialisiert einen Future für das Ergebnis. Standardmäßig wird beim Kombinieren von Stufen der materialisierte Wert ganz links beibehalten (NotUsed bei einfachen Sources).
Verwenden Sie toMat und Keep, um auszuwählen, welchen Wert der beiden Seiten Sie behalten möchten.
import akka.stream.scaladsl.Keep
import scala.concurrent.Future
val g: RunnableGraph[Future[Int]] =
Source(1 to 100)
.toMat(Sink.fold(0)(_ + _))(Keep.right)Wiederverwendbare Komponenten
Da Sources, Flows und Sinks unveränderliche Werte sind, können Sie sie einmal definieren und in vielen Pipelines wiederverwenden. Das fördert eine Bibliothek aus kleinen, benannten Verarbeitungsstufen.
Ein zum Parsen definierter Flow kann sowohl in eine Datei- als auch in eine HTTP-Pipeline eingefügt werden.
val parse: Flow[String, Int, akka.NotUsed] =
Flow[String].map(_.trim.toInt)
val fromFile = lines.via(parse)
val fromHttp = requestBody.via(parse)Gängige Source-Konstruktoren
Akka Streams liefert zahlreiche Source-Factories: Source.single, Source.repeat, Source.tick für zeitgesteuerte Ausgaben, Source.future aus einem Future und Source.empty.
Mit dem passenden Konstruktor wird die Absicht eines Datenproduzenten eindeutig ausgedrückt.
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")Gängige Sink-Konstruktoren
Zu den Sinks gehören entsprechend Sink.head (erstes Element als Future), Sink.seq (alle Elemente in einer Seq sammeln), Sink.ignore (auslesen und verwerfen) und Sink.last.
Für Pipelines, die ein Ergebnis an Ihren Code zurückgeben, werden am häufigsten Sink.seq und Sink.fold verwendet.
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.ignoreSchnelltest
Testen Sie Ihr Verständnis der linearen Stufentypen.
Zusammenfassung
Sie haben die drei linearen Bausteine kennengelernt: Source erzeugt Daten, Flow transformiert sie und Sink verarbeitet sie. Es handelt sich um unveränderliche, wiederverwendbare Baupläne, die mit via und to verbunden werden.
Wenn Sie eine Source mit einem Sink verbinden, erhalten Sie einen RunnableGraph, der einen materialisierten Wert enthält, aber erst beim Ausführen Daten verarbeitet. Als Nächstes transformieren Sie Streams mit leistungsfähigeren Operatoren.
Lerne Scala mit einem KI-Tutor — kostenlos
Schreibe und führe echten Code in deinem Browser aus, bekomme sofortige Hilfe von einem 24/7 KI-Tutor und setze dein Lernen im Web oder in der App fort.
- Kurse
- 39
- Lektionen
- 143
Häufig gestellte Fragen
Ist die Lektion „Source, Flow und Sink“ kostenlos?
Ja — der vollständige Text von „Source, Flow und Sink“ 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 „Source, Flow und Sink“?
Die Grundbausteine von Streams. 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 1 von 4.
Wie lange dauert die Lektion „Source, Flow und Sink“?
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.
Alle Lektionen in diesem Kurs
- Source, Flow und Sink
- Streams transformieren
- Backpressure
- Eine Pipeline ausführen