Programmation réactive avec Reactor
Introduction à la programmation réactive
- Évolution des systèmes:
- Besoin croissant: Les applications modernes nécessitent une gestion efficace des ressources et une haute disponibilité.
- Charges variables: Les systèmes doivent s’adapter dynamiquement aux variations de la charge utilisateur.
- Événements en temps réel: Capacité à traiter et à réagir aux événements dès qu’ils se produisent.
- Scalabilité: Les systèmes doivent pouvoir évoluer pour gérer un grand nombre d’utilisateurs et de données.
- Efficacité: Utilisation optimale des ressources pour maximiser les performances et réduire les coûts.
- Reactive Manifesto:
- Un document fondateur publié en 2013.
- Définit les principes de la programmation réactive pour construire des systèmes réactifs.
Définition de la programmation réactive
Programmation réactive: Un paradigme de programmation orienté vers les flux de données et la propagation des changements.
Réactivité: Capacité d’un système à réagir aux événements en temps réel.
Les principes de la programmation réactive
- Responsive: Le système doit répondre rapidement et de manière cohérente.
- Resilient: Le système doit rester fonctionnel même en cas de défaillance.
- Elastic: Le système doit s’adapter à la variation de la charge.
- Message Driven: Le système doit utiliser des messages asynchrones pour garantir une communication efficace.
Comparaison avec la programmation impérative et asynchrone
- Programmation impérative: Basée sur des instructions séquentielles, moins adaptée aux systèmes réactifs.
- Programmation asynchrone: Utilise des callbacks, des promesses ou des futures pour gérer les opérations non bloquantes.
- Programmation réactive: Étend la programmation asynchrone avec des flux de données et une gestion plus fine des événements.
Des API réactives dans différents langages
- Java:
- Project Reactor: Une bibliothèque pour la programmation réactive basée sur la Reactive Streams API.
- RxJava: Une bibliothèque pour la programmation réactive basée sur le modèle des Observables.
- JavaScript:
- RxJS: Une bibliothèque pour la programmation réactive en JavaScript, utilisée notamment avec Angular.
- Bacon.js: Une bibliothèque pour la programmation réactive fonctionnelle en JavaScript.
- Python:
- RxPY: Une implémentation de Reactive Extensions pour Python.
- Streamz: Une bibliothèque pour la manipulation de flux de données en temps réel.
- C#:
- Reactive Extensions (Rx.NET): Une bibliothèque pour la programmation réactive en .NET.
- System.Reactive: Une autre implémentation de Reactive Extensions pour .NET.
- Kotlin:
- Kotlin Flow: Une API pour la programmation réactive dans Kotlin, intégrée à Kotlin Coroutines.
- RxKotlin: Une extension de RxJava pour Kotlin.
Introduction à la Reactive Streams API
- Reactive Streams API:
- Une spécification standard pour le traitement des flux de données asynchrones avec gestion de la pression (backpressure).
- Conçue pour assurer une communication fluide et efficace entre les composants réactifs.
- Objectifs:
- Interopérabilité: Permettre l’interopérabilité entre différentes bibliothèques réactives, telles que Project Reactor, RxJava, et Akka Streams.
- Gestion efficace des flux de données: Assurer que les producteurs de données ne submergent pas les consommateurs, en régulant le flux de données pour éviter les surcharges et les blocages.
- Asynchronisme: Faciliter le traitement asynchrone des flux de données pour améliorer les performances et la réactivité des applications.
- Backpressure: Introduire des mécanismes pour gérer la pression exercée par les producteurs de données sur les consommateurs, permettant ainsi une meilleure gestion des ressources et une prévention des surcharges.
- Comparaison avec Java Streams API:
- Java Streams API: Conçue pour le traitement de collections de données en mémoire.
- Reactive Streams API: Conçue pour le traitement de flux de données asynchrones avec gestion de la pression.
- Push vs Pull: Java Streams utilise un modèle pull, tandis que Reactive Streams utilise un modèle push avec gestion de la pression.
Exemple de Reactive Streams API en Java avec Project Reactor
Introduction à Project Reactor
- Présentation de Project Reactor:
- Project Reactor est une bibliothèque pour la programmation réactive en Java, basée sur la Reactive Streams API.
- Elle permet de créer des applications réactives, non bloquantes et hautement performantes.
- Installation et configuration de Reactor dans un projet Java:
Ajouter la dépendance Maven suivante à votre fichier
pom.xml:<dependency> <groupId>io.projectreactor</groupId> <artifactId>reactor-core</artifactId> <version>3.7.0</version> </dependency>
Mono
- représente une séquence asynchrone qui émet zéro ou un élément.
- Utilisé pour les opérations qui retournent un seul résultat ou aucun.
Flux
- Représente une séquence asynchrone qui émet zéro, un ou plusieurs éléments.
- Utilisé pour les opérations qui retournent plusieurs résultats.
Publisher interface
- Interface fondamentale de Reactive Streams
- Source de données qui émet des éléments de manière asynchrone
- Point de départ du flux de données réactif
MonoetFluximplantent l’interface Publisher
public interface Publisher<T> {
void subscribe(Subscriber<? super T> subscriber);
}Génération de Flux
- Création depuis des Données
- Valeurs Discrètes
Flux.just(1, 2, 3)Flux.fromArray(array)Flux.fromIterable(list)Flux.fromStream(stream)
- Génération Programmée
Flux.range(1, 10): Séquence de nombresFlux.interval(Duration): Émission périodiqueFlux.generate(): Génération synchroneFlux.create(): Génération avec backpressure
- Valeurs Discrètes
- Sources Asynchrones
- Temps
Flux.interval(Duration.ofSeconds(1))Flux.timer(Duration.ofMillis(500))
- Callbacks
Flux.create(sink -> {...})Flux.push(sink -> {...})
- CompletableFuture
Mono.fromFuture(future)Flux.fromCompletionStage(stage)
- Temps
Opérateurs de Transformation
- map: Transformation synchrone élément par élément
- Conversion de types
- Formatage de données
- Calculs simples
- flatMap: Transformation asynchrone avec aplatissement
- Appels services externes
- Opérations I/O
- Flux imbriqués
- transform: Modification globale du flux
- Composition de transformations
- Réutilisation de logique
- Modification de la chaîne
- switchMap: Gestion des changements rapides
- Annulation automatique
- Dernier flux uniquement
- Recherche interactive
Opérateurs de Filtrage dans Project Reactor
Filtrage Simple et Déduplication - filter : Filtre selon prédicat - filter(x -> x > 0) - filter(String::isNotEmpty)
distinct: Élimine doublons- Utilise equals/hashCode
- Option keySelector
distinctUntilChanged- Élimine doublons consécutifs
- Compare éléments adjacents
Limitation - take(n) : n premiers éléments - takeLast(n) : n derniers - takeWhile : condition vraie - takeUntil : jusqu’à condition
Saut d’Éléments - skip(n) : Ignore n premiers - skipLast(n) : Ignore n derniers - skipWhile : Ignore si vrai - skipUntil : Ignore jusqu’à vrai
Sélection Spécifique - elementAt : Position précise - next : Premier élément - last : Dernier élément - single : Vérifie unicité
Opérateurs de Combinaison de Flux
Fusion de Flux
merge: Combine en intercalant- Ordre d’arrivée préservé
- Pas de synchronisation
- Parallèle
concat: Concatène séquentiellement- Ordre strict respecté
- Séquentiel
- Sans entrelacement
mergeSequential: Fusion ordonnée- Souscription parallèle
- Émission ordonnée
- Préserve ordre
Combinaison Synchronisée
zip: Combine par paires- Attend les deux sources
- Transformation possible
- 1:1 matching
combineLatest: Dernières valeurs- Combine valeurs récentes
- Réactif aux changements
- N:M possible
Utilitaires
withLatestFrom: EnrichissementstartWith: Valeurs initialessample: Échantillonnage
Interface Subscriber
- Interface fondamentale de Reactive Streams
- Consommateur de données qui reçoit les éléments de manière asynchrone
- Point de terminaison du flux de données réactif
public interface Subscriber<T> {
// Called first when subscribing to a Publisher
// Provides the Subscription object to control the flow
// Should call subscription.request() to start receiving items
void onSubscribe(Subscription s);
// Called for each item emitted by the Publisher
// Can be called 0 to N times based on the number of items
// Will not be called after onError or onComplete
void onNext(T t);
// Called when an error occurs in the stream
// Terminal operation - no more calls will happen after this
// Must not be called after onComplete
void onError(Throwable t);
// Called when the Publisher has no more items to emit
// Terminal operation - no more calls will happen after this
// Must not be called after onError
void onComplete();
}Gestion des erreurs avec onErrorReturn
Gestion des erreurs avec onErrorResume
Gestion du Backpressure dans Project Reactor
Concept du Backpressure
- Régulation du flux producteur/consommateur
- Évite la surcharge du consommateur
- Communication bidirectionnelle
Contrôle par le Subscriber
request(n): Demande n élémentsrequest(Long.MAX_VALUE): Flux illimitécancel(): Arrêt immédiat
Stratégies de Gestion
buffer()- Stockage temporaire
- Mémoire tampon configurable
- Préserve les données
drop()- Suppression des excédents
- Sans impact mémoire
- Perte acceptée
latest()- Garde valeur récente
- Mise à jour continue
- Sampling
error()- Signalement erreur
- Protection système
- Circuit breaker
ConnectableFlux - Flux Partagé et Contrôlé
Concept et Utilisation
- Variante de
Fluxavec émission contrôlée - Partage un flux entre multiples abonnés
- “Hot” publisher - émission indépendante des abonnements
Méthodes de Contrôle
connect(): Démarre l’émissionautoConnect(n): Connexion après n abonnésrefCount(n): Gestion automatique connexion/déconnexion
Patterns d’Utilisation
publish(): Crée un ConnectableFluxreplay(): Met en cache les élémentsshare(): Combine publish et refCount
Cas d’Usage
- Diffusion d’événements
- Partage de données en temps réel
- Cache partagé
- Multicast de flux
Introduction aux tests réactifs avec StepVerifier
Principes Fondamentaux - Test déclaratif des flux réactifs - Vérification pas à pas du comportement - Support du backpressure - Tests synchrones et asynchrones
Méthodes Principales - expectNext(): Vérifie valeur suivante - expectNextCount(): Nombre d’éléments - expectComplete(): Vérifie fin normale - expectError(): Vérifie erreur - verifyTimeout(): Test avec timeout - ajouter la dépendance suivante à votre fichier pom.xml:
<dependency>
<groupId>io.projectreactor</groupId>
<artifactId>reactor-test</artifactId>
<version>3.7.0</version>
<scope>test</scope>
</dependency>Débogage et log pour le code réactif
Opérateur log() - Trace tous les événements du flux - Signaux: onNext, onError, onComplete - Requêtes backpressure - Options de configuration: - Niveau de log - Logger personnalisé - Catégorie
Points de Contrôle - checkpoint() - Capture stack trace - Identifie source d’erreurs - Impact performance minimal
Points Clés à Retenir
- Évolution des Paradigmes
- Du synchrone à l’asynchrone
- Des threads aux flux réactifs
- Performance et scalabilité
- Approche Moderne
- Non-blocking par défaut
- Gestion des ressources optimisée
- Tests et monitoring intégrés
- Bonnes Pratiques
- Gestion des erreurs systématique
- Contrôle du backpressure
- Découpage en flux atomiques
- Monitoring continu