Source, flux et puits
Les briques de base du traitement en flux.
Source, flux et puits est une leçon Scala for Backend Engineering & Functional Programming gratuite sur CoddyKit. Ceci est la leçon 1 sur 4. Tu peux lire la leçon complète ci-dessous gratuitement — puis la pratiquer en direct dans le navigateur avec un éditeur de code intégré et un tuteur IA 24/7. Elle fait partie du parcours d'apprentissage Scala for Backend Engineering & Functional Programming, et ta progression se synchronise sur le web et l'application CoddyKit. Le cours Scala for Backend Engineering & Functional Programming comprend 4 leçons au total.
Les trois éléments fondamentaux
Akka Streams modélise un pipeline de données comme un graphe d'étapes de traitement. Les trois étapes linéaires fondamentales sont Source (produit des éléments), Flow (les transforme) et Sink (les consomme).
Une Source possède une sortie, un Sink une entrée, et un Flow possède exactement une entrée et une sortie. Les relier décrit ce qui doit se produire, et non quand.
Définir une Source
Une Source[Out, Mat] émet des éléments de type Out et expose une valeur matérialisée de type Mat. Les Sources les plus simples proviennent de collections en mémoire ou d'intervalles.
Tant que le flux n'est pas exécuté, une Source n'est qu'un modèle immuable qui peut être librement réutilisé.
import akka.stream.scaladsl.Source
val numbers: Source[Int, akka.NotUsed] =
Source(1 to 100)
val single: Source[String, akka.NotUsed] =
Source.single("hello")Définir un Sink
Un Sink[In, Mat] consomme des éléments de type In. La valeur matérialisée contient souvent le résultat de la consommation, par exemple un Future qui se termine lorsque le flux s'achève.
Sink.foreach exécute un effet secondaire pour chaque élément ; Sink.fold accumule un résultat unique.
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)(_ + _)Définir un Flow
Un Flow[In, Out, Mat] se place entre une Source et un Sink et transforme chaque élément. Les Flows sont réutilisables seuls et peuvent être composés avant d'être associés à un point d'extrémité.
Ici, un Flow double des entiers et les convertit en chaînes.
import akka.stream.scaladsl.Flow
val doubleToString: Flow[Int, String, akka.NotUsed] =
Flow[Int]
.map(_ * 2)
.map(n => s"value=$n")Relier une Source à un Sink
L'opérateur via associe un Flow à une Source, et to associe un Sink. Relier directement une Source à un Sink avec to produit un RunnableGraph : un modèle fermé et exécutable.
Aucun élément ne circule encore ; cela ne décrit que la topologie.
import akka.stream.scaladsl.{Source, Sink, RunnableGraph}
val graph: RunnableGraph[akka.NotUsed] =
Source(1 to 10).to(Sink.foreach(println))via : insérer un Flow
Utilisez via pour insérer un Flow dans le pipeline. Un Source.via(flow) produit une nouvelle Source dont le type de sortie correspond à celui de la sortie du Flow.
En enchaînant les appels à via, vous pouvez construire de longs pipelines de transformation à partir de petits Flows testables.
val pipeline =
Source(1 to 10)
.via(Flow[Int].filter(_ % 2 == 0))
.via(Flow[Int].map(_ * 10))
.to(Sink.foreach(println))Sûreté des types entre les étapes
Le compilateur vérifie que le type de sortie de chaque étape correspond au type d'entrée de la suivante. Une Source[Int] ne peut pas être reliée à un Sink[String] sans un Flow intermédiaire qui convertit le type.
Cette vérification statique détecte les erreurs de câblage du pipeline avant l'exécution.
// Source[Int] -> Flow[Int, String] -> Sink[String]
val ok =
Source(1 to 3)
.via(Flow[Int].map(_.toString))
.to(Sink.foreach[String](println))Valeurs matérialisées
Chaque modèle porte une valeur matérialisée : une référence produite lorsque le flux s'exécute. Un Sink.fold matérialise un Future contenant le résultat. Par défaut, la combinaison d'étapes conserve la valeur matérialisée la plus à gauche (NotUsed pour les Sources simples).
Utilisez toMat et Keep pour sélectionner la valeur du côté souhaité.
import akka.stream.scaladsl.Keep
import scala.concurrent.Future
val g: RunnableGraph[Future[Int]] =
Source(1 to 100)
.toMat(Sink.fold(0)(_ + _))(Keep.right)Composants réutilisables
Les Sources, les Flows et les Sinks étant des valeurs immuables, vous pouvez les définir une fois et les réutiliser dans de nombreux pipelines. Cela favorise la création d'une bibliothèque de petites étapes de traitement nommées.
Un Flow défini pour l'analyse syntaxique peut être inséré à la fois dans un pipeline de fichier et dans un pipeline HTTP.
val parse: Flow[String, Int, akka.NotUsed] =
Flow[String].map(_.trim.toInt)
val fromFile = lines.via(parse)
val fromHttp = requestBody.via(parse)Constructeurs courants de Source
Akka Streams fournit de nombreuses fabriques de Source : Source.single, Source.repeat, Source.tick pour les émissions temporisées, Source.future à partir d'un Future et Source.empty.
Choisir le bon constructeur rend explicite l'intention du producteur de données.
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")Constructeurs courants de Sink
De même, les sinks comprennent Sink.head (premier élément sous forme de Future), Sink.seq (tout regrouper dans une Seq), Sink.ignore (vider et ignorer) et Sink.last.
Pour les pipelines qui renvoient un résultat à votre code, Sink.seq et Sink.fold sont les plus courants.
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.ignoreVérification rapide
Vérifiez votre compréhension des types d'étapes linéaires.
Récapitulatif
Vous avez découvert les trois éléments linéaires fondamentaux : Source produit, Flow transforme et Sink consomme. Ce sont des modèles immuables et réutilisables, reliés avec via et to.
Relier une Source à un Sink produit un RunnableGraph qui porte une valeur matérialisée, mais ne déplace aucune donnée avant son exécution. Vous allez ensuite transformer les flux avec des opérateurs plus riches.
Questions Fréquemment Posées
La leçon « Source, flux et puits » est-elle gratuite ?
Oui — le texte complet de « Source, flux et puits » est gratuit à lire ici sur le web. Pour la pratiquer de manière interactive (un éditeur de code intégré et un tuteur IA 24/7) et déverrouiller le reste du cours Scala for Backend Engineering & Functional Programming, passe à CoddyKit PRO. Le cours Scala for Backend Engineering & Functional Programming comprend 4 leçons au total.
Qu'est-ce que j'apprendrai dans « Source, flux et puits » ?
Les briques de base du traitement en flux. Tu pratiques Scala for Backend Engineering & Functional Programming avec du code pratique que tu exécutes directement dans le navigateur, et un tuteur IA 24/7 répond à tes questions au fur et à mesure que tu avances dans la leçon.
Dois-je avoir de l'expérience pour commencer Scala for Backend Engineering & Functional Programming ?
Aucune expérience préalable n'est requise. Scala for Backend Engineering & Functional Programming sur CoddyKit est structuré pour les débutants jusqu'aux apprenants avancés, donc tu peux commencer ici ou depuis le début et avancer à ton rythme. Ceci est la leçon 1 sur 4.
Combien de temps prend la leçon « Source, flux et puits » ?
La plupart des leçons CoddyKit prennent environ 5–10 minutes. Chacune est courte et interactive, tu progresses régulièrement et tu repiques exactement où tu t'es arrêté sur le web et l'app.
Peux-tu écrire et exécuter du code dans cette leçon Scala for Backend Engineering & Functional Programming ?
Oui. Chaque leçon Scala for Backend Engineering & Functional Programming inclut un éditeur de code intégré, tu écris et exécutes du vrai code directement dans ton navigateur et tu reçois des retours IA instantanés — aucune configuration locale requise.
Toutes les leçons de ce cours
- Source, flux et puits
- Transformer des flux
- Contre-pression
- Exécuter un pipeline