Eine Pipeline ausführen
Materialisieren und führen Sie einen Graphen aus.
Eine Pipeline ausführen ist eine kostenlose Scala for Backend Engineering & Functional Programming-Lektion auf CoddyKit. Dies ist Lektion 4 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.
Vom Bauplan zur Ausführung
Bisher war die Pipeline ein reiner Bauplan. Materialisierung bezeichnet den Vorgang, bei dem dieser Bauplan in laufende Actors umgewandelt wird, die tatsächlich Daten verarbeiten.
Es geschieht nichts, bis Sie den Graphen ausdrücklich ausführen. Genau das macht Akka Streams gut zusammensetzbar und wiederverwendbar.
Das ActorSystem
Für die Materialisierung ist ein ActorSystem erforderlich. Es stellt die Threads und den Dispatcher bereit, auf denen die Stufen des Streams ausgeführt werden. In modernem Akka dient das System außerdem als impliziter Materializer.
Ein ActorSystem versorgt normalerweise eine gesamte Anwendung und viele gleichzeitig laufende Streams.
import akka.actor.ActorSystem
implicit val system: ActorSystem =
ActorSystem("data-pipeline")
import system.dispatcher // ExecutionContextrunWith
Die direkteste Möglichkeit, eine Source auszuführen, ist runWith. Dabei wird ein Sink angehängt und die Materialisierung in einem Schritt durchgeführt; zurückgegeben wird der materialisierte Wert dieses Sinks.
Hier ist das Ergebnis ein Future[Int], der beim Abschluss des Streams mit der Summe abgeschlossen wird.
import akka.stream.scaladsl.{Source, Sink}
import scala.concurrent.Future
val total: Future[Int] =
Source(1 to 100).runWith(Sink.fold(0)(_ + _))Eine RunnableGraph ausführen
Wenn Sie bereits einen geschlossenen RunnableGraph mit to oder toMat erstellt haben, rufen Sie run() auf, um ihn zu materialisieren. Der Rückgabewert ist der materialisierte Wert, den der Graph beibehalten hat.
Dadurch werden Erstellung und Ausführung der Pipeline sauber voneinander getrennt.
import akka.stream.scaladsl.Keep
import scala.concurrent.Future
val graph =
Source(1 to 100)
.toMat(Sink.fold(0)(_ + _))(Keep.right)
val result: Future[Int] = graph.run()Praktische run-Operatoren
Sources bieten Kurzformen an: runForeach, runFold und runReduce hängen jeweils den entsprechenden Sink an und führen den Stream sofort aus.
Für gängige abschließende Operationen auf einer Source sind sie besonders kompakt.
import scala.concurrent.Future
val printed: Future[akka.Done] =
Source(1 to 10).runForeach(println)
val sum: Future[Int] =
Source(1 to 10).runFold(0)(_ + _)Mit dem Ergebnis-Future arbeiten
Abschließende Sinks geben einen Future zurück, der abgeschlossen wird, wenn der Stream endet oder fehlschlägt. Registrieren Sie mit onComplete Callbacks, um auf Erfolg oder Fehler zu reagieren.
Verwenden Sie den Dispatcher des ActorSystems als impliziten ExecutionContext für diese Callbacks.
import scala.util.{Success, Failure}
total.onComplete {
case Success(value) => println(s"Sum = $value")
case Failure(ex) => println(s"Failed: ${ex.getMessage}")
}Eine realistische Pipeline
Eine typische Datenpipeline liest aus einer Source, transformiert mit Flows, führt asynchrone I/O mit mapAsync aus, bündelt mit grouped und schreibt in einen Sink.
Jede Stufe ist klein, und die gesamte Pipeline wird mit einem einzigen run materialisiert.
val done =
lineSource
.map(parse)
.mapAsync(4)(validate)
.grouped(500)
.runWith(bulkWriteSink)Fehlgeschlagene Streams neu starten
Für mehr Resilienz umschließen Sie eine Source oder einen Flow mit RestartSource.withBackoff, sodass vorübergehende Fehler (z. B. eine unterbrochene Verbindung) automatisch einen Neustart mit exponentiellem Backoff auslösen.
So bleiben langlebige Ingestion-Pipelines ohne manuelle Überwachung aktiv.
import akka.stream.scaladsl.RestartSource
import akka.stream.RestartSettings
import scala.concurrent.duration._
val resilient = RestartSource.withBackoff(
RestartSettings(1.second, 30.seconds, 0.2))(() => flakySource)Geordnetes Herunterfahren mit KillSwitch
Ein KillSwitch ermöglicht es externem Code, einen laufenden Stream kontrolliert zu beenden. Fügen Sie KillSwitches.single über viaMat ein und speichern Sie den materialisierten Wert, um später shutdown() aufzurufen.
Das ist für langlebige Streams unerlässlich, die beim Herunterfahren der Anwendung beendet werden müssen.
import akka.stream.{KillSwitches, KillSwitch}
import akka.stream.scaladsl.Keep
val (switch, done) =
source
.viaMat(KillSwitches.single)(Keep.right)
.toMat(Sink.ignore)(Keep.both)
.run()
// later: switch.shutdown()Ressourcen freigeben
Beenden Sie beim Beenden der Anwendung das ActorSystem, um seine Threads freizugeben. Verketten Sie den Aufruf von terminate mit dem Completion-Future des Streams, damit das Herunterfahren geordnet abläuft.
Ein nicht freigegebenes ActorSystem hält die JVM am Leben und lässt Ressourcen weiterhin geöffnet.
done.onComplete { _ =>
system.terminate()
}Den Materializer wiederverwenden
Wenn Sie denselben Blueprint mehrfach materialisieren, entstehen unabhängige laufende Streams, die sich die Ressourcen des ActorSystem teilen. Der Blueprint selbst bleibt unveränderlich und frei von Seiteneffekten.
Dadurch können Sie eine Pipeline sicher einmal definieren und sie für jeden eingehenden Auftrag bei Bedarf ausführen.
val blueprint =
Source(1 to 5).toMat(Sink.seq)(Keep.right)
val run1 = blueprint.run()
val run2 = blueprint.run() // independent executionKurzer Test
Überlegen Sie, was erforderlich ist, damit ein Stream tatsächlich Elemente verarbeitet.
Zusammenfassung
Eine Pipeline auszuführen bedeutet, einen Blueprint mit einem ActorSystem über run, runWith oder Komfortoperatoren zu materialisieren, wobei jeweils ein Future-Ergebnis zurückgegeben wird.
Sie haben realistische Pipelines mit mehreren Stufen, automatische Neustarts mit Backoff, ein geordnetes Herunterfahren über KillSwitch, die Freigabe von Ressourcen mit system.terminate() sowie die sichere Wiederverwendung eines unveränderlichen Blueprints über unabhängige Ausführungen hinweg kennengelernt.
Häufig gestellte Fragen
Ist die Lektion „Eine Pipeline ausführen“ kostenlos?
Ja — der vollständige Text von „Eine Pipeline ausführen“ 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 „Eine Pipeline ausführen“?
Materialisieren und führen Sie einen Graphen aus. 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 4 von 4.
Wie lange dauert die Lektion „Eine Pipeline ausführen“?
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