Transformer des flux
Mappez et filtrez les données en circulation.
Transformer des flux est une leçon Scala for Backend Engineering & Functional Programming gratuite sur CoddyKit. Ceci est la leçon 2 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 opérateurs comme transformations
Akka Streams fournit un riche ensemble d'opérateurs sur les Sources et les Flows, qui rappellent l'API des collections de Scala, mais s'exécutent de manière asynchrone et respectent la contre-pression.
Chaque opérateur retourne un nouveau modèle ; les transformations sont donc composées de manière déclarative avant même l'exécution du flux.
map et filter
map applique une fonction synchrone à chaque élément ; filter supprime les éléments qui ne satisfont pas un prédicat. Ce sont les outils essentiels des transformations élément par élément.
Tous deux préservent l'ordre et propagent l'achèvement ainsi que les échecs vers l'aval.
val flow =
Flow[Int]
.filter(_ % 2 == 0)
.map(n => n * n)mapConcat pour produire plusieurs éléments
Lorsqu'une entrée doit produire plusieurs sorties, utilisez mapConcat. Il prend une fonction qui renvoie un itérable et aplatit les résultats dans le flux.
Retourner une collection vide revient à supprimer l'élément.
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 et sliding
grouped(n) regroupe les éléments consécutifs dans un Seq contenant au plus n éléments, ce qui est utile pour les écritures en masse dans une base de données. sliding(n) émet des fenêtres qui se chevauchent.
Le traitement par lots réduit le coût par élément dans les pipelines fortement dépendants des entrées-sorties.
val batches: Source[Seq[Int], akka.NotUsed] =
Source(1 to 1000).grouped(100)
val windows =
Source(1 to 10).sliding(3, step = 1)scan et fold
scan émet l'accumulateur courant après chaque élément, ce qui produit un flux d'état évolutif. fold n'émet la valeur accumulée finale qu'une fois la source amont terminée.
Utilisez scan pour les compteurs en temps réel et fold pour les agrégats finaux.
val running =
Source(1 to 5).scan(0)(_ + _) // 0,1,3,6,10,15
val total =
Source(1 to 5).fold(0)(_ + _) // 15mapAsync pour le travail asynchrone
mapAsync(parallelism) appelle une fonction qui renvoie un Future et émet les résultats dans l'ordre, en exécutant simultanément jusqu'à parallelism futures.
Utilisez-le pour les appels asynchrones, comme les recherches en base de données ou les requêtes HTTP, lorsque l'ordre est important.
import scala.concurrent.Future
val enriched =
Flow[UserId]
.mapAsync(parallelism = 4)(id => lookup(id))
def lookup(id: UserId): Future[User] = ???mapAsyncUnordered
mapAsyncUnordered se comporte comme mapAsync, mais émet chaque résultat dès qu'il est terminé, sans tenir compte de l'ordre des entrées.
Il peut améliorer le débit lorsque l'étape en aval ne dépend pas de l'ordre, car un traitement asynchrone lent ne bloque plus les traitements plus rapides.
val fast =
Flow[UserId]
.mapAsyncUnordered(parallelism = 8)(id => lookup(id))Transformation avec état avec statefulMapConcat
Pour les transformations élément par élément qui nécessitent un état local mutable, statefulMapConcat crée un nouvel état à chaque matérialisation et renvoie un itérable de sorties.
C'est la manière sûre de conserver des compteurs ou des tampons sans partager l'état entre les exécutions du flux.
val withIndex: Flow[String, (Int, String), akka.NotUsed] =
Flow[String].statefulMapConcat { () =>
var i = 0
elem => { i += 1; List((i, elem)) }
}Opérateurs fondés sur le temps
Les flux peuvent effectuer des transformations en fonction du temps aussi bien que du contenu. throttle plafonne le débit d'émission, groupedWithin regroupe les éléments selon une taille ou une durée écoulée, et takeWithin limite la durée.
Ces opérateurs sont essentiels pour limiter le débit des API externes.
import scala.concurrent.duration._
val limited =
Source(1 to 1000)
.throttle(10, 1.second)
.groupedWithin(100, 500.millis)Gérer les erreurs dans les transformations
Par défaut, une exception levée dans un opérateur fait échouer l'ensemble du flux. Une stratégie de supervision peut à la place resume (ignorer l'élément incorrect) ou restart l'étape.
Associez la stratégie au flux avec withAttributes.
import akka.stream.{ActorAttributes, Supervision}
val safe =
Flow[String].map(_.toInt)
.withAttributes(
ActorAttributes.supervisionStrategy(_ => Supervision.Resume))Composer des flux
De petits flux se composent en flux plus grands avec via, ce qui produit un flux unique et réutilisable. Chaque transformation reste ainsi ciblée et peut être testée indépendamment.
Le flux composé possède le type d'entrée du premier flux et le type de sortie du dernier.
val parse = Flow[String].map(_.toInt)
val square = Flow[Int].map(n => n * n)
val parseAndSquare: Flow[String, Int, akka.NotUsed] =
parse.via(square)Vérification rapide
Réfléchissez aux transformations asynchrones et à leurs garanties concernant l'ordre.
Récapitulatif
Vous avez étudié les opérateurs de transformation : map/filter élément par élément, mapConcat pour produire plusieurs éléments à partir d'un seul, grouped pour regrouper les éléments, scan et fold pour les accumuler, ainsi que mapAsync pour les traitements asynchrones.
Vous avez également découvert les transformations avec état, les opérateurs fondés sur le temps comme throttle, les stratégies de supervision des erreurs et la composition des flux avec via. Ensuite : comment la contre-pression sécurise ces étapes.
Questions Fréquemment Posées
La leçon « Transformer des flux » est-elle gratuite ?
Oui — le texte complet de « Transformer des flux » 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 « Transformer des flux » ?
Mappez et filtrez les données en circulation. 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 2 sur 4.
Combien de temps prend la leçon « Transformer des flux » ?
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