Contre-pression
Gérez les producteurs rapides en toute sécurité.
Contre-pression est une leçon Scala for Backend Engineering & Functional Programming gratuite sur CoddyKit. Ceci est la leçon 3 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.
Qu'est-ce que la contre-pression ?
La contre-pression est un mécanisme de contrôle du flux qui empêche un producteur rapide de submerger un consommateur lent. Au lieu de mettre les données en mémoire sans limite ou de les supprimer, le consommateur indique la quantité qu'il peut traiter.
Akka Streams implémente la norme Reactive Streams, dans laquelle la demande remonte vers l'amont et les éléments circulent vers l'aval.
Flux piloté par la demande
Chaque étape n'émet des éléments que lorsque l'étape suivante a signalé une demande. Un puits demande N éléments ; cette demande se propage vers l'amont jusqu'à ce qu'une source produise exactement la quantité demandée.
Ce protocole fondé sur l'extraction signifie que les producteurs n'envoient jamais plus d'éléments que les consommateurs ne peuvent en traiter.
Pourquoi est-ce important pour les pipelines ?
Sans contre-pression, un consommateur Kafka rapide qui alimente une base de données lente accumulerait des millions d'enregistrements en cours de traitement, épuisant la mémoire et faisant planter le processus.
La contre-pression ralentit naturellement l'amont jusqu'au rythme de l'étape la plus lente, ce qui garantit une utilisation stable de la mémoire sous charge.
// Fast source, slow sink: backpressure slows the source
val g =
Source(1 to 1000000)
.map(_ * 2)
.to(slowDatabaseSink)Mise en mémoire tampon interne
Entre les limites asynchrones, Akka Streams conserve un petit tampon interne (16 éléments par défaut). Il absorbe les brèves pointes d'activité afin que les étapes n'aient pas à se synchroniser sur chaque élément.
Lorsque le tampon est plein, la contre-pression s'active et l'amont cesse de produire jusqu'à ce qu'une place se libère.
import akka.stream.Attributes
val buffered =
Flow[Int]
.map(identity)
.addAttributes(Attributes.inputBuffer(initial = 32, max = 32))Tampon explicite avec stratégie de débordement
L'opérateur buffer insère un tampon explicite d'une taille choisie, avec une OverflowStrategy qui détermine le comportement lorsqu'il est plein.
Vous pouvez ainsi échanger de la mémoire contre la possibilité de découpler les vitesses du producteur et du consommateur.
import akka.stream.OverflowStrategy
val withBuffer =
Source(1 to 1000)
.buffer(size = 100, OverflowStrategy.backpressure)Stratégies de débordement
Les stratégies comprennent backpressure (ralentir l'amont), dropHead/dropTail (supprimer les éléments les plus anciens ou les plus récents), dropBuffer, dropNew et fail (terminer avec une erreur).
Les stratégies de suppression conviennent aux données en direct, comme les relevés de capteurs, lorsque les valeurs obsolètes peuvent être supprimées sans risque.
import akka.stream.OverflowStrategy
val latestWins =
liveTicks.buffer(1, OverflowStrategy.dropHead)
val strict =
liveTicks.buffer(50, OverflowStrategy.fail)Résumer avec conflate
Lorsqu'un consommateur est lent, conflate fusionne les éléments en attente en un seul à l'aide d'une fonction de combinaison, au lieu de tous les conserver dans un tampon.
Par exemple, vous pouvez réduire de nombreuses mises à jour numériques à leur somme afin que le consommateur voie toujours un agrégat des éléments qu'il n'a pas pu traiter.
val summarized =
fastMetrics
.conflate((acc, next) => acc + next)
// Slow downstream receives summed batchesÉtendre pour satisfaire la demande
expand est le dual de conflate : lorsque l'aval demande des éléments plus rapidement que l'amont n'en produit, il synthétise des éléments supplémentaires à partir de la dernière valeur reçue.
C'est utile pour continuer à émettre régulièrement la mesure la plus récente.
val repeated =
sensor.expand(last => Iterator.continually(last))
// Downstream always gets the latest sensor valueLimites asynchrones
Par défaut, les étapes fusionnées s'exécutent dans un seul acteur, sans tampon entre elles. Insérer async place une étape dans son propre acteur, ajoute un tampon et permet le parallélisme en pipeline.
Les limites asynchrones sont les endroits où se trouvent réellement les tampons de contre-pression.
val pipelined =
Source(1 to 1000)
.map(slowStep).async
.map(anotherSlowStep).async
.to(Sink.ignore)Throttle comme contrôle explicite du débit
throttle impose un débit maximal défini, ce qui génère une contre-pression vers l'amont pour le respecter. Cela protège les services externes soumis à une limitation de débit, même lorsque le consommateur pourrait aller plus vite.
Un paramètre de rafale permet de brèves pointes au-dessus du débit stable.
import scala.concurrent.duration._
val limited =
requests
.throttle(
elements = 100, per = 1.second, maximumBurst = 20,
akka.stream.ThrottleMode.Shaping)Observer la contre-pression
Vous pouvez détecter la contre-pression en surveillant les ralentissements de l'amont ou en mesurant le taux d'occupation du tampon. L'opérateur log et les attributs de flux d'Akka permettent de déterminer où un pipeline se bloque.
Un tampon constamment plein indique l'étape la plus lente, celle qui limite le débit.
val traced =
Source(1 to 100)
.log("after-source")
.map(_ * 2)
.log("after-map")
.to(Sink.ignore)Vérification rapide
Réfléchissez à la manière dont Akka Streams empêche un producteur rapide d'inonder un consommateur lent.
Récapitulatif
La contre-pression est le mécanisme central d'Akka Streams piloté par la demande : les consommateurs signalent leur demande vers l'amont afin que les producteurs ne puissent pas les submerger, ce qui maintient une mémoire limitée.
Vous avez découvert les tampons internes, l'opérateur buffer explicite avec ses stratégies de débordement, le résumé avec conflate, la génération d'éléments pour satisfaire la demande avec expand, les limites asynchrones et le contrôle volontaire du débit avec throttle. Ensuite, vous exécuterez un pipeline complet.
Questions Fréquemment Posées
La leçon « Contre-pression » est-elle gratuite ?
Oui — le texte complet de « Contre-pression » 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 « Contre-pression » ?
Gérez les producteurs rapides en toute sécurité. 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 3 sur 4.
Combien de temps prend la leçon « Contre-pression » ?
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.