Source, Flow og Sink
Byggestenene i streaming.
Source, Flow og Sink er en gratis Scala til backendudvikling og funktionel programmering-lektion på CoddyKit. Dette er lektion 1 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.
De tre byggesten
Akka Streams modellerer en datapipeline som en graf af behandlingstrin. De tre centrale lineære trin er Source (producerer elementer), Flow (transformerer dem) og Sink (forbruger dem).
En Source har ét output, en Sink har ét input, og en Flow har præcis ét input og ét output. Når du forbinder dem, beskriver du hvad der skal ske, ikke hvornår.
Definition af en Source
En Source[Out, Mat] udsender elementer af typen Out og stiller en materialiseret værdi af typen Mat til rådighed. De enkleste kilder kommer fra samlinger i hukommelsen eller områder.
Indtil datastrømmen køres, er en Source blot en uforanderlig skabelon, som frit kan genbruges.
import akka.stream.scaladsl.Source
val numbers: Source[Int, akka.NotUsed] =
Source(1 to 100)
val single: Source[String, akka.NotUsed] =
Source.single("hello")Definition af en Sink
En Sink[In, Mat] forbruger elementer af typen In. Den materialiserede værdi indeholder ofte resultatet af forbruget, f.eks. en Future, der fuldføres, når datastrømmen afsluttes.
Sink.foreach udfører en bivirkning pr. element, mens Sink.fold akkumulerer ét resultat.
import akka.stream.scaladsl.Sink
import scala.concurrent.Future
val printSink: Sink[Int, Future[akka.Done]] =
Sink.foreach(println)
val sumSink: Sink[Int, Future[Int]] =
Sink.fold(0)(_ + _)Definition af en Flow
En Flow[In, Out, Mat] ligger mellem en Source og en Sink og transformerer hvert element. Flows kan genbruges selvstændigt og sammensættes, før de tilknyttes et endepunkt.
Her fordobler en Flow heltal og konverterer dem til strenge.
import akka.stream.scaladsl.Flow
val doubleToString: Flow[Int, String, akka.NotUsed] =
Flow[Int]
.map(_ * 2)
.map(n => s"value=$n")Forbindelse fra Source til Sink
Operatoren via knytter en Flow til en Source, og to knytter en Sink til den. Når du forbinder en Source direkte med en Sink med to, oprettes en RunnableGraph: en lukket, kørbar skabelon.
Der flyttes endnu ingen elementer; dette beskriver kun topologien.
import akka.stream.scaladsl.{Source, Sink, RunnableGraph}
val graph: RunnableGraph[akka.NotUsed] =
Source(1 to 10).to(Sink.foreach(println))via: Indsættelse af en Flow
Brug via til at indsætte et Flow i behandlingskæden. En Source.via(flow) giver en ny Source, hvis outputtype svarer til Flowets output.
Ved at kæde kald til via sammen kan du opbygge lange transformationskæder af små, testbare Flow-komponenter.
val pipeline =
Source(1 to 10)
.via(Flow[Int].filter(_ % 2 == 0))
.via(Flow[Int].map(_ * 10))
.to(Sink.foreach(println))Typesikkerhed på tværs af trin
Compileren sikrer, at outputtypen fra hvert trin svarer til inputtypen for det næste. En Source[Int] kan ikke forbindes med en Sink[String] uden et mellemliggende Flow, der konverterer typen.
Denne statiske kontrol opdager fejl i forbindelserne mellem trinene, før programmet kører.
// Source[Int] -> Flow[Int, String] -> Sink[String]
val ok =
Source(1 to 3)
.via(Flow[Int].map(_.toString))
.to(Sink.foreach[String](println))Materialiserede værdier
Hver blueprint indeholder en materialiseret værdi: et håndtag, der oprettes, når strømmen kører. En Sink.fold materialiserer en Future med resultatet. Som standard bevarer en kombination af trin den materialiserede værdi længst til venstre (NotUsed for almindelige kilder).
Brug toMat og Keep til at vælge, hvilken sides værdi du vil have.
import akka.stream.scaladsl.Keep
import scala.concurrent.Future
val g: RunnableGraph[Future[Int]] =
Source(1 to 100)
.toMat(Sink.fold(0)(_ + _))(Keep.right)Genanvendelige komponenter
Fordi Sources, Flows og Sinks er uforanderlige værdier, kan du definere dem én gang og genbruge dem i mange behandlingskæder. Det fremmer et bibliotek af små navngivne behandlingstrin.
Et Flow, der er defineret til parsing, kan indsættes i både en filbaseret behandlingskæde og en HTTP-behandlingskæde.
val parse: Flow[String, Int, akka.NotUsed] =
Flow[String].map(_.trim.toInt)
val fromFile = lines.via(parse)
val fromHttp = requestBody.via(parse)Almindelige Source-konstruktører
Akka Streams leveres med mange Source-fabrikker: Source.single, Source.repeat, Source.tick til tidsstyret udsendelse, Source.future fra en Future og Source.empty.
Ved at vælge den rigtige konstruktør bliver hensigten med en dataproducent tydelig.
import scala.concurrent.duration._
val ticks = Source.tick(0.seconds, 1.second, "tick")
val onceF = Source.future(scala.concurrent.Future.successful(42))
val forever = Source.repeat("x")Almindelige Sink-konstruktører
Tilsvarende omfatter sinks Sink.head (første element som en Future), Sink.seq (indsamler alt i en Seq), Sink.ignore (tømmer og kasserer) og Sink.last.
I behandlingskæder, der sender et resultat tilbage til din kode, er Sink.seq og Sink.fold de mest almindelige.
import scala.concurrent.Future
val collect: Sink[Int, Future[Seq[Int]]] = Sink.seq
val firstOne: Sink[Int, Future[Int]] = Sink.head
val drain: Sink[Int, Future[akka.Done]] = Sink.ignoreHurtigt tjek
Test din forståelse af de lineære typer for behandlingstrin.
Opsummering
Du har lært de tre lineære byggesten: Source producerer, Flow transformerer, og Sink forbruger. De er uforanderlige, genanvendelige blueprints, der forbindes med via og to.
Hvis du forbinder en Source med en Sink, får du en RunnableGraph, der indeholder en materialiseret værdi, men ikke flytter data, før den køres. Dernæst skal du transformere strømme med mere avancerede operatorer.
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 “Source, Flow og Sink” gratis?
Ja — hele teksten til “Source, Flow og Sink” 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 “Source, Flow og Sink”?
Byggestenene i streaming. 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 1 af 4.
Hvor lang tid tager lektionen “Source, Flow og Sink”?
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