EIP (Enterprise Integration Patterns)
Les patterns d’intégration enterprise implémentés dans GoCamel.
Split
Divise un message en plusieurs parties traitées individuellement.
Propriétés de l’Exchange:
CamelSplitIndex— Index actuel (0-based)CamelSplitSize— Nombre total de partiesCamelSplitComplete— Dernière partie ?
Ces propriétés sont limitées à la branche
Elles décrivent une partie, pas le message entier, et sont retirées au
moment où le résultat agrégé est publié sur l’exchange parent. Un
processeur placé après le Split ne voit donc pas un CamelSplitComplete=true
résiduel de la dernière partie, et un Split imbriqué n’écrase pas la
comptabilité de celui qui l’englobe. La même règle s’applique à
CamelMulticast*, CamelRecipientList* et CamelRoutingSlip*.
Un corps []byte n’est pas une collection
Découper un corps []byte produit une seule partie, comme Apache Camel
qui ne découpe pas les byte[]. Le traiter comme une liste produisait un
exchange par octet — une charge utile de 10 Ko devenait 10 000 exchanges.
Pour découper délibérément du contenu binaire, découpez sur une valeur
décodée :
Avec agrégation :
Aggregate
Combine plusieurs messages en un seul.
Stratégie personnalisée:
Conditions de complétion :
| Méthode | Description |
|---|---|
SetCompletionSize(n) | Compléter après n messages |
SetCompletionTimeout(ms) | Compléter après un délai (timeout) |
SetCompletionPredicate(fn) | Compléter quand le prédicat retourne vrai |
Options de stockage :
Multicast
Envoie une copie à plusieurs destinations.
Options:
ParallelProcessing()— Exécute les branches en parallèleAggregationStrategy— Fusionne les résultats
Propriétés:
CamelMulticastIndexCamelMulticastSizeCamelMulticastComplete
Wire Tap
Envoie une copie de l’exchange à un endpoint secondaire de manière asynchrone (fire-and-forget), sans bloquer la route principale. La copie est indépendante : les modifications sur l’exchange tapé n’affectent pas la route principale et vice-versa.
Fonctionnement :
- Crée une copie profonde (deep copy) de l’exchange.
- Lance un goroutine qui envoie la copie à l’endpoint spécifié.
- Retourne immédiatement — la route principale continue sans attendre.
- L’exchange tapé utilise un contexte indépendant (
context.Background()) : il n’est pas annulé quand la route principale se termine. - Les erreurs lors du tap sont loggées mais ne sont pas propagées à la route principale.
Notes :
- Idéal pour l’audit, la journalisation ou la surveillance sans impact sur les performances.
Route.Stop()attend que tous les taps en cours soient terminés avant de s’arrêter (aucun message n’est perdu lors du shutdown).- Disponible également dans les contextes
Multicast,Pipeline,SplitetLoop.
Pipeline
Chaîne séquentielle de processors.
Choice
Content-Based Router - route les messages vers différentes destinations selon des conditions sur le contenu.
Expressions supportées :
| Expression | Description | Exemple |
|---|---|---|
${header.name} | Valeur d’en-tête | ${header.Content-Type == 'json'} |
${body} | Corps du message | ${body == 'active'} |
${body['key']} | Accès par clé map | ${body['status'] == 'pending'} |
${body[0]} | Accès par index | ${body[0] > 50} |
${exchangeProperty.prop} | Propriété Exchange | ${exchangeProperty.userId != ''} |
${date:now} | Date/heure actuelle | — |
${random(100)} | Nombre aléatoire | — |
${uuid} | UUID unique | — |
Opérateurs de comparaison :
==— Égal à!=— Différent de>— Supérieur à<— Inférieur à>=— Supérieur ou égal<=— Inférieur ou égal
Méthodes disponibles dans When/Otherwise :
Ordre d’évaluation
Les clauses When sont évaluées dans l’ordre. La première condition vraie déclenche son processor, les autres sont
ignorées. Si aucune condition ne correspond et qu’Otherwise est présent, il est exécuté.
Filter
Filtre les messages selon une condition. Les exchanges qui ne correspondent pas sont silencieusement supprimés (
ErrStopRouting).
Throttle
Limite le débit de traitement à N messages maximum par fenêtre de temps. Les exchanges excédentaires sont retardés ( jamais perdus) jusqu’à ce que la fenêtre se réinitialise.
Notes :
- Utilise un compteur à fenêtre fixe — pour un lissage plus fin, réduire la durée de la fenêtre
- Bloquant : le goroutine attend que la fenêtre autorise le passage
maxRequestsettimePerioddoivent être supérieurs à zéro ; les valeurs invalides retournent une erreur au lieu de bloquer indéfiniment- Si le contexte de l’exchange est annulé pendant l’attente de la fenêtre suivante,
Processretourne une erreur
Delay
Met le traitement de l’exchange en pause pour une durée fixe ou calculée dynamiquement. Le sleep respecte l’annulation du contexte.
Notes :
- Bloquant : le goroutine dort pendant la durée calculée
- Si le contexte de l’exchange est annulé pendant le sleep,
Processretourne une erreur - Retourner
0depuisDelayFuncpour ignorer le délai pour un exchange donné - Passer
nilàDelayFuncfait échouer la construction de la route ;NewDelayerFunc(nil)reste constructible mais retourne une erreur de traitement au lieu de paniquer
Content Enricher
Enrichit l’exchange courant avec des données provenant d’une source externe.
Fonctionnement :
- Une copie de l’exchange courant est envoyée à l’URI d’enrichissement
- La réponse est récupérée
AggregationStrategy.Aggregate(original, réponse)fusionne les deux- Sans stratégie, le corps de la réponse remplace le corps original
Idempotent Consumer
Empêche le traitement de messages en double.
Fonctionnement :
- Évalue l’expression pour obtenir une clé unique.
- Vérifie si la clé existe déjà dans le dépôt (repository).
- Si elle n’est pas présente, ajoute la clé immédiatement (bloquant les doublons concurrents) et continue le traitement.
- Si elle est déjà présente, arrête le routage (
ErrStopRouting). - Si l’exchange échoue ensuite en aval, la clé est retirée du dépôt (via une
Synchronizationde l’exchange) : une relivraison du même message sera donc traitée au lieu d’être perdue — sémantique at-least-once.
Dépôts (Repositories) :
NewMemoryIdempotentRepository(): Stockage en mémoire vive (par défaut).NewRedisIdempotentRepository(): Stockage distribué via Redis.- Des dépôts basés sur SQL peuvent être implémentés via l’interface
IdempotentRepository.
Circuit Breaker
Protège les routes contre les pannes en cascade lors de l’appel de services externes instables.
États :
- Closed (Fermé) : Les requêtes passent normalement. Les échecs sont comptabilisés.
- Open (Ouvert) : Les requêtes échouent immédiatement ou passent par le fallback.
- Half-Open (Semi-ouvert) : Après le
OpenTimeout, une seule requête d’essai à la fois sonde le service aval ; les requêtes concurrentes sont redirigées vers le fallback comme si le circuit était encore ouvert. QuandSuccessThresholdsondes réussissent, le circuit se referme. Si une sonde échoue, il s’ouvre à nouveau.
Load Balancer
Répartit les messages sur plusieurs endpoints.
Stratégies :
RoundRobin(): Distribution circulaire (par défaut).Random(): Distribution aléatoire.Custom(strategy): Fournir une implémentation personnalisée de l’interfaceLoadBalancerStrategy.
SEDA
Staged Event-Driven Architecture. Découple les producteurs et les consommateurs via une file d’attente interne (Go channel).
Fonctionnement :
- Le producteur envoie une copie de l’exchange dans une file d’attente interne.
- Le producteur rend la main immédiatement (asynchrone).
- Un ou plusieurs workers lisent la file et traitent l’exchange de manière indépendante.
Paramètres :
size: La capacité de la file d’attente (défaut 1000).concurrentConsumers: Le nombre de goroutines esclaves traitant la file (défaut 1).
Les valeurs invalides (size négatif, concurrentConsumers non positif) retombent sur
les valeurs par défaut avec un avertissement. Les paramètres sont appliqués à la
première création de l’endpoint ; les références ultérieures au même nom de file
réutilisent l’endpoint existant.
Resequencer
Remet dans l’ordre les messages arrivés de manière désordonnée en se basant sur un numéro de séquence ou un horodatage.
Fonctionnement :
- Évalue l’expression pour obtenir le numéro de séquence.
- Si la séquence est celle attendue, elle est traitée immédiatement.
- S’il s’agit d’une séquence future, elle est mise en tampon et triée.
- Si un trou (gap) persiste plus longtemps que le
Timeout, le resequenceur saute la séquence manquante et continue avec la suivante disponible.
Notes :
- Les messages sont délivrés en aval via le pipeline interne du resequenceur (les étapes
entre
Resequence(...)etEnd()). L’exchange original s’arrête là (ErrStopRouting) : les étapes placées aprèsEnd()ne voient jamais de messages hors ordre. - La livraison est sérialisée pour garantir l’ordre ; les étapes en aval ne doivent pas rerouter vers le même resequenceur.
- Les numéros de séquence dupliqués ou obsolètes sont ignorés (avec un avertissement dans les logs).
Claim Check
Réduit la taille des messages en stockant le contenu (payload) dans un dépôt externe et en ne faisant circuler qu’une référence (une clé).
Opérations :
ClaimCheckPush: Stocke le corps et le remplace par une clé. Si aucune clé n’est fournie, elle est générée.ClaimCheckPop: Récupère les données par clé, restaure le corps et les supprime du dépôt.ClaimCheckGet: Récupère les données par clé et restaure le corps (les garde dans le dépôt).ClaimCheckDiscard: Supprime les données du dépôt sans les restaurer.
Stop
Arrête le routage actuel sans erreur.
ToD (Dynamic To)
URI résolue dynamiquement à l’exécution.
Expressions supportées:
${header.<name>}— Valeur d’en-tête${property.<name>}— Valeur de propriété${body}— Corps du message
Recipient List
Envoie un message à plusieurs destinataires calculés dynamiquement à l’exécution.
Dynamic Router
Route un exchange à travers une séquence d’endpoints résolus dynamiquement en appelant une fonction de contrôle après
chaque étape. Le routage s’arrête quand la fonction retourne une chaîne vide ou nil.
Fonctionnement :
- Appelle la fonction ou l’expression de contrôle.
- Si le résultat n’est pas vide, envoie l’exchange à cet URI.
- Propage la sortie vers l’entrée et recommence à l’étape 1.
Routing Slip
Route un exchange à travers une séquence dynamique d’endpoints. Contrairement à Recipient List, chaque étape reçoit la sortie de l’étape précédente.
Propriétés Exchange pendant Routing Slip :
| Propriété | Type | Description |
|---|---|---|
CamelRoutingSlipIndex | int | Index de l’étape courante (base 0) |
CamelRoutingSlipSize | int | Nombre total d’étapes |
CamelRoutingSlipComplete | bool | Indique la dernière étape |
Headers & Properties
Set/Remove Headers
Set/Remove Properties
Gestion des erreurs
Do-Try-Catch-Finally
Gestion structurée des erreurs au sein d’une route.
Fonctionnement :
- Exécute les processeurs du bloc
DoTry. - Si une erreur survient, recherche un bloc
DoCatchcorrespondant (une chaîne vide intercepte tout). - Si une correspondance est trouvée, l’erreur est considérée comme “gérée” et le routage continue après le bloc
DoTry(sauf si le bloc catch échoue lui-même). - Le bloc
DoFinallyest toujours exécuté, qu’il y ait eu une erreur ou non.
Process
Exécute un processeur personnalisé.
Loop
Répète un segment de route un nombre de fois statique ou dynamique.
Propriétés de l’Exchange pendant la boucle :
| Propriété | Type | Description |
|---|---|---|
CamelLoopIndex | int | Index de l’itération actuelle (commence à 0) |
CamelLoopSize | int | Nombre total d’itérations de la boucle |
CamelLoopComplete | bool | Booléen indiquant s’il s’agit de la dernière itération |
Saga
Le pattern Saga fournit un mécanisme pour coordonner une série de transactions locales à travers des services distribués sans verrouiller les bases de données. Il définit un ensemble d’actions et leurs actions de compensation correspondantes (annulations).
Si une action dans le bloc Saga échoue ou retourne une erreur, l’orchestrateur de la Saga déclenche automatiquement les étapes de compensation pour toutes les actions terminées, dans l’ordre inverse (LIFO - Last In, First Out).
Vous pouvez également utiliser des fonctions Go personnalisées pour les actions et les compensations directement avec ActionFunc :
Transformer
Transforme le contenu du message.
Transform
Un EIP dédié à la transformation de message qui modifie dynamiquement le corps (body) du message à l’aide d’une fonction Go ou d’une expression Simple Language.
Message History
L’EIP Message History permet de retracer le parcours chronologique d’un message lorsqu’il traverse les différents nœuds et processeurs de la route.
Pour activer le suivi d’une route, appelez .EnableMessageHistory() sur votre route builder. Vous pouvez ensuite récupérer et inspecter l’historique des messages à l’aide de gocamel.GetMessageHistory(exchange).
Récapitulatif des patterns EIP
| Pattern | Catégorie | Description |
|---|---|---|
| Choice | Routage | Routage basé sur le contenu |
| Filter | Routage | Filtrage conditionnel |
| Idempotent Consumer | Routage | Empêche les messages en double |
| Circuit Breaker | Routage | Pattern de résilience pour les pannes |
| Load Balancer | Routage | Répartition de charge entre endpoints |
| Dynamic Router | Routage | Routage dynamique basé sur décision |
| Resequencer | Routage | Réordonnancer les messages par séquence |
| Claim Check | Routage | Stocker le payload et passer une référence |
| Throttle | Routage | Limitation de débit |
| Delay | Routage | Pause avant la prochaine étape |
| Multicast | Routage | Destinations multiples |
| Recipient List | Routage | Routage dynamique vers plusieurs destinataires |
| Routing Slip | Routage | Routage séquentiel dynamique |
| Wire Tap | Routage | Copie asynchrone fire-and-forget |
| Split | Transformation | Découpage de message |
| Aggregate | Transformation | Agrégation de messages |
| Content Enricher | Transformation | Enrichissement avec données externes |
| Transform | Transformation | Transformation de contenu |
| ToD | Endpoint | Endpoint dynamique |
| Stop | Contrôle | Arrêter le routage |
| DoTry | Contrôle | Gestion d’erreurs Try-Catch-Finally |
| SetHeader | En-têtes | Manipulation d’en-têtes |
| SetProperty | Propriétés | Propriétés d’exchange |
| Loop | Contrôle | Itération avec compteur statique/dynamique |
| Saga | Contrôle | Pattern transactionnel Saga |
| Message History | Diagnostics | Suivi du parcours de traitement des messages |