Een pipeline uitvoeren
Materialiseer en voer een graaf uit.
Een pipeline uitvoeren is een gratis Scala voor backend-engineering en functioneel programmeren-les op CoddyKit. Dit is les 4 van 4. Je kunt de volledige les hieronder gratis lezen en daarna in de browser praktisch oefenen met een ingebouwde code-editor en een AI-begeleider die 24/7 beschikbaar is. Deze les maakt deel uit van het leertraject Scala voor backend-engineering en functioneel programmeren. Je voortgang wordt gesynchroniseerd op het web en in de CoddyKit-app. De cursus Scala voor backend-engineering en functioneel programmeren bevat in totaal 4 lessen.
Van blauwdruk naar uitvoering
Tot nu toe was de pijplijn een zuivere blauwdruk. Materialisatie is het proces waarbij die blauwdruk wordt omgezet in actieve actors die daadwerkelijk gegevens verplaatsen.
Er gebeurt niets totdat u de grafiek expliciet uitvoert. Dat maakt Akka Streams samenstelbaar en herbruikbaar.
Het ActorSystem
Voor materialisatie is een ActorSystem nodig. Dit levert de threads en dispatcher waarop de fasen van de stream draaien. In modern Akka fungeert het systeem ook als de impliciete materialisator.
Eén ActorSystem bedient doorgaans een volledige applicatie en veel gelijktijdige streams.
import akka.actor.ActorSystem
implicit val system: ActorSystem =
ActorSystem("data-pipeline")
import system.dispatcher // ExecutionContextrunWith
De meest directe manier om een Source uit te voeren is runWith. Hiermee koppelt u een Sink en materialiseert u in één stap, waarbij de gematerialiseerde waarde van die Sink wordt geretourneerd.
Hier is het resultaat een Future[Int] die wordt voltooid met de som wanneer de stream eindigt.
import akka.stream.scaladsl.{Source, Sink}
import scala.concurrent.Future
val total: Future[Int] =
Source(1 to 100).runWith(Sink.fold(0)(_ + _))Een RunnableGraph uitvoeren
Als u al een gesloten RunnableGraph met to of toMat hebt opgebouwd, roept u run() aan om deze te materialiseren. De retourwaarde is de gematerialiseerde waarde die de grafiek heeft behouden.
Zo worden het opbouwen en uitvoeren van de pijplijn netjes van elkaar gescheiden.
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()Handige uitvoeringsoperatoren
Sources bieden snelkoppelingen: runForeach, runFold en runReduce koppelen elk de bijbehorende Sink en voeren deze onmiddellijk uit.
Ze zijn beknopt voor veelgebruikte eindbewerkingen op een Source.
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)(_ + _)Werken met de resultaat-Future
Eindsinks retourneren een Future die wordt voltooid wanneer de stream eindigt of mislukt. Registreer terugbelmethoden met onComplete om op succes of een fout te reageren.
Gebruik de dispatcher van het ActorSystem als de impliciete ExecutionContext voor deze terugbelmethoden.
import scala.util.{Success, Failure}
total.onComplete {
case Success(value) => println(s"Sum = $value")
case Failure(ex) => println(s"Failed: ${ex.getMessage}")
}Een realistische pijplijn
Een typische gegevenspijplijn leest uit een Source, transformeert met Flows, voert asynchrone I/O uit met mapAsync, bundelt elementen in groepen met grouped en schrijft naar een Sink.
Elke fase is klein en het geheel wordt met één run gematerialiseerd.
val done =
lineSource
.map(parse)
.mapAsync(4)(validate)
.grouped(500)
.runWith(bulkWriteSink)Mislukte streams opnieuw starten
Voor veerkracht wikkel je een Source of Flow in met RestartSource.withBackoff, zodat tijdelijke fouten (zoals een verbroken verbinding) automatisch een herstart met exponentieel oplopende wachttijd activeren.
Zo blijven langlopende gegevensinnamepijplijnen actief zonder handmatige bewaking.
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)Gecontroleerd afsluiten met KillSwitch
Met een KillSwitch kan externe code een actieve stream netjes stoppen. Voeg KillSwitches.single in via viaMat en bewaar de gematerialiseerde waarde om later shutdown() aan te roepen.
Dit is essentieel voor langlevende streams die moeten stoppen wanneer de toepassing wordt afgesloten.
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()Resources vrijgeven
Wanneer de toepassing wordt afgesloten, beëindig je het ActorSystem om de threads ervan vrij te geven. Koppel de aanroep van terminate aan de Future voor de voltooiing van de stream, zodat het afsluiten ordelijk verloopt.
Een niet-vrijgegeven ActorSystem houdt de JVM actief en laat systeembronnen openstaan.
done.onComplete { _ =>
system.terminate()
}De Materializer hergebruiken
Als je dezelfde blauwdruk meerdere keren materialiseert, ontstaan onafhankelijke actieve streams die de resources van het ActorSystem delen. De blauwdruk zelf blijft onveranderlijk en vrij van bijwerkingen.
Zo kun je veilig één keer een pijplijn definiëren en die op verzoek uitvoeren voor elke binnenkomende taak.
val blueprint =
Source(1 to 5).toMat(Sink.seq)(Keep.right)
val run1 = blueprint.run()
val run2 = blueprint.run() // independent executionKorte controle
Bedenk wat nodig is om een stream daadwerkelijk elementen te laten verwerken.
Samenvatting
Een pijplijn uitvoeren betekent dat je met een ActorSystem een blauwdruk materialiseert via run, runWith of handige operatoren, die elk een Future-resultaat opleveren.
Je hebt realistische pijplijnen met meerdere fasen gezien, automatische herstarts met oplopende wachttijd, gecontroleerd afsluiten via KillSwitch, het vrijgeven van systeembronnen met system.terminate() en het veilig hergebruiken van een onveranderlijke blauwdruk in onafhankelijke uitvoeringen.
Leer Scala met een AI-tutor — gratis
Schrijf echte code en voer die uit in je browser, krijg direct hulp van een AI-tutor die 24/7 beschikbaar is en ga verder waar je gebleven bent op het web of in de app.
- Cursussen
- 39
- Lessen
- 143
Veelgestelde vragen
Is de les “Een pipeline uitvoeren” gratis?
Ja — de volledige tekst van “Een pipeline uitvoeren” kun je hier gratis op het web lezen. Als je interactief wilt oefenen met een ingebouwde code-editor en een AI-begeleider die 24/7 beschikbaar is, en de rest van de cursus Scala voor backend-engineering en functioneel programmeren wilt ontgrendelen, kun je upgraden naar CoddyKit PRO. De cursus Scala voor backend-engineering en functioneel programmeren bevat in totaal 4 lessen.
Wat leer ik in “Een pipeline uitvoeren”?
Materialiseer en voer een graaf uit. Je oefent met Scala voor backend-engineering en functioneel programmeren door code rechtstreeks in de browser uit te voeren. Een AI-begeleider die 24/7 beschikbaar is beantwoordt je vragen terwijl je de les doorwerkt.
Heb ik ervaring nodig om met Scala voor backend-engineering en functioneel programmeren te beginnen?
Ervaring vooraf is niet nodig. Scala voor backend-engineering en functioneel programmeren op CoddyKit is opgebouwd voor beginners tot gevorderden, zodat je hier of bij het begin kunt starten en in je eigen tempo kunt leren. Dit is les 4 van 4.
Hoe lang duurt de les “Een pipeline uitvoeren”?
De meeste lessen van CoddyKit duren ongeveer 5–10 minuten. Elke les is kort en interactief, zodat je gestaag vooruitgaat en op het web en in de app precies verdergaat waar je was gebleven.
Kan ik code schrijven en uitvoeren in deze les over Scala voor backend-engineering en functioneel programmeren?
Ja. Elke les over Scala voor backend-engineering en functioneel programmeren bevat een ingebouwde code-editor, zodat je rechtstreeks in je browser echte code kunt schrijven en uitvoeren en direct feedback van AI krijgt — lokale installatie is niet nodig.
Alle lessen in deze cursus
- Source, flow en sink
- Streams transformeren
- Backpressure
- Een pipeline uitvoeren