Kör en pipeline
Materialisera och kör en graf.
Kör en pipeline är en gratis lektion i Scala för backendutveckling och funktionell programmering på CoddyKit. Detta är lektion 4 av 4. Ni kan läsa hela lektionen gratis nedan och sedan öva praktiskt i webbläsaren med en inbyggd kodredigerare och en AI-handledare som är tillgänglig dygnet runt. Den ingår i lärvägen för Scala för backendutveckling och funktionell programmering, och Era framsteg synkroniseras mellan webben och CoddyKit-appen. Kursen i Scala för backendutveckling och funktionell programmering innehåller totalt 4 lektioner.
Från ritning till körning
Hittills har pipelinen varit en ren ritning. Materialisering är processen som omvandlar ritningen till körande actors som faktiskt flyttar data.
Ingenting händer förrän Ni uttryckligen kör grafen, vilket gör Akka Streams komponerbart och återanvändbart.
ActorSystem
Materialisering kräver ett ActorSystem, som tillhandahåller trådarna och den dispatcher som driver strömmens steg. I moderna Akka fungerar systemet även som implicit materialiserare.
Ett ActorSystem betjänar vanligtvis en hel applikation och många samtidiga strömmar.
import akka.actor.ActorSystem
implicit val system: ActorSystem =
ActorSystem("data-pipeline")
import system.dispatcher // ExecutionContextrunWith
Det enklaste sättet att köra en Source är runWith. Det kopplar till en Sink och materialiserar i ett enda steg, och returnerar Sink-objektets materialiserade värde.
Här är resultatet en Future[Int] som slutförs med summan när strömmen är klar.
import akka.stream.scaladsl.{Source, Sink}
import scala.concurrent.Future
val total: Future[Int] =
Source(1 to 100).runWith(Sink.fold(0)(_ + _))Kör en RunnableGraph
Om Ni redan har byggt en sluten RunnableGraph med to eller toMat anropar Ni run() för att materialisera den. Returvärdet är det materialiserade värde som grafen behöll.
Detta skiljer på ett tydligt sätt konstruktionen av pipelinen och körningen.
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()Praktiska run-operatorer
Sources erbjuder genvägar: runForeach, runFold och runReduce kopplar var och en till motsvarande Sink och kör den direkt.
De är kortfattade alternativ för vanliga avslutande operationer på en 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)(_ + _)Arbeta med resultatets Future
Avslutande sinks returnerar en Future som slutförs när strömmen avslutas eller misslyckas. Registrera återanrop med onComplete för att reagera på lyckat resultat eller fel.
Använd ActorSystem:s dispatcher som implicit ExecutionContext för dessa återanrop.
import scala.util.{Success, Failure}
total.onComplete {
case Success(value) => println(s"Sum = $value")
case Failure(ex) => println(s"Failed: ${ex.getMessage}")
}En realistisk pipeline
En typisk datapipeline läser från en Source, transformerar med Flows, utför asynkron I/O med mapAsync, grupperar i batchar med grouped och skriver till en Sink.
Varje steg är litet, och hela pipelinen materialiseras med ett enda run.
val done =
lineSource
.map(parse)
.mapAsync(4)(validate)
.grouped(500)
.runWith(bulkWriteSink)Starta om misslyckade strömmar
För bättre feltålighet kan Ni omsluta en Source eller Flow med RestartSource.withBackoff, så att tillfälliga fel, till exempel en bruten anslutning, utlöser en automatisk omstart med exponentiell backoff.
Detta håller långkörande inläsningspipelines igång utan manuell övervakning.
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)Smidig avstängning med KillSwitch
En KillSwitch låter extern kod stoppa en pågående stream på ett ordnat sätt. Infoga KillSwitches.single via viaMat och spara dess materialiserade värde för att senare anropa shutdown().
Detta är viktigt för långlivade strömmar som måste stoppas när applikationen stängs av.
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()Frigöra resurser
När applikationen avslutas ska Ni terminera ActorSystem för att frigöra dess trådar. Kedja anropet till terminate efter streamens completion Future så att avstängningen sker ordnat.
Om Ni läcker ett ActorSystem hålls JVM:en igång och resurser förblir öppna.
done.onComplete { _ =>
system.terminate()
}Återanvända Materializer
Om samma blueprint materialiseras flera gånger skapas oberoende körande strömmar som delar ActorSystems resurser. Själva blueprinten förblir oföränderlig och fri från sidoeffekter.
Det gör det säkert att definiera en pipeline en gång och köra den vid behov för varje inkommande jobb.
val blueprint =
Source(1 to 5).toMat(Sink.seq)(Keep.right)
val run1 = blueprint.run()
val run2 = blueprint.run() // independent executionSnabbtest
Fundera på vad som krävs för att en stream faktiskt ska bearbeta element.
Sammanfattning
Att köra en pipeline innebär att materialisera en blueprint med ett ActorSystem via run, runWith eller bekvämlighetsoperatorer, där var och en returnerar ett Future-resultat.
Ni har sett realistiska pipelines med flera steg, automatiska omstarter med backoff, smidig avstängning via KillSwitch, resursfrigöring med system.terminate() samt säker återanvändning av en oföränderlig blueprint i oberoende körningar.
Lär dig Scala med en AI-lärare – gratis
Skriv och kör riktig kod i webbläsaren, få omedelbar hjälp av en AI-lärare dygnet runt och fortsätt där du slutade – på webben eller i appen.
- Kurser
- 39
- Lektioner
- 143
Vanliga frågor
Är lektionen ”Kör en pipeline” gratis?
Ja – hela texten till ”Kör en pipeline” kan läsas gratis här på webben. Om Ni vill öva interaktivt med en inbyggd kodredigerare och en AI-handledare som är tillgänglig dygnet runt och låsa upp resten av kursen i Scala för backendutveckling och funktionell programmering, kan Ni uppgradera till CoddyKit PRO. Kursen i Scala för backendutveckling och funktionell programmering innehåller totalt 4 lektioner.
Vad lär jag mig i ”Kör en pipeline”?
Materialisera och kör en graf. Ni övar på Scala för backendutveckling och funktionell programmering med praktisk kod som körs direkt i webbläsaren, medan en AI-handledare som är tillgänglig dygnet runt svarar på Era frågor under lektionen.
Behöver jag någon erfarenhet för att börja lära mig Scala för backendutveckling och funktionell programmering?
Du behöver inga förkunskaper. Utbildningen i Scala för backendutveckling och funktionell programmering på CoddyKit är upplagd för allt från nybörjare till avancerade elever, så att du kan börja här eller från början och gå fram i din egen takt. Detta är lektion 4 av 4.
Hur lång tid tar lektionen ”Kör en pipeline”?
De flesta CoddyKit-lektioner tar cirka 5–10 minuter. Varje lektion är kort och interaktiv, så att du gör stadiga framsteg och kan fortsätta precis där du slutade – på webben eller i appen.
Kan jag skriva och köra kod i den här Scala för backendutveckling och funktionell programmering-lektionen?
Ja. Varje Scala för backendutveckling och funktionell programmering-lektion innehåller en inbyggd kodredigerare, så att du kan skriva och köra riktig kod direkt i webbläsaren och få omedelbar AI-feedback – utan lokal installation.
Alla lektioner i den här kursen
- Source, Flow och Sink
- Omforma strömmar
- Backpressure
- Kör en pipeline