Scala til backendudvikling og funktionel programmering · Lektion

Transformér streams

Map og filtrér data i bevægelse.

Lektion 2 af 413 trin

Transformér streams er en gratis Scala til backendudvikling og funktionel programmering-lektion på CoddyKit. Dette er lektion 2 af 4. Du kan læse hele lektionen gratis nedenfor — og derefter øve dig praktisk i browseren med en indbygget kodeeditor og en AI-vejleder, der er tilgængelig døgnet rundt. Den er en del af læringsforløbet i Scala til backendudvikling og funktionel programmering, og dine fremskridt synkroniseres på tværs af nettet og CoddyKit-appen. Scala til backendudvikling og funktionel programmering-kurset indeholder 4 lektioner i alt.

Operatorer som transformationer

Akka Streams indeholder et omfattende sæt operatorer til Source og Flow, som minder om Scalas samlings-API, men kører asynkront og respekterer backpressure.

Hver operator returnerer en ny blueprint, så transformationer kan sammensættes deklarativt, før strømmen overhovedet kører.

map og filter

map anvender en synkron funktion på hvert element, mens filter fjerner elementer, der ikke opfylder et prædikat. De er grundværktøjerne til transformation af enkelte elementer.

Begge bevarer rækkefølgen og viderefører afslutning og fejl nedstrøms.

val flow =
  Flow[Int]
    .filter(_ % 2 == 0)
    .map(n => n * n)

mapConcat til én-til-mange

Når ét input skal give flere outputs, skal du bruge mapConcat. Den tager en funktion, der returnerer noget itererbart, og flader resultaterne ud i strømmen.

Hvis du returnerer en tom samling, fjernes elementet effektivt.

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 og sliding

grouped(n) samler fortløbende elementer i en Seq med op til n elementer, hvilket er nyttigt ved massevise databaseskrivninger. sliding(n) udsender overlappende vinduer.

Gruppering reducerer omkostningen pr. element i I/O-tunge behandlingskæder.

val batches: Source[Seq[Int], akka.NotUsed] =
  Source(1 to 1000).grouped(100)

val windows =
  Source(1 to 10).sliding(3, step = 1)

scan og fold

scan udsender den løbende akkumulator efter hvert element og giver dermed en strøm af tilstande under udvikling. fold udsender kun den endelige akkumulerede værdi, når strømmen opstrøms er afsluttet.

Brug scan til tællere i realtid og fold til endelige aggregater.

val running =
  Source(1 to 5).scan(0)(_ + _) // 0,1,3,6,10,15

val total =
  Source(1 to 5).fold(0)(_ + _) // 15

mapAsync til asynkront arbejde

mapAsync(parallelism) kalder en funktion, der returnerer en Future, og udsender resultaterne i rækkefølge, mens op til parallelism futures kører samtidigt.

Brug den til asynkrone kald som databaseopslag eller HTTP-anmodninger, hvor rækkefølgen er vigtig.

import scala.concurrent.Future

val enriched =
  Flow[UserId]
    .mapAsync(parallelism = 4)(id => lookup(id))

def lookup(id: UserId): Future[User] = ???

mapAsyncUnordered

mapAsyncUnordered fungerer som mapAsync, men udsender hvert resultat, så snart det er færdigt, uden hensyn til inputrækkefølgen.

Det kan øge gennemløbet, når det efterfølgende trin er ligeglad med rækkefølgen, fordi en langsom future ikke længere blokerer hurtigere futures.

val fast =
  Flow[UserId]
    .mapAsyncUnordered(parallelism = 8)(id => lookup(id))

Tilstandsbaseret transformation med statefulMapConcat

Til transformationer af enkelte elementer, der kræver foranderlig lokal tilstand, opretter statefulMapConcat en ny tilstand for hver materialisering og returnerer en itererbar samling af outputs.

Det er den sikre måde at bevare tællere eller buffere på uden at dele tilstand mellem kørsler af strømmen.

val withIndex: Flow[String, (Int, String), akka.NotUsed] =
  Flow[String].statefulMapConcat { () =>
    var i = 0
    elem => { i += 1; List((i, elem)) }
  }

Tidsbaserede operatorer

Strømme kan transformeres ud fra tid såvel som indhold. throttle begrænser udsendelsesraten, groupedWithin grupperer efter størrelse eller forløbet tid, og takeWithin begrænser varigheden.

De er afgørende ved begrænsning af raten til eksterne API'er.

import scala.concurrent.duration._

val limited =
  Source(1 to 1000)
    .throttle(10, 1.second)
    .groupedWithin(100, 500.millis)

Håndtering af fejl i transformationer

En undtagelse, der kastes inde i en operator, får som standard hele strømmen til at fejle. En supervisionsstrategi kan i stedet resume (fjerne det ugyldige element) eller restart behandlingstrinnet.

Tilknyt strategien med withAttributes på Flowet.

import akka.stream.{ActorAttributes, Supervision}

val safe =
  Flow[String].map(_.toInt)
    .withAttributes(
      ActorAttributes.supervisionStrategy(_ => Supervision.Resume))

Sammensætning af Flows

Små Flows kan sammensættes til større med via, så der opstår ét genanvendeligt Flow. Det holder hver transformation fokuseret og gør den mulig at teste uafhængigt.

Det sammensatte Flow har inputtypen fra det første og outputtypen fra det sidste trin.

val parse  = Flow[String].map(_.toInt)
val square = Flow[Int].map(n => n * n)

val parseAndSquare: Flow[String, Int, akka.NotUsed] =
  parse.via(square)

Hurtigt tjek

Overvej asynkrone transformationer og deres garantier for rækkefølgen.

Opsummering

Du har udforsket transformationsoperatorer: map/filter til enkelte elementer, mapConcat til én-til-mange, grouped til gruppering, akkumulering med scan og fold samt asynkront arbejde via mapAsync.

Du har også set tilstandsbaserede transformationer, tidsbaserede operatorer som throttle, supervisionsstrategier til fejl og hvordan Flows sammensættes med via. Næste emne er, hvordan backpressure holder disse trin sikre.

Gratis at komme i gang

Lær Scala med en AI-underviser — gratis

Skriv og kør rigtig kode i din browser, få øjeblikkelig hjælp fra en AI-underviser døgnet rundt, og fortsæt, hvor du slap, på web eller i appen.

Kurser
39
Lektioner
143

Ofte stillede spørgsmål

Er lektionen “Transformér streams” gratis?

Ja — hele teksten til “Transformér streams” kan læses gratis her på nettet. Hvis du vil øve dig interaktivt med en indbygget kodeeditor og en AI-vejleder døgnet rundt og få adgang til resten af Scala til backendudvikling og funktionel programmering-kurset, skal du opgradere til CoddyKit PRO. Scala til backendudvikling og funktionel programmering-kurset indeholder 4 lektioner i alt.

Hvad lærer jeg i “Transformér streams”?

Map og filtrér data i bevægelse. Du øver dig i Scala til backendudvikling og funktionel programmering med praktisk kode, som du kører direkte i browseren, og en AI-vejleder døgnet rundt besvarer dine spørgsmål, mens du arbejder dig gennem lektionen.

Skal jeg have erfaring for at begynde på Scala til backendudvikling og funktionel programmering?

Der kræves ingen tidligere erfaring. Scala til backendudvikling og funktionel programmering på CoddyKit er tilrettelagt for både begyndere og øvede, så du kan starte her eller fra begyndelsen og lære i dit eget tempo. Dette er lektion 2 af 4.

Hvor lang tid tager lektionen “Transformér streams”?

De fleste CoddyKit-lektioner tager cirka 5–10 minutter. Hver lektion er kort og interaktiv, så du gør løbende fremskridt og kan fortsætte, hvor du slap – på både web og app.

Kan jeg skrive og køre kode i denne Scala til backendudvikling og funktionel programmering-lektion?

Ja. Alle Scala til backendudvikling og funktionel programmering-lektioner har en indbygget kodeeditor, så du kan skrive og køre rigtig kode direkte i din browser og få øjeblikkelig feedback fra AI – uden lokal opsætning.

Alle lektioner i dette kursus

  1. Source, Flow og Sink
  2. Transformér streams
  3. Backpressure
  4. Kør en pipeline
← Tilbage til Scala til backendudvikling og funktionel programmering