Scala for Backend Engineering & Functional Programming · Lekcja

Przekształcanie strumieni

Zmapuje i odfiltruje Pan/Pani przepływające dane.

Lekcja 2 z 413 kroki

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)(_ + _) // 15

mapAsync 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.

Bezpłatny start

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

  1. Source, Flow i Sink
  2. Przekształcanie strumieni
  3. Backpressure
  4. Uruchamianie potoku
← Powrót do Scala for Backend Engineering & Functional Programming