0Pricing
Scala for Backend Engineering & Functional Programming · Lekcja

Backpressure

Bezpiecznie obsłuży Pan/Pani szybkich producentów danych.

Backpressure to bezpłatna lekcja Scala for Backend Engineering & Functional Programming na CoddyKit. To lekcja 3 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.

Czym jest backpressure?

Backpressure to mechanizm sterowania przepływem, który zapobiega przeciążeniu wolnego konsumenta przez szybkiego producenta. Zamiast bez ograniczeń buforować dane lub je odrzucać, konsument sygnalizuje, ile danych może obsłużyć.

Akka Streams implementuje standard Reactive Streams, w którym zapotrzebowanie przepływa w górę strumienia, a elementy w dół.

Przepływ sterowany zapotrzebowaniem

Każdy etap emituje dane dopiero wtedy, gdy kolejny etap zasygnalizuje zapotrzebowanie. Sink żąda N elementów, a zapotrzebowanie propaguje się w górę strumienia, aż Source wyprodukuje dokładnie tyle elementów, ile zażądano.

Ten protokół oparty na pobieraniu oznacza, że producenci nigdy nie wysyłają większej liczby elementów, niż konsumenci mogą przetworzyć.

Dlaczego ma to znaczenie dla potoków

Bez backpressure szybki konsument Kafki zasilający wolną bazę danych zgromadziłby miliony rekordów będących w trakcie przetwarzania, wyczerpując pamięć i doprowadzając do awarii procesu.

Backpressure w naturalny sposób spowalnia wcześniejsze etapy do tempa najwolniejszego etapu, zapewniając stabilne zużycie pamięci pod obciążeniem.

// Fast source, slow sink: backpressure slows the source
val g =
  Source(1 to 1000000)
    .map(_ * 2)
    .to(slowDatabaseSink)

Buforowanie wewnętrzne

Pomiędzy granicami asynchronicznymi Akka Streams utrzymuje mały bufor wewnętrzny (domyślnie 16 elementów). Pochłania on krótkie skoki napływu danych, dzięki czemu etapy nie muszą synchronizować się przy każdym elemencie.

Gdy bufor się zapełni, uruchamia się backpressure, a wcześniejszy etap przestaje produkować dane do czasu zwolnienia miejsca.

import akka.stream.Attributes

val buffered =
  Flow[Int]
    .map(identity)
    .addAttributes(Attributes.inputBuffer(initial = 32, max = 32))

Jawny bufor ze strategią przepełnienia

Operator buffer wstawia jawny bufor o wybranym rozmiarze wraz z obiektem OverflowStrategy, który określa, co ma się stać po jego zapełnieniu.

Pozwala to wymienić dodatkową pamięć na możliwość uniezależnienia szybkości producenta od szybkości konsumenta.

import akka.stream.OverflowStrategy

val withBuffer =
  Source(1 to 1000)
    .buffer(size = 100, OverflowStrategy.backpressure)

Strategie przepełnienia

Strategie obejmują backpressure (spowolnienie wcześniejszego etapu), dropHead/dropTail (odrzucenie najstarszych lub najnowszych elementów), dropBuffer, dropNew oraz fail (zakończenie z błędem).

Strategie odrzucania sprawdzają się w przypadku danych na żywo, takich jak odczyty z czujników, gdy nieaktualne wartości można bezpiecznie odrzucić.

import akka.stream.OverflowStrategy

val latestWins =
  liveTicks.buffer(1, OverflowStrategy.dropHead)

val strict =
  liveTicks.buffer(50, OverflowStrategy.fail)

Conflate do podsumowywania

Gdy konsument działa wolno, conflate łączy oczekujące elementy w jeden za pomocą funkcji łączącej, zamiast przechowywać je wszystkie w buforze.

Można na przykład połączyć wiele aktualizacji liczbowych w ich sumę, aby konsument zawsze widział agregat tego, co zostało pominięte.

val summarized =
  fastMetrics
    .conflate((acc, next) => acc + next)

// Slow downstream receives summed batches

Expand do zaspokajania zapotrzebowania

expand jest przeciwieństwem conflate: gdy kolejny etap zgłasza zapotrzebowanie szybciej, niż wcześniejszy etap produkuje dane, operator tworzy dodatkowe elementy na podstawie ostatniej zaobserwowanej wartości.

Przydaje się to do emitowania ostatniego odczytu w stałym tempie.

val repeated =
  sensor.expand(last => Iterator.continually(last))

// Downstream always gets the latest sensor value

Granice asynchroniczne

Domyślnie połączone etapy działają w ramach jednego aktora, bez buforowania między nimi. Wstawienie async umieszcza etap we własnym aktorze, dodając bufor i umożliwiając potokowe przetwarzanie równoległe.

Granice asynchroniczne to miejsca, w których faktycznie znajdują się bufory backpressure.

val pipelined =
  Source(1 to 1000)
    .map(slowStep).async
    .map(anotherSlowStep).async
    .to(Sink.ignore)

Throttle jako jawne sterowanie szybkością

throttle nakłada określony maksymalny limit szybkości, generując backpressure w górę strumienia, aby go respektować. Chroni to zewnętrzne usługi z ograniczeniem częstotliwości wywołań, nawet gdy konsument mógłby działać szybciej.

Parametr burst pozwala na krótkie skoki powyżej stałego limitu.

import scala.concurrent.duration._

val limited =
  requests
    .throttle(
      elements = 100, per = 1.second, maximumBurst = 20,
      akka.stream.ThrottleMode.Shaping)

Obserwowanie backpressure

Backpressure można wykryć, obserwując spowolnienia wcześniejszych etapów lub mierząc zajętość buforów. Operator log oraz atrybuty strumienia Akka pomagają śledzić miejsce zatrzymania potoku.

Trwale zapełniony bufor wskazuje najwolniejszy etap ograniczający przepustowość.

val traced =
  Source(1 to 100)
    .log("after-source")
    .map(_ * 2)
    .log("after-map")
    .to(Sink.ignore)

Szybkie sprawdzenie

Zastanów się, jak Akka Streams chroni wolnego konsumenta przed zalaniem danymi przez szybkiego producenta.

Podsumowanie

Backpressure to sterowany zapotrzebowaniem fundament Akka Streams: konsumenci sygnalizują zapotrzebowanie w górę strumienia, dzięki czemu producenci nie mogą ich przeciążyć, a zużycie pamięci pozostaje ograniczone.

Poznali Państwo bufory wewnętrzne, jawny operator buffer ze strategiami przepełnienia, podsumowywanie za pomocą conflate, zaspokajanie zapotrzebowania za pomocą expand, granice asynchroniczne oraz jawne sterowanie szybkością za pomocą throttle. Następnie uruchomią Państwo kompletny potok.

Często zadawane pytania

Czy lekcja „Backpressure” jest bezpłatna?

Tak — pełny tekst „Backpressure” 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 „Backpressure”?

Bezpiecznie obsłuży Pan/Pani szybkich producentów danych. Ć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 3 z 4.

Ile czasu zajmuje lekcja „Backpressure”?

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