Source, Flow i Sink
Podstawowe elementy budowy strumieni.
Source, Flow i Sink to bezpłatna lekcja Scala for Backend Engineering & Functional Programming na CoddyKit. To lekcja 1 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.
Trzy podstawowe elementy
Akka Streams modeluje potok danych jako graf etapów przetwarzania. Trzy podstawowe etapy liniowe to Source (produkuje elementy), Flow (przekształca je) i Sink (konsumuje je).
Source ma jedno wyjście, Sink ma jedno wejście, a Flow ma dokładnie jedno wejście i jedno wyjście. Połączenie ich opisuje co ma się wydarzyć, a nie kiedy.
Definiowanie Source
Obiekt Source[Out, Mat] emituje elementy typu Out i udostępnia zmaterializowaną wartość typu Mat. Najprostsze źródła pochodzą z kolekcji przechowywanych w pamięci lub z zakresów.
Do momentu uruchomienia strumienia Source jest tylko niezmiennym szablonem, który można swobodnie ponownie wykorzystywać.
import akka.stream.scaladsl.Source
val numbers: Source[Int, akka.NotUsed] =
Source(1 to 100)
val single: Source[String, akka.NotUsed] =
Source.single("hello")Definiowanie Sink
Obiekt Sink[In, Mat] konsumuje elementy typu In. Zmaterializowana wartość często przechowuje wynik konsumowania, na przykład obiekt Future kończący się po zakończeniu strumienia.
Sink.foreach wykonuje efekt uboczny dla każdego elementu, a Sink.fold akumuluje pojedynczy wynik.
import akka.stream.scaladsl.Sink
import scala.concurrent.Future
val printSink: Sink[Int, Future[akka.Done]] =
Sink.foreach(println)
val sumSink: Sink[Int, Future[Int]] =
Sink.fold(0)(_ + _)Definiowanie Flow
Obiekt Flow[In, Out, Mat] znajduje się między Source i Sink, przekształcając każdy element. Elementy Flow można ponownie wykorzystywać samodzielnie i komponować przed dołączeniem do dowolnego punktu końcowego.
W tym przykładzie Flow podwaja liczby całkowite i konwertuje je na napisy.
import akka.stream.scaladsl.Flow
val doubleToString: Flow[Int, String, akka.NotUsed] =
Flow[Int]
.map(_ * 2)
.map(n => s"value=$n")Łączenie Source z Sink
Operator via dołącza Flow do Source, a to dołącza Sink. Bezpośrednie połączenie Source z Sink za pomocą to tworzy RunnableGraph — zamknięty szablon gotowy do uruchomienia.
Żadne elementy nie są jeszcze przesyłane; opisuje to jedynie topologię.
import akka.stream.scaladsl.{Source, Sink, RunnableGraph}
val graph: RunnableGraph[akka.NotUsed] =
Source(1 to 10).to(Sink.foreach(println))via: wstawianie Flow
Użyj via, aby włączyć Flow do potoku. Wywołanie Source.via(flow) zwraca nowy Source, którego typ wyjściowy odpowiada typowi wyjściowemu Flow.
Łączenie wywołań via pozwala budować długie potoki transformacji z małych, łatwych do testowania elementów Flow.
val pipeline =
Source(1 to 10)
.via(Flow[Int].filter(_ % 2 == 0))
.via(Flow[Int].map(_ * 10))
.to(Sink.foreach(println))Bezpieczeństwo typów między etapami
Kompilator wymusza zgodność typu wyjściowego każdego etapu z typem wejściowym kolejnego. Source[Int] nie może zostać połączony z Sink[String] bez pośredniego Flow, który konwertuje typ.
To sprawdzanie statyczne wykrywa błędy w połączeniach potoku jeszcze przed uruchomieniem programu.
// Source[Int] -> Flow[Int, String] -> Sink[String]
val ok =
Source(1 to 3)
.via(Flow[Int].map(_.toString))
.to(Sink.foreach[String](println))Wartości materializowane
Każdy blueprint zawiera wartość materializowaną — uchwyt tworzony podczas uruchamiania strumienia. Sink.fold materializuje obiekt Future zawierający wynik. Domyślnie łączenie etapów zachowuje skrajną lewą wartość materializowaną (NotUsed dla zwykłych źródeł).
Użyj toMat i Keep, aby wybrać wartość z odpowiedniej strony.
import akka.stream.scaladsl.Keep
import scala.concurrent.Future
val g: RunnableGraph[Future[Int]] =
Source(1 to 100)
.toMat(Sink.fold(0)(_ + _))(Keep.right)Komponenty wielokrotnego użytku
Ponieważ Sources, Flows i Sinks są niezmiennymi wartościami, można zdefiniować je raz i ponownie wykorzystywać w wielu potokach. Sprzyja to tworzeniu biblioteki małych, nazwanych etapów przetwarzania.
Flow zdefiniowany do parsowania można wstawić zarówno do potoku plikowego, jak i do potoku HTTP.
val parse: Flow[String, Int, akka.NotUsed] =
Flow[String].map(_.trim.toInt)
val fromFile = lines.via(parse)
val fromHttp = requestBody.via(parse)Typowe konstruktory Source
Akka Streams udostępnia wiele fabryk Source: Source.single, Source.repeat, Source.tick do emisji w określonych odstępach czasu, Source.future na podstawie obiektu Future oraz Source.empty.
Wybór właściwego konstruktora jasno określa przeznaczenie producenta danych.
import scala.concurrent.duration._
val ticks = Source.tick(0.seconds, 1.second, "tick")
val onceF = Source.future(scala.concurrent.Future.successful(42))
val forever = Source.repeat("x")Typowe konstruktory Sink
Podobnie Sinki obejmują Sink.head (pierwszy element jako Future), Sink.seq (zebranie wszystkich elementów do Seq), Sink.ignore (opróżnianie i odrzucanie elementów) oraz Sink.last.
W potokach, które przekazują wynik z powrotem do kodu, najczęściej używa się Sink.seq i Sink.fold.
import scala.concurrent.Future
val collect: Sink[Int, Future[Seq[Int]]] = Sink.seq
val firstOne: Sink[Int, Future[Int]] = Sink.head
val drain: Sink[Int, Future[akka.Done]] = Sink.ignoreSzybkie sprawdzenie
Sprawdź swoją wiedzę na temat liniowych typów etapów.
Podsumowanie
Poznali Państwo trzy liniowe elementy składowe: Source produkuje dane, Flow je przekształca, a Sink je konsumuje. Są to niezmienne, wielokrotnego użytku blueprinty, łączone za pomocą via i to.
Połączenie Source z Sink daje obiekt RunnableGraph, który zawiera wartość materializowaną, ale nie przesyła żadnych danych, dopóki nie zostanie uruchomiony. Następnie poznają Państwo transformowanie strumieni za pomocą bardziej rozbudowanych operatorów.
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 „Source, Flow i Sink” jest bezpłatna?
Tak — pełny tekst „Source, Flow i Sink” 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 „Source, Flow i Sink”?
Podstawowe elementy budowy strumieni. Ć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 1 z 4.
Ile czasu zajmuje lekcja „Source, Flow i Sink”?
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