Transformér streams
Map og filtrér data i bevægelse.
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)(_ + _) // 15mapAsync 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.
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
- Source, Flow og Sink
- Transformér streams
- Backpressure
- Kør en pipeline