Scala for Backend Engineering & Functional Programming · Lektion

Source, Flow und Sink

Die Grundbausteine von Streams.

Lektion 1 von 413 Schritte

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.ignore

Schnelltest

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.

Kostenlos starten

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

  1. Source, Flow und Sink
  2. Streams transformieren
  3. Backpressure
  4. Eine Pipeline ausführen
← Zurück zu Scala for Backend Engineering & Functional Programming