Uruchamianie potoku
Zmaterializuje Pan/Pani graf i wykona go.
Uruchamianie potoku to bezpłatna lekcja Scala for Backend Engineering & Functional Programming na CoddyKit. To lekcja 4 z 4. Możesz przeczytać całą lekcję poniżej za darmo — a potem ćwiczyć ją interaktywnie w przeglądarce z wbudowanym edytorem kodu i tutorem AI dostępnym 24/7. To część ścieżki edukacyjnej Scala for Backend Engineering & Functional Programming, a Twój postęp synchronizuje się między webem a aplikacją CoddyKit. Kurs Scala for Backend Engineering & Functional Programming zawiera 4 lekcji w sumie.
Od blueprintu do wykonania
Do tej pory potok był wyłącznie blueprintem. Materializacja to proces przekształcania tego blueprintu w działających aktorów, którzy faktycznie przesyłają dane.
Nic się nie dzieje, dopóki jawnie nie uruchomią Państwo grafu — właśnie to sprawia, że Akka Streams można komponować i wielokrotnie wykorzystywać.
ActorSystem
Materializacja wymaga obiektu ActorSystem, który udostępnia wątki i dispatcher obsługujące etapy strumienia. We współczesnej Akka system pełni również rolę niejawnego materializera.
Jeden ActorSystem zwykle obsługuje całą aplikację i wiele współbieżnych strumieni.
import akka.actor.ActorSystem
implicit val system: ActorSystem =
ActorSystem("data-pipeline")
import system.dispatcher // ExecutionContextrunWith
Najbardziej bezpośrednim sposobem uruchomienia Source jest runWith, które dołącza Sink i wykonuje materializację w jednym kroku, zwracając wartość materializowaną tego Sink.
W tym przypadku wynikiem jest Future[Int], który zostaje ukończony sumą po zakończeniu strumienia.
import akka.stream.scaladsl.{Source, Sink}
import scala.concurrent.Future
val total: Future[Int] =
Source(1 to 100).runWith(Sink.fold(0)(_ + _))Uruchamianie RunnableGraph
Jeśli mają już Państwo zbudowany zamknięty obiekt RunnableGraph za pomocą to lub toMat, należy wywołać run(), aby go zmaterializować. Zwracana wartość jest dowolną wartością materializowaną zachowaną przez graf.
W ten sposób konstrukcja potoku i jego wykonanie są wyraźnie rozdzielone.
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()Wygodne operatory uruchamiające
Sources udostępniają skróty: runForeach, runFold i runReduce dołączają odpowiedni Sink i od razu uruchamiają strumień.
Są zwięzłe i przydatne przy typowych operacjach końcowych na 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)(_ + _)Praca z wynikowym obiektem Future
Końcowe Sinki zwracają obiekt Future, który zostaje ukończony po zakończeniu strumienia lub w przypadku błędu. Zarejestruj funkcje zwrotne za pomocą onComplete, aby reagować na powodzenie lub błąd.
Użyj dispatchera obiektu ActorSystem jako niejawnego ExecutionContext dla tych funkcji zwrotnych.
import scala.util.{Success, Failure}
total.onComplete {
case Success(value) => println(s"Sum = $value")
case Failure(ex) => println(s"Failed: ${ex.getMessage}")
}Realistyczny potok
Typowy potok danych odczytuje dane ze źródła Source, przekształca je za pomocą Flows, wykonuje asynchroniczne operacje wejścia-wyjścia przy użyciu mapAsync, grupuje je za pomocą grouped i zapisuje w Sink.
Każdy etap jest niewielki, a całość można zmaterializować jednym wywołaniem run.
val done =
lineSource
.map(parse)
.mapAsync(4)(validate)
.grouped(500)
.runWith(bulkWriteSink)Ponowne uruchamianie strumieni zakończonych błędem
Aby zwiększyć odporność, należy opakować Source lub Flow za pomocą RestartSource.withBackoff, tak aby przejściowe awarie (na przykład zerwanie połączenia) powodowały automatyczne ponowne uruchomienie z wykładniczym opóźnieniem.
Dzięki temu długotrwałe potoki pozyskiwania danych mogą działać bez ręcznego nadzorowania.
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)Kontrolowane zamykanie za pomocą KillSwitch
KillSwitch pozwala zewnętrznemu kodowi zatrzymać działający strumień w uporządkowany sposób. Należy wstawić KillSwitches.single za pomocą viaMat i zachować jego zmaterializowaną wartość, aby później wywołać shutdown().
Jest to niezbędne w przypadku długotrwałych strumieni, które muszą zostać zatrzymane podczas zamykania aplikacji.
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()Zwalnianie zasobów
Po zakończeniu działania aplikacji należy zakończyć ActorSystem, aby zwolnić jego wątki. Należy połączyć wywołanie terminate z Future zakończenia strumienia, aby zamykanie przebiegło w uporządkowany sposób.
Pozostawienie niezamkniętego ActorSystem utrzymuje JVM przy życiu i powoduje dalsze zajmowanie zasobów.
done.onComplete { _ =>
system.terminate()
}Ponowne używanie Materializer
Wielokrotne materializowanie tego samego schematu tworzy niezależne działające strumienie, które współdzielą zasoby ActorSystem. Sam schemat pozostaje niezmienny i pozbawiony efektów ubocznych.
Dzięki temu można bezpiecznie zdefiniować potok raz i uruchamiać go na żądanie dla każdego przychodzącego zadania.
val blueprint =
Source(1 to 5).toMat(Sink.seq)(Keep.right)
val run1 = blueprint.run()
val run2 = blueprint.run() // independent executionSzybkie sprawdzenie
Proszę zastanowić się, co jest wymagane, aby strumień rzeczywiście przetwarzał elementy.
Podsumowanie
Uruchomienie potoku oznacza zmaterializowanie schematu za pomocą ActorSystem i run, runWith lub operatorów pomocniczych, z których każdy zwraca wynik typu Future.
Omówiono realistyczne potoki wieloetapowe, automatyczne ponowne uruchamianie z wykładniczym opóźnieniem, kontrolowane zamykanie za pomocą KillSwitch, zwalnianie zasobów przy użyciu system.terminate() oraz bezpieczne ponowne używanie niezmiennego schematu w niezależnych uruchomieniach.
Ucz się Scala dzięki korepetycjom AI — za darmo
Pisz i uruchamiaj kod w przeglądarce, otrzymuj natychmiastową pomoc od korepetytora AI dostępnego 24/7 i kontynuuj naukę w sieci lub w aplikacji.
- Kursy
- 39
- Lekcje
- 143
Często zadawane pytania
Czy lekcja „Uruchamianie potoku” jest bezpłatna?
Tak — pełny tekst „Uruchamianie potoku” jest dostępny za darmo tutaj w sieci. Aby ćwiczyć ją interaktywnie (wbudowany edytor kodu i tutor AI dostępny 24/7) i odblokować resztę kursu Scala for Backend Engineering & Functional Programming, przejdź na CoddyKit PRO. Kurs Scala for Backend Engineering & Functional Programming zawiera 4 lekcji w sumie.
Co nauczysz się w „Uruchamianie potoku”?
Zmaterializuje Pan/Pani graf i wykona go. Ćwiczysz Scala for Backend Engineering & Functional Programming z praktycznym kodem, który uruchamiasz bezpośrednio w przeglądarce, a tutor AI dostępny 24/7 odpowiada na Twoje pytania podczas pracy nad lekcją.
Czy potrzebuję doświadczenia, aby zacząć Scala for Backend Engineering & Functional Programming?
Nie wymagamy żadnego doświadczenia. Scala for Backend Engineering & Functional Programming w CoddyKit jest strukturyzowany dla początkujących i zaawansowanych użytkowników, więc możesz zacząć tutaj lub od początku i uczyć się w swoim tempie. To lekcja 4 z 4.
Ile czasu zajmuje lekcja „Uruchamianie potoku”?
Większość lekcji CoddyKit trwa około 5–10 minut. Każda lekcja to mały, interaktywny krok, dzięki czemu robisz systematyczne postępy i zawsze wracasz dokładnie do tego samego miejsca — na webie i w aplikacji.
Czy mogę pisać i uruchamiać kod w tej lekcji Scala for Backend Engineering & Functional Programming?
Tak. Każda lekcja Scala for Backend Engineering & Functional Programming zawiera wbudowany edytor kodu, więc piszesz i uruchamiasz prawdziwy kod bezpośrednio w przeglądarce i od razu otrzymujesz sprzężenie zwrotne od AI — bez konfiguracji na komputerze.
Wszystkie lekcje w tym kursie
- Source, Flow i Sink
- Przekształcanie strumieni
- Backpressure
- Uruchamianie potoku