Exécuter un pipeline
Matérialisez et exécutez un graphe.
Exécuter un pipeline est une leçon Scala for Backend Engineering & Functional Programming gratuite sur CoddyKit. Ceci est la leçon 4 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.
Du modèle à l'exécution
Jusqu'à présent, le pipeline n'était qu'un modèle abstrait. La matérialisation est le processus qui transforme ce modèle en acteurs en cours d'exécution qui déplacent réellement les données.
Rien ne se passe tant que vous n'exécutez pas explicitement le graphe, ce qui rend Akka Streams composable et réutilisable.
L'ActorSystem
La matérialisation nécessite un ActorSystem, qui fournit les threads et le répartiteur utilisés par les étapes du flux. Dans les versions modernes d'Akka, le système sert également de matérialiseur implicite.
Un ActorSystem sert généralement une application entière et de nombreux flux concurrents.
import akka.actor.ActorSystem
implicit val system: ActorSystem =
ActorSystem("data-pipeline")
import system.dispatcher // ExecutionContextrunWith
La manière la plus directe d'exécuter une source est d'utiliser runWith, qui lui associe un puits et la matérialise en une seule étape, en renvoyant la valeur matérialisée de ce puits.
Ici, le résultat est un Future[Int] qui se termine avec la somme lorsque le flux s'achève.
import akka.stream.scaladsl.{Source, Sink}
import scala.concurrent.Future
val total: Future[Int] =
Source(1 to 100).runWith(Sink.fold(0)(_ + _))Exécuter un RunnableGraph
Si vous avez déjà construit un RunnableGraph fermé avec to ou toMat, appelez run() pour le matérialiser. La valeur renvoyée est la valeur matérialisée conservée par le graphe.
Cela sépare clairement la construction du pipeline de son exécution.
import akka.stream.scaladsl.Keep
import scala.concurrent.Future
val graph =
Source(1 to 100)
.toMat(Sink.fold(0)(_ + _))(Keep.right)
val result: Future[Int] = graph.run()Opérateurs d'exécution pratiques
Les sources proposent des raccourcis : runForeach, runFold et runReduce associent chacun le puits correspondant et l'exécutent immédiatement.
Ils sont concis pour les opérations terminales courantes sur une source.
import scala.concurrent.Future
val printed: Future[akka.Done] =
Source(1 to 10).runForeach(println)
val sum: Future[Int] =
Source(1 to 10).runFold(0)(_ + _)Utiliser la Future de résultat
Les puits terminaux renvoient une Future qui se termine lorsque le flux s'achève ou échoue. Enregistrez des fonctions de rappel avec onComplete pour réagir au succès ou à l'erreur.
Utilisez le répartiteur de l'ActorSystem comme ExecutionContext implicite pour ces fonctions de rappel.
import scala.util.{Success, Failure}
total.onComplete {
case Success(value) => println(s"Sum = $value")
case Failure(ex) => println(s"Failed: ${ex.getMessage}")
}Un pipeline réaliste
Un pipeline de données classique lit depuis une source, transforme les données avec des flux, effectue des entrées-sorties asynchrones avec mapAsync, regroupe les éléments avec grouped et écrit dans un puits.
Chaque étape est petite et l'ensemble est matérialisé par un seul run.
val done =
lineSource
.map(parse)
.mapAsync(4)(validate)
.grouped(500)
.runWith(bulkWriteSink)Redémarrer les flux en échec
Pour renforcer la résilience, entourez une source ou un flux avec RestartSource.withBackoff afin que les défaillances transitoires (comme une connexion interrompue) déclenchent un redémarrage automatique avec temporisation exponentielle.
Les pipelines d'ingestion de longue durée restent ainsi actifs sans supervision manuelle.
import akka.stream.scaladsl.RestartSource
import akka.stream.RestartSettings
import scala.concurrent.duration._
val resilient = RestartSource.withBackoff(
RestartSettings(1.second, 30.seconds, 0.2))(() => flakySource)Arrêt progressif avec KillSwitch
Un KillSwitch permet à du code externe d'arrêter proprement un flux en cours d'exécution. Insérez KillSwitches.single avec viaMat et conservez sa valeur matérialisée pour appeler shutdown() ultérieurement.
C'est essentiel pour les flux de longue durée qui doivent s'arrêter lors de l'arrêt de l'application.
import akka.stream.{KillSwitches, KillSwitch}
import akka.stream.scaladsl.Keep
val (switch, done) =
source
.viaMat(KillSwitches.single)(Keep.right)
.toMat(Sink.ignore)(Keep.both)
.run()
// later: switch.shutdown()Libérer les ressources
Lorsque l'application se termine, arrêtez l'ActorSystem pour libérer ses threads. Enchaînez l'appel d'arrêt après la Future d'achèvement du flux afin que l'arrêt se déroule correctement.
Une fuite d'ActorSystem maintient la JVM en vie et conserve les ressources ouvertes.
done.onComplete { _ =>
system.terminate()
}Réutiliser le matérialiseur
Matérialiser plusieurs fois le même modèle crée des flux indépendants en cours d'exécution qui partagent les ressources de l'ActorSystem. Le modèle lui-même reste immuable et dépourvu d'effets de bord.
Vous pouvez ainsi définir un pipeline une seule fois et l'exécuter à la demande pour chaque tâche entrante.
val blueprint =
Source(1 to 5).toMat(Sink.seq)(Keep.right)
val run1 = blueprint.run()
val run2 = blueprint.run() // independent executionVérification rapide
Réfléchissez à ce qui est nécessaire pour qu'un flux traite réellement des éléments.
Récapitulatif
Exécuter un pipeline signifie matérialiser un modèle avec un ActorSystem via run, runWith ou des opérateurs pratiques, chacun renvoyant un résultat Future.
Vous avez découvert des pipelines réalistes à plusieurs étapes, les redémarrages automatiques avec temporisation, l'arrêt progressif avec KillSwitch, le nettoyage des ressources avec system.terminate() et la réutilisation sûre d'un modèle immuable lors d'exécutions indépendantes.
Apprends Scala avec un tuteur IA — gratuit
Écris et exécute du vrai code dans ton navigateur, obtiens de l'aide instantanée d'un tuteur IA disponible 24h/24, et reprends là où tu t'es arrêté sur le web ou dans l'app.
- Cours
- 39
- Leçons
- 143
Questions Fréquemment Posées
La leçon « Exécuter un pipeline » est-elle gratuite ?
Oui — le texte complet de « Exécuter un pipeline » 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 « Exécuter un pipeline » ?
Matérialisez et exécutez un graphe. 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 4 sur 4.
Combien de temps prend la leçon « Exécuter un pipeline » ?
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