Przekształcanie strumieni
Zmapuje i odfiltruje Pan/Pani przepływające dane.
Przekształcanie strumieni to bezpłatna lekcja Scala for Backend Engineering & Functional Programming na CoddyKit. To lekcja 2 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.
Operatory jako transformacje
Akka Streams udostępnia bogaty zestaw operatorów dla Source i Flow, które przypominają API kolekcji Scali, ale działają asynchronicznie i respektują backpressure.
Każdy operator zwraca nowy blueprint, dzięki czemu transformacje można deklaratywnie komponować, zanim strumień zostanie uruchomiony.
map i filter
map stosuje synchroniczną funkcję do każdego elementu, a filter odrzuca elementy niespełniające predykatu. Są to podstawowe operatory transformacji element po elemencie.
Oba zachowują kolejność oraz przekazują sygnały zakończenia i błędy w dół strumienia.
val flow =
Flow[Int]
.filter(_ % 2 == 0)
.map(n => n * n)mapConcat dla relacji jeden-do-wielu
Jeśli jedno wejście powinno wygenerować kilka wyjść, należy użyć mapConcat. Przyjmuje on funkcję zwracającą obiekt iterowalny i spłaszcza wyniki do postaci strumienia.
Zwrócenie pustej kolekcji skutecznie usuwa dany element.
val explode: Flow[String, String, akka.NotUsed] =
Flow[String].mapConcat(line => line.split(",").toList)
val words = Source(List("a,b", "c,d,e"))
.via(explode)grouped i sliding
grouped(n) grupuje kolejne elementy w obiekt Seq zawierający do n elementów, co jest przydatne przy zbiorczych zapisach do bazy danych. sliding(n) emituje nakładające się na siebie okna.
Grupowanie zmniejsza narzut przypadający na element w potokach intensywnie korzystających z operacji wejścia-wyjścia.
val batches: Source[Seq[Int], akka.NotUsed] =
Source(1 to 1000).grouped(100)
val windows =
Source(1 to 10).sliding(3, step = 1)scan i fold
scan emituje bieżący akumulator po każdym elemencie, tworząc strumień zmieniającego się stanu. fold emituje tylko końcową wartość akumulatora po zakończeniu strumienia nadrzędnego.
Użyj scan dla liczników działających na bieżąco, a fold dla agregacji końcowych.
val running =
Source(1 to 5).scan(0)(_ + _) // 0,1,3,6,10,15
val total =
Source(1 to 5).fold(0)(_ + _) // 15mapAsync dla pracy asynchronicznej
mapAsync(parallelism) wywołuje funkcję zwracającą obiekt Future i emituje wyniki w kolejności, wykonując jednocześnie maksymalnie parallelism obiektów Future.
Użyj go do wywołań asynchronicznych, takich jak wyszukiwanie w bazie danych lub żądania HTTP, gdy kolejność ma znaczenie.
import scala.concurrent.Future
val enriched =
Flow[UserId]
.mapAsync(parallelism = 4)(id => lookup(id))
def lookup(id: UserId): Future[User] = ???mapAsyncUnordered
mapAsyncUnordered działa podobnie jak mapAsync, ale emituje każdy wynik natychmiast po jego ukończeniu, ignorując kolejność danych wejściowych.
Może zwiększyć przepustowość, gdy kolejność nie ma znaczenia dla kolejnego etapu, ponieważ wolny obiekt Future nie blokuje już szybszych.
val fast =
Flow[UserId]
.mapAsyncUnordered(parallelism = 8)(id => lookup(id))Transformacja stanowa z statefulMapConcat
W przypadku transformacji wykonywanych dla poszczególnych elementów, które wymagają zmiennego stanu lokalnego, statefulMapConcat tworzy nowy stan przy każdej materializacji i zwraca iterowalny zbiór wyników.
Jest to bezpieczny sposób przechowywania liczników lub buforów bez współdzielenia stanu między uruchomieniami strumienia.
val withIndex: Flow[String, (Int, String), akka.NotUsed] =
Flow[String].statefulMapConcat { () =>
var i = 0
elem => { i += 1; List((i, elem)) }
}Operatory zależne od czasu
Strumienie mogą być transformowane na podstawie czasu, a nie tylko zawartości. throttle ogranicza częstotliwość emisji, groupedWithin grupuje elementy według rozmiaru lub upływu czasu, a takeWithin ogranicza czas trwania.
Operatory te są niezbędne do ograniczania częstotliwości wywołań zewnętrznych interfejsów API.
import scala.concurrent.duration._
val limited =
Source(1 to 1000)
.throttle(10, 1.second)
.groupedWithin(100, 500.millis)Obsługa błędów w transformacjach
Wyjątek rzucony wewnątrz operatora domyślnie kończy cały strumień błędem. Strategia nadzoru może zamiast tego wykonać resume (odrzucić niepoprawny element) albo restart etapu.
Dołącz strategię za pomocą withAttributes na obiekcie Flow.
import akka.stream.{ActorAttributes, Supervision}
val safe =
Flow[String].map(_.toInt)
.withAttributes(
ActorAttributes.supervisionStrategy(_ => Supervision.Resume))Komponowanie obiektów Flow
Małe obiekty Flow można łączyć w większe za pomocą via, uzyskując pojedynczy obiekt Flow wielokrotnego użytku. Dzięki temu każda transformacja pozostaje skoncentrowana na jednym zadaniu i może być testowana niezależnie.
Połączony Flow ma typ wejściowy pierwszego elementu i typ wyjściowy ostatniego.
val parse = Flow[String].map(_.toInt)
val square = Flow[Int].map(n => n * n)
val parseAndSquare: Flow[String, Int, akka.NotUsed] =
parse.via(square)Szybkie sprawdzenie
Przeanalizuj transformacje asynchroniczne i gwarancje dotyczące ich kolejności.
Podsumowanie
Poznali Państwo operatory transformacji: działające na poszczególnych elementach map/filter, relację jeden-do-wielu za pomocą mapConcat, grupowanie za pomocą grouped, akumulację za pomocą scan i fold oraz pracę asynchroniczną za pomocą mapAsync.
Poznali Państwo również transformacje stanowe, operatory zależne od czasu, takie jak throttle, strategie nadzoru błędów oraz sposób komponowania obiektów Flow za pomocą via. Następnie dowiedzą się Państwo, jak backpressure zapewnia bezpieczeństwo tych etapó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 „Przekształcanie strumieni” jest bezpłatna?
Tak — pełny tekst „Przekształcanie strumieni” 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 „Przekształcanie strumieni”?
Zmapuje i odfiltruje Pan/Pani przepływające dane. Ć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 2 z 4.
Ile czasu zajmuje lekcja „Przekształcanie strumieni”?
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