Streams transformieren
Mappen und filtern Sie fließende Daten.
Streams transformieren ist eine kostenlose Scala for Backend Engineering & Functional Programming-Lektion auf CoddyKit. Dies ist Lektion 2 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.
Operatoren als Transformationen
Akka Streams bietet eine umfangreiche Sammlung von Operatoren für Source und Flow, die der API von Scala-Sammlungen ähnelt, aber asynchron ausgeführt wird und Backpressure berücksichtigt.
Jeder Operator gibt einen neuen Bauplan zurück. Dadurch werden Transformationen deklarativ zusammengestellt, bevor der Stream überhaupt ausgeführt wird.
map und filter
map wendet eine synchrone Funktion auf jedes Element an; filter verwirft Elemente, für die ein Prädikat nicht erfüllt ist. Diese beiden Operatoren bilden das grundlegende Werkzeug für elementweise Transformationen.
Beide behalten die Reihenfolge bei und leiten Abschluss und Fehler an nachgelagerte Stufen weiter.
val flow =
Flow[Int]
.filter(_ % 2 == 0)
.map(n => n * n)mapConcat für Eins-zu-viele-Transformationen
Wenn aus einer Eingabe mehrere Ausgaben entstehen sollen, verwenden Sie mapConcat. Der Operator erwartet eine Funktion, die ein iterierbares Ergebnis zurückgibt, und flacht die Ergebnisse zu einem Stream ab.
Das Zurückgeben einer leeren Sammlung verwirft das Element effektiv.
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 und sliding
grouped(n) bündelt aufeinanderfolgende Elemente in einer Seq mit bis zu n Elementen, was beispielsweise für umfangreiche Datenbankschreibvorgänge nützlich ist. sliding(n) gibt sich überlappende Fenster aus.
Das Bündeln verringert den Aufwand pro Element in I/O-intensiven Pipelines.
val batches: Source[Seq[Int], akka.NotUsed] =
Source(1 to 1000).grouped(100)
val windows =
Source(1 to 10).sliding(3, step = 1)scan und fold
scan gibt nach jedem Element den aktuellen Akkumulator aus und erzeugt so einen sich verändernden Zustandsstream. fold gibt erst dann einmalig den endgültig akkumulierten Wert aus, wenn der Upstream abgeschlossen ist.
Verwenden Sie scan für Live-Zähler und fold für abschließende Aggregate.
val running =
Source(1 to 5).scan(0)(_ + _) // 0,1,3,6,10,15
val total =
Source(1 to 5).fold(0)(_ + _) // 15mapAsync für asynchrone Verarbeitung
mapAsync(parallelism) ruft eine Funktion auf, die einen Future zurückgibt, und gibt die Ergebnisse in der richtigen Reihenfolge aus. Dabei werden bis zu parallelism Futures gleichzeitig ausgeführt.
Verwenden Sie den Operator für asynchrone Aufrufe wie Datenbankabfragen oder HTTP-Anfragen, wenn die Reihenfolge wichtig ist.
import scala.concurrent.Future
val enriched =
Flow[UserId]
.mapAsync(parallelism = 4)(id => lookup(id))
def lookup(id: UserId): Future[User] = ???mapAsyncUnordered
mapAsyncUnordered verhält sich wie mapAsync, gibt jedoch jedes Ergebnis sofort nach dessen Fertigstellung aus und ignoriert dabei die Reihenfolge der Eingaben.
Das kann den Durchsatz erhöhen, wenn die nachgelagerte Verarbeitung keine bestimmte Reihenfolge benötigt, da ein langsamer Future schnellere Ergebnisse nicht mehr blockiert.
val fast =
Flow[UserId]
.mapAsyncUnordered(parallelism = 8)(id => lookup(id))Zustandsbehaftete Transformation mit statefulMapConcat
Für elementweise Transformationen, die einen veränderlichen lokalen Zustand benötigen, erstellt statefulMapConcat pro Materialisierung einen neuen Zustand und gibt ein iterierbares Ergebnis von Ausgaben zurück.
Damit können Sie Zähler oder Puffer sicher verwalten, ohne den Zustand zwischen verschiedenen Stream-Ausführungen zu teilen.
val withIndex: Flow[String, (Int, String), akka.NotUsed] =
Flow[String].statefulMapConcat { () =>
var i = 0
elem => { i += 1; List((i, elem)) }
}Zeitbasierte Operatoren
Streams können sowohl anhand der Zeit als auch anhand ihres Inhalts transformiert werden. throttle begrenzt die Ausgaberate, groupedWithin bündelt nach Größe oder verstrichener Zeit, und takeWithin begrenzt die Dauer.
Diese Operatoren sind für die Begrenzung der Rate externer APIs unverzichtbar.
import scala.concurrent.duration._
val limited =
Source(1 to 1000)
.throttle(10, 1.second)
.groupedWithin(100, 500.millis)Fehlerbehandlung bei Transformationen
Eine Ausnahme innerhalb eines Operators schlägt standardmäßig im gesamten Stream fehl. Eine Supervisionsstrategie kann stattdessen resume verwenden, um das fehlerhafte Element zu verwerfen, oder die Stufe mit restart neu starten.
Hängen Sie die Strategie mit withAttributes an den Flow an.
import akka.stream.{ActorAttributes, Supervision}
val safe =
Flow[String].map(_.toInt)
.withAttributes(
ActorAttributes.supervisionStrategy(_ => Supervision.Resume))Flows zusammenstellen
Kleine Flows lassen sich mit via zu größeren Flows zusammensetzen, wodurch ein einzelner wiederverwendbarer Flow entsteht. So bleibt jede Transformation überschaubar und kann unabhängig getestet werden.
Der zusammengesetzte Flow hat den Eingabetyp des ersten und den Ausgabetyp des letzten Flows.
val parse = Flow[String].map(_.toInt)
val square = Flow[Int].map(n => n * n)
val parseAndSquare: Flow[String, Int, akka.NotUsed] =
parse.via(square)Schnelltest
Betrachten Sie asynchrone Transformationen und ihre Garantien bezüglich der Reihenfolge.
Zusammenfassung
Sie haben Transformationsoperatoren kennengelernt: elementweise map/filter-Transformationen, Eins-zu-viele-Transformationen mit mapConcat, Bündelung mit grouped, Akkumulation mit scan und fold sowie asynchrone Verarbeitung mit mapAsync.
Außerdem haben Sie zustandsbehaftete Transformationen, zeitbasierte Operatoren wie throttle, Supervisionsstrategien für Fehler und das Zusammensetzen von Flows mit via kennengelernt. Als Nächstes geht es darum, wie Backpressure diese Stufen sicher hält.
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 „Streams transformieren“ kostenlos?
Ja — der vollständige Text von „Streams transformieren“ 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 „Streams transformieren“?
Mappen und filtern Sie fließende Daten. 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 2 von 4.
Wie lange dauert die Lektion „Streams transformieren“?
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