2026-10-06
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
Ajouter la dépendance Maven suivante à votre fichier pom.xml:
Mono et Flux implantent l’interface PublisherFlux.just(1, 2, 3)Flux.fromArray(array)Flux.fromIterable(list)Flux.fromStream(stream)Flux.range(1, 10) : Séquence de nombresFlux.interval(Duration) : Émission périodiqueFlux.generate() : Génération synchroneFlux.create() : Génération avec backpressureFlux.interval(Duration.ofSeconds(1))Flux.timer(Duration.ofMillis(500))Flux.create(sink -> {...})Flux.push(sink -> {...})Mono.fromFuture(future)Flux.fromCompletionStage(stage)Filtrage Simple et Déduplication - filter : Filtre selon prédicat - filter(x -> x > 0) - filter(String::isNotEmpty)
distinct : Élimine doublons
distinctUntilChanged
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é
Fusion de Flux
merge : Combine en intercalant
concat : Concatène séquentiellement
mergeSequential : Fusion ordonnée
Combinaison Synchronisée
zip : Combine par paires
combineLatest : Dernières valeurs
Utilitaires
withLatestFrom : EnrichissementstartWith : Valeurs initialessample : Échantillonnagepublic 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();
}Concept du Backpressure
Contrôle par le Subscriber
request(n): Demande n élémentsrequest(Long.MAX_VALUE): Flux illimitécancel(): Arrêt immédiatStratégies de Gestion
buffer()
drop()
latest()
error()
Concept et Utilisation
Flux avec émission contrôléeMéthodes de Contrôle
connect(): Démarre l’émissionautoConnect(n): Connexion après n abonnésrefCount(n): Gestion automatique connexion/déconnexionPatterns d’Utilisation
publish(): Crée un ConnectableFluxreplay(): Met en cache les élémentsshare(): Combine publish et refCountCas d’Usage
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:
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
E. Bruno - Programmation réactive avec Reactor