Scala for backendutvikling og funksjonell programmering · leksjon

Kjør en pipeline

Materialiser og kjør en graf.

Leksjon 4 av 413 trinn

Kjør en pipeline er en gratis leksjon i Scala for backendutvikling og funksjonell programmering på CoddyKit. Dette er leksjon 4 av 4. Du kan lese hele leksjonen gratis nedenfor – og deretter øve praktisk i nettleseren med en innebygd kodeeditor og en AI-veileder som er tilgjengelig døgnet rundt. Den er en del av læringsløpet i Scala for backendutvikling og funksjonell programmering, og fremdriften din synkroniseres mellom nettet og CoddyKit-appen. Kurset i Scala for backendutvikling og funksjonell programmering inneholder totalt 4 leksjoner.

Fra blueprint til kjøring

Så langt har pipelinen vært en ren blueprint. Materialisering er prosessen som gjør denne blueprint-en om til kjørende aktører som faktisk flytter data.

Ingenting skjer før De eksplisitt kjører grafen, og det er dette som gjør Akka Streams mulig å sette sammen og gjenbruke.

ActorSystem

Materialisering krever et ActorSystem, som tilbyr trådene og dispatcheren som driver trinnene i strømmen. I moderne Akka fungerer systemet også som den implisitte materialisatoren.

Ett ActorSystem betjener vanligvis en hel applikasjon og mange samtidige strømmer.

import akka.actor.ActorSystem

implicit val system: ActorSystem =
  ActorSystem("data-pipeline")
import system.dispatcher // ExecutionContext

runWith

Den mest direkte måten å kjøre en Source på er runWith. Den kobler til en Sink og materialiserer i ett trinn, og returnerer Sink-ens materialiserte verdi.

Her er resultatet en Future[Int] som fullføres med summen når strømmen er ferdig.

import akka.stream.scaladsl.{Source, Sink}
import scala.concurrent.Future

val total: Future[Int] =
  Source(1 to 100).runWith(Sink.fold(0)(_ + _))

Kjøre en RunnableGraph

Hvis De allerede har bygget en lukket RunnableGraph med to eller toMat, kaller De run() for å materialisere den. Returverdien er den materialiserte verdien grafen beholdt.

Dette skiller konstruksjon av pipelinen tydelig fra kjøringen.

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()

Praktiske run-operatorer

Sources tilbyr snarveier: runForeach, runFold og runReduce kobler hver til den tilsvarende Sink-en og kjører umiddelbart.

De er kortfattede ved vanlige terminaloperasjoner 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)(_ + _)

Arbeide med resultatets Future

Terminale sinks returnerer en Future som fullføres når strømmen avsluttes eller feiler. Registrer tilbakekallinger med onComplete for å reagere på suksess eller feil.

Bruk ActorSystem-ets dispatcher som implisitt ExecutionContext for disse tilbakekallingene.

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 leser fra en Source, transformerer med Flows, utfører asynkron I/O med mapAsync, samler i batcher med grouped og skriver til en Sink.

Hvert trinn er lite, og hele pipelinen materialiseres med én enkelt run.

val done =
  lineSource
    .map(parse)
    .mapAsync(4)(validate)
    .grouped(500)
    .runWith(bulkWriteSink)

Starte mislykkede strømmer på nytt

For bedre robusthet kan De pakke inn en Source eller Flow med RestartSource.withBackoff, slik at midlertidige feil (for eksempel en brutt forbindelse) utløser en automatisk omstart med eksponentiell backoff.

Dette holder langvarige datapipelines i gang uten manuell overvåking.

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)

Kontrollert avslutning med KillSwitch

En KillSwitch lar ekstern kode stoppe en kjørende strøm på en ryddig måte. Sett inn KillSwitches.single via viaMat, og ta vare på den materialiserte verdien for å kunne kalle shutdown() senere.

Dette er avgjørende for langvarige strømmer som må stoppe når applikasjonen avsluttes.

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()

Frigjøre ressurser

Når applikasjonen avsluttes, terminerer De ActorSystem for å frigjøre trådene. Koble terminate-kallet etter streamens completion Future, slik at avslutningen skjer på en ryddig måte.

Hvis en ActorSystem lekker, holdes JVM-en i gang og ressurser forblir åpne.

done.onComplete { _ =>
  system.terminate()
}

Gjenbruke Materializer

Når den samme blueprinten materialiseres flere ganger, opprettes uavhengige kjørende strømmer som deler ActorSystem-ens ressurser. Blueprinten selv forblir uforanderlig og uten sideeffekter.

Dette gjør det trygt å definere en pipeline én gang og kjøre den ved behov for hver innkommende jobb.

val blueprint =
  Source(1 to 5).toMat(Sink.seq)(Keep.right)

val run1 = blueprint.run()
val run2 = blueprint.run() // independent execution

Rask kontroll

Vurder hva som kreves for at en strøm faktisk skal behandle elementer.

Oppsummering

Å kjøre en pipeline innebærer å materialisere en blueprint med et ActorSystem via run, runWith eller bekvemmelighetsoperatorer, der hver av disse returnerer et Future-resultat.

Her gjennomgås realistiske pipelines med flere trinn, automatiske omstarter med backoff, kontrollert avslutning via KillSwitch, opprydding av ressurser med system.terminate() og trygg gjenbruk av en uforanderlig blueprint på tvers av uavhengige kjøringer.

Gratis å komme i gang

Lær deg Scala med en AI-veileder – gratis

Skriv og kjør ekte kode i nettleseren, få umiddelbar hjelp fra en AI-veileder som er tilgjengelig døgnet rundt, og fortsett der du slapp – på nettet eller i appen.

Kurs
39
Leksjoner
143

Ofte stilte spørsmål

Er leksjonen «Kjør en pipeline» gratis?

Ja – hele teksten i «Kjør en pipeline» er gratis å lese her på nettet. For å øve interaktivt med en innebygd kodeeditor og en AI-veileder som er tilgjengelig døgnet rundt, og for å låse opp resten av Scala for backendutvikling og funksjonell programmering-kurset, kan du oppgradere til CoddyKit PRO. Kurset i Scala for backendutvikling og funksjonell programmering inneholder totalt 4 leksjoner.

Hva lærer jeg i «Kjør en pipeline»?

Materialiser og kjør en graf. Du øver på Scala for backendutvikling og funksjonell programmering med praktisk kode som du kjører direkte i nettleseren, mens en AI-veileder som er tilgjengelig døgnet rundt, svarer på spørsmålene dine mens du jobber deg gjennom leksjonen.

Trenger jeg erfaring for å begynne med Scala for backendutvikling og funksjonell programmering?

Ingen tidligere erfaring er nødvendig. Scala for backendutvikling og funksjonell programmering på CoddyKit er lagt opp for både nybegynnere og viderekomne, så De kan begynne her eller helt fra start og lære i Deres eget tempo. Dette er leksjon 4 av 4.

Hvor lang tid tar leksjonen «Kjør en pipeline»?

De fleste CoddyKit-leksjoner tar omtrent 5–10 minutter. Hver leksjon er kort og interaktiv, slik at De gjør jevne fremskritt og kan fortsette akkurat der De slapp – både på nettet og i appen.

Kan jeg skrive og kjøre kode i denne Scala for backendutvikling og funksjonell programmering-leksjonen?

Ja. Alle Scala for backendutvikling og funksjonell programmering-leksjoner har en innebygd kodeeditor, slik at De kan skrive og kjøre ekte kode direkte i nettleseren og få umiddelbar tilbakemelding fra AI – uten lokal konfigurering.

Alle leksjonene i dette kurset

  1. Source, Flow og Sink
  2. Transformer strømmer
  3. Mottrykk
  4. Kjør en pipeline
← Tilbake til Scala for backendutvikling og funksjonell programmering